Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -152,7 +152,7 @@ func main() {

// Collect the ids of any messages that failed. They won't be deleted,
// so they return to the queue once their visibility timeout expires.
var failed []interface{}
var failed []any

for _, msg := range msgs {
// Assert the type of message to get the body or any other attributes
Expand Down
64 changes: 64 additions & 0 deletions v2/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
# Formigo - distributed SQS worker pools.

> This library relies heavily on [`github.com/alitto/pond/v2`](https://github.com/alitto/pond), a very flexible and well-tested worker pool library.

Formigo is a fast, reliable worker pool consumer for SQS.

Basic features of Pond:

- automatic scaling to available resources/limits based on incoming queue pressure.
- fire & forget queue submission
- fire & **wait for a response** submission to a queue (e.g. for dependent jobs, or deferred behaviour).
- clean, graceful exit when shutting down.

## Usage

See `cmd/echo/main.go` for a complete example of how to use v2.

### Consumers

Your consumer should adhere to the following:

```go
import formigo "github.com/Pod-Point/go-queue-worker/v2"

func (context.Context, formigo.Message) error
```

A `Decode` method is provided on formigo.Message to help ensure proper decoding with JSON SQS messages:

```go
type Message struct {
Foo string `json:"foo"`
}

func (ctx context.Context, msg formigo.Message) error {
var body Message
if err := msg.Decode(&body); err != nil {
return err
}
// pass to your internal consumer, etc.
}
```

### Error Reporting

If you need errors to be sent to Sentry, structured logging, etc, feel free to use `WithReporter` when setting up a manager.

```go
formigo.WithReporter(func(err error) {
sentry.CaptureException(err)
})
```

### Other Options

All other options should be fairly self-explanatory, and have godocs and defaults.

## Shutting Down

Ensure the context passed to `manager.Run(ctx)` is cancelable, preferrably via `signal.NotifyContext`.

If it isn't, the program will have no exit condition until fully terminated by the operating system.

See the example consumer to see how this is done.
82 changes: 82 additions & 0 deletions v2/cmd/echo/main.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,82 @@
// Package main is a test consumer for a localstack queue.
// This is not meant to be a production config, though it is set up exactly as you would use it there.
package main

import (
"context"
"errors"
"fmt"
"log"
"log/slog"
"os"
"os/signal"
"syscall"
"time"

"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/service/sqs"

formigo "github.com/Pod-Point/go-queue-worker/v2"
)

type Message struct {
Level string `json:"level"`
Message string `json:"message"`
}

func main() {
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGILL, syscall.SIGHUP)
defer cancel()

cfg, err := config.LoadDefaultConfig(ctx)
if err != nil {
log.Printf("unable to load config: %v\n", err)
os.Exit(2)
}

client := formigo.NewSQSClient(sqs.NewFromConfig(cfg), &sqs.ReceiveMessageInput{
QueueUrl: aws.String("http://sqs.eu-west-1.localhost.localstack.cloud:4566/000000000000/formigo-echo-testing"),
MaxNumberOfMessages: 5,
VisibilityTimeout: 6,
WaitTimeSeconds: 10,
})

manager := formigo.NewManager(client,
formigo.WithDeadline(time.Second*4),
formigo.WithFetchConcurrency(2),
formigo.WithFetchDelay(time.Second*1),
formigo.WithWorkerConcurrency(3),
formigo.WithReporter(func(err error) {
slog.Error(err.Error(), slog.String("type", fmt.Sprintf("%T", err)))
}),
formigo.WithConsumer(func(ctx context.Context, msg formigo.Message) error {
var body Message
if err := msg.Decode(&body); err != nil {
return err
}

switch body.Level {
case "error":
return errors.New(body.Message)
case "info":
slog.InfoContext(ctx, body.Message)
case "warn", "warning":
slog.WarnContext(ctx, body.Message)
case "debug":
slog.DebugContext(ctx, body.Message)
}

return nil
}),
formigo.WithLogger(slog.New(
slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{
Level: slog.LevelInfo,
}),
)),
)

if err := manager.Run(ctx); err != nil {
log.Fatal(err)
}
}
57 changes: 57 additions & 0 deletions v2/cmd/producer/main.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
// Package main is a test producer for a localstack queue.
// This is not meant to be a production config, though it is set up exactly as you would use it there.
package main

import (
"context"
"encoding/json"
"fmt"
"log"
"os"

"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/service/sqs"
)

func main() {
ctx := context.Background()

cfg, err := config.LoadDefaultConfig(ctx)
if err != nil {
log.Printf("unable to load config: %v\n", err)
os.Exit(2)
}

if len(os.Args[1:]) != 2 {
fmt.Println("USAGE: ./echo [level] \"[message]\"")
}

level := os.Args[1]
body := os.Args[2]

b, err := json.Marshal(struct {
Level string `json:"level"`
Body string `json:"message"`
}{
Level: level,
Body: body,
})
if err != nil {
log.Fatal(err)
}

client := sqs.NewFromConfig(cfg)

// duplicate for testing
for range 10 {
output, err := client.SendMessage(ctx, &sqs.SendMessageInput{
MessageBody: aws.String(string(b)),
QueueUrl: aws.String("http://sqs.eu-west-1.localhost.localstack.cloud:4566/000000000000/formigo-echo-testing"),
})
if err != nil {
log.Fatal(err)
}
fmt.Printf("message id: %s\n", *output.MessageId)
}
}
29 changes: 29 additions & 0 deletions v2/go.mod
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
module github.com/Pod-Point/go-queue-worker/v2

go 1.26

require (
github.com/alitto/pond/v2 v2.7.1
github.com/aws/aws-sdk-go-v2 v1.47.1
github.com/aws/aws-sdk-go-v2/config v1.33.6
github.com/aws/aws-sdk-go-v2/service/sqs v1.34.5
github.com/stretchr/testify v1.9.0
)

require (
github.com/aws/aws-sdk-go-v2/credentials v1.20.6 // indirect
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.20.1 // indirect
github.com/aws/aws-sdk-go-v2/internal/configsources v1.5.4 // indirect
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.4 // indirect
github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.4 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.19 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.14.4 // indirect
github.com/aws/aws-sdk-go-v2/service/signin v1.10.1 // indirect
github.com/aws/aws-sdk-go-v2/service/sso v1.38.1 // indirect
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.43.1 // indirect
github.com/aws/aws-sdk-go-v2/service/sts v1.51.1 // indirect
github.com/aws/smithy-go v1.28.1 // indirect
github.com/davecgh/go-spew v1.1.1 // indirect
github.com/pmezard/go-difflib v1.0.0 // indirect
gopkg.in/yaml.v3 v3.0.1 // indirect
)
42 changes: 42 additions & 0 deletions v2/go.sum
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
github.com/alitto/pond/v2 v2.7.1 h1:QxMbcfjcVTa0pyxX5Ib1226mM8u8D7gKUVkCUU4DYIw=
github.com/alitto/pond/v2 v2.7.1/go.mod h1:xkjYEgQ05RSpWdfSd1nM3OVv7TBhLdy7rMp3+2Nq+yE=
github.com/aws/aws-sdk-go-v2 v1.47.1 h1:uOIZnp4PK3ZhKI0dNrJrhTEsLxbpXHTAJlwoS1pvAtw=
github.com/aws/aws-sdk-go-v2 v1.47.1/go.mod h1:bttEH6JqnUL8LepvDVfdrds/fZ5bCIxzpe3abyUrhDU=
github.com/aws/aws-sdk-go-v2/config v1.33.6 h1:MBjkSTLczek/UgiK+EYPIoRTqE7gP8vtW3OFbFo7Nug=
github.com/aws/aws-sdk-go-v2/config v1.33.6/go.mod h1:grRAFzdAZJrwcbasJRg2MPvIrVjtlfXllHssN6+E1JE=
github.com/aws/aws-sdk-go-v2/credentials v1.20.6 h1:NpAFXCU7NzXNkdGK3zQTtsRJ+3v9tZQV0xcdRw8uBdw=
github.com/aws/aws-sdk-go-v2/credentials v1.20.6/go.mod h1:mcZCoiPnyMvP8VMNbygNX5lLqSlkYJIMPODylQMurOk=
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.20.1 h1:8gALAAmacnIXh+z6VkdDanv4/IkG5APdg4DZLDTmLog=
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.20.1/go.mod h1:Z7IJhJU+poOdJjUR2wpyY21ossQ1XS/R3Lk9Msq5kM4=
github.com/aws/aws-sdk-go-v2/internal/configsources v1.5.4 h1:CLq4+8UHCI+ZZYl/EuJxXovaIVN2xeeT8JV+dsApQ5E=
github.com/aws/aws-sdk-go-v2/internal/configsources v1.5.4/go.mod h1:Wv4q5sAM04xAMkoOedxLx2inVf6K5FdxYp+A61L+q/0=
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.4 h1:dD4MR81I7YkpEBRk6UP9rocC2QnT3qVuXwzlYTtfGEs=
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.4/go.mod h1:EcXV1kAFd5XwSkDHlj94gnF3q5CkJyYiIJfH8N0VmrE=
github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.4 h1:7Wo47d/xn/7KttCSBd8EGYeZ7ULRFRkUHr6vkZPBzVQ=
github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.4/go.mod h1:tDB2IVC1xC3vX8o+6uRlzhTxP3g1b77CZXFX/oD2FnQ=
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.19 h1:bAdDl/HkGCcGPoe25ToSHEw23VIxt6CT5fLcg111BKg=
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.19/go.mod h1:KaUzbLxv4CeSxh6ZCl9B4m7CuFenS8kUEaDs+f/DQr4=
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.14.4 h1:29SvnfGhXjTl8ONxFwbj2rs6lbhiFXD2CgFQmbT/bXY=
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.14.4/go.mod h1:wm04I5DMuNVvZHFe/dHnUxincvNbbK7AiNBbYsQivek=
github.com/aws/aws-sdk-go-v2/service/signin v1.10.1 h1:DzCCWLzcIRQ77F3DEUljud7bEjTgFOIKXP52NmVRyhU=
github.com/aws/aws-sdk-go-v2/service/signin v1.10.1/go.mod h1:xpo/geVldu8payT375WekctUzopG/hBU7miiqItMUlw=
github.com/aws/aws-sdk-go-v2/service/sqs v1.34.5 h1:HYyVDOC2/PIg+3oBX1q0wtDU5kONki6lrgIG0afrBkY=
github.com/aws/aws-sdk-go-v2/service/sqs v1.34.5/go.mod h1:7idt3XszF6sE9WPS1GqZRiDJOxw4oPtlRBXodWnCGjU=
github.com/aws/aws-sdk-go-v2/service/sso v1.38.1 h1:Umtl/0YZhng4xndfW3lKJrYYP7NLEjI6bGXVomwLcs0=
github.com/aws/aws-sdk-go-v2/service/sso v1.38.1/go.mod h1:rRD/dnm7q0HYE/I5TMaPgkWyyUGLcwuxHLABsLnQ3e0=
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.43.1 h1:orIWdNiLgzrhu/11RcPPKO/SBzUUymbUQuZbSPImghg=
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.43.1/go.mod h1:skwM/xsbR/1ReUTesv9BhpJp1VjajR7DWQnuVLwiXsQ=
github.com/aws/aws-sdk-go-v2/service/sts v1.51.1 h1:0HOqZXRvMytH6bFHVIc0oJX07sZjfhz0zXtjs6gdE8s=
github.com/aws/aws-sdk-go-v2/service/sts v1.51.1/go.mod h1:26zA0GhDrLo+yiLI2yXWxqB1PdsShfLikoI7GOEgugM=
github.com/aws/smithy-go v1.28.1 h1:R/nXH00c8qcfCzQVELtRw+eLQWtzv+VAIEFJ1/xxXlQ=
github.com/aws/smithy-go v1.28.1/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqxX9i1Fc=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/stretchr/testify v1.9.0 h1:HtqpIVDClZ4nwg75+f6Lvsy/wHu+3BoSGCbBAcpTsTg=
github.com/stretchr/testify v1.9.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
Loading
Loading