From 83aace75b14321a8e384b944fdb76d35f91432e6 Mon Sep 17 00:00:00 2001 From: Lee Graham-Merrett Date: Tue, 22 Sep 2026 10:47:05 +0100 Subject: [PATCH 1/2] docs: update readme --- .gitignore | 3 + README.md | 230 ++++++++++++++++++++++++++++------------------------- 2 files changed, 124 insertions(+), 109 deletions(-) diff --git a/.gitignore b/.gitignore index 3b735ec..71fc0e4 100644 --- a/.gitignore +++ b/.gitignore @@ -14,6 +14,9 @@ # Output of the go coverage tool, specifically when used with LiteIDE *.out +# IDEs and editors +/.idea + # Dependency directories (remove the comment below to include it) # vendor/ diff --git a/README.md b/README.md index eef1373..e7ff1b1 100644 --- a/README.md +++ b/README.md @@ -33,61 +33,59 @@ Let's create some simple examples to demonstrate how to use this library to proc ### Basic example ```go +package main + import ( - "context" - "fmt" - "log" - - formigo "github.com/Pod-Point/go-queue-worker" - workerSqs "github.com/Pod-Point/go-queue-worker/clients/sqs" - - "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" - "github.com/aws/aws-sdk-go-v2/service/sqs/types" + "context" + "log" + + formigo "github.com/Pod-Point/go-queue-worker" + + "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/service/sqs" + "github.com/aws/aws-sdk-go-v2/service/sqs/types" ) func main() { - ctx := context.Background() - - awsCfg, err := config.LoadDefaultConfig(ctx) - if err != nil { - log.Fatalln("Unable to create AWS config", err) - } - - sqsSvc := sqs.NewFromConfig(awsCfg) - sqsClient, err := formigo.NewSqsClient(ctx, formigo.SqsClientConfiguration{ - Svc: sqsSvc, - ReceiveMessageInput: &sqs.ReceiveMessageInput{ - QueueUrl: &queueUrl, - MaxNumberOfMessages: 1, - VisibilityTimeout: 30, - WaitTimeSeconds: 20, - }, - }) - if err != nil { - return fmt.Errorf("unable to create sqs client: %w", err) - } - - wkr := formigo.NewWorker(formigo.Configuration{ - Client: sqsClient, - Concurrency: 100, - Consumer: formigo.NewMessageConsumer(formigo.MessageConsumerConfiguration{ - Handler: func(ctx context.Context, msg formigo.Message) error { - log.Println("Got Message", msg.Content()) - - // Assert the type of message to get the body or any other attributes - log.Println("Message body", *msg.Content().(types.Message).Body) - - return nil - }, - }), - }) - - err = wkr.Run(ctx) - if err != nil { - log.Fataln("Worker stopped with error", err) - } + ctx := context.Background() + + queueUrl := "https://sqs.eu-west-1.amazonaws.com/123456789012/my-queue" + + awsCfg, err := config.LoadDefaultConfig(ctx) + if err != nil { + log.Fatalln("Unable to create AWS config", err) + } + + sqsSvc := sqs.NewFromConfig(awsCfg) + sqsClient, err := formigo.NewSqsClient(ctx, formigo.SqsClientConfiguration{ + Svc: sqsSvc, + ReceiveMessageInput: &sqs.ReceiveMessageInput{ + QueueUrl: &queueUrl, + MaxNumberOfMessages: 10, + VisibilityTimeout: 30, + WaitTimeSeconds: 20, + }, + }) + if err != nil { + log.Fatalln("Unable to create sqs client", err) + } + + wkr := formigo.NewWorker(formigo.Configuration{ + Client: sqsClient, + Concurrency: 100, + Consumer: formigo.NewMessageConsumer(formigo.MessageConsumerConfiguration{ + Handler: func(ctx context.Context, msg formigo.Message) error { + // Assert the type of message to get the body or any other attributes + log.Println("Message body", *msg.Content().(types.Message).Body) + + return nil + }, + }), + }) + + if err := wkr.Run(ctx); err != nil { + log.Fatalln("Worker stopped with error", err) + } } ``` @@ -95,7 +93,7 @@ In this example, we have created a worker that consumes messages one at a time f By default, the worker's concurrency is set to 100, meaning it can process up to 100 messages concurrently, optimizing throughput and efficiency. -If any errors occur during message handling, the worker will log them using log.PrintLn by default. Additionally, the worker is configured to stop if it encounters more than 3 errors within any 120-second interval. +If any errors occur during message handling, the worker will log them using log.Println by default. Additionally, the worker is configured to stop if it encounters more than 3 errors within any 120-second interval. Please note that these are the default settings, and you can customize the concurrency level, error handling, and other parameters to suit your specific requirements. @@ -103,75 +101,88 @@ Please note that these are the default settings, and you can customize the concu ### Batching ```go +package main + import ( - "context" - "fmt" - "log" + "context" + "log" + "time" - formigo "github.com/Pod-Point/go-queue-worker" - workerSqs "github.com/Pod-Point/go-queue-worker/clients/sqs" + formigo "github.com/Pod-Point/go-queue-worker" - "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" - "github.com/aws/aws-sdk-go-v2/service/sqs/types" + "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/service/sqs" + "github.com/aws/aws-sdk-go-v2/service/sqs/types" ) func main() { - ctx := context.Background() - - awsCfg, err := config.LoadDefaultConfig(ctx) - if err != nil { - log.Fatalln("Unable to create AWS config", err) - } - - sqsSvc := sqs.NewFromConfig(awsCfg) - sqsClient, err := formigo.NewSqsClient(ctx, formigo.SqsClientConfiguration{ - Svc: sqsSvc, - ReceiveMessageInput: &sqs.ReceiveMessageInput{ - QueueUrl: &queueUrl, - MaxNumberOfMessages: 1, - VisibilityTimeout: 30, - WaitTimeSeconds: 20, - }, - }) - if err != nil { - return fmt.Errorf("unable to create sqs client: %w", err) - } - - wkr := formigo.NewWorker(formigo.Configuration{ - Client: sqsClient, - Concurrency: 100, - Consumer: formigo.BatchConsumer(formigo.BatchConsumerConfiguration{ - BufferConfig: formigo.BatchBufferConfiguration{ - Size: 100, - Timeout: time.Second * 5, - }, - Handler: func(ctx context.Context, msgs []formigo.Message) error { - log.Printf("Got %d messages to process\n", len(msgs) - - // Assert the type of message to get the body or any other attributes - - for i, msg := range msgs { - log.Printf("Message %d body: %s", i, *msg.Content().(types.Message).Body) - } - - return nil - }, - }), - }) - - err = wkr.Run(ctx) - if err != nil { - log.Fataln("Worker stopped with error", err) - } + ctx := context.Background() + + queueUrl := "https://sqs.eu-west-1.amazonaws.com/123456789012/my-queue" + + awsCfg, err := config.LoadDefaultConfig(ctx) + if err != nil { + log.Fatalln("Unable to create AWS config", err) + } + + sqsSvc := sqs.NewFromConfig(awsCfg) + sqsClient, err := formigo.NewSqsClient(ctx, formigo.SqsClientConfiguration{ + Svc: sqsSvc, + ReceiveMessageInput: &sqs.ReceiveMessageInput{ + QueueUrl: &queueUrl, + MaxNumberOfMessages: 10, + VisibilityTimeout: 30, + WaitTimeSeconds: 20, + }, + }) + if err != nil { + log.Fatalln("Unable to create sqs client", err) + } + + wkr := formigo.NewWorker(formigo.Configuration{ + Client: sqsClient, + Concurrency: 100, + Consumer: formigo.NewBatchConsumer(formigo.BatchConsumerConfiguration{ + BufferConfig: formigo.BatchConsumerBufferConfiguration{ + Size: 100, + Timeout: time.Second * 5, + }, + Handler: func(ctx context.Context, msgs []formigo.Message) (formigo.BatchResponse, error) { + log.Printf("Got %d messages to process\n", len(msgs)) + + // 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{} + + for _, msg := range msgs { + // Assert the type of message to get the body or any other attributes + body := *msg.Content().(types.Message).Body + + if err := process(ctx, body); err != nil { + log.Println("Unable to process message", err) + failed = append(failed, msg.Id()) + } + } + + return formigo.BatchResponse{FailedMessagesId: failed}, nil + }, + }), + }) + + if err := wkr.Run(ctx); err != nil { + log.Fatalln("Worker stopped with error", err) + } } + +func process(ctx context.Context, body string) error { return nil } ``` In this example, we have created a worker that efficiently consumes batches of messages from an AWS SQS queue. The handler will be invoked either when the buffer is full or when a specified timeout expires. It's essential to note that the timer starts as soon as the first message is added to the buffer. +Returning message ids in `BatchResponse.FailedMessagesId` marks those messages as failed: they are not deleted from the queue, so they become visible again once their visibility timeout expires. Every other message in the batch is deleted. Returning an error from the handler instead fails the whole batch and reports it to the `ErrorConfig.ReportFunc`. + By processing messages in batches, the worker can significantly enhance throughput for specific use cases or reduce resource consumption. For instance, it can be leveraged for batch insertions or deletions. ## Configuration @@ -182,7 +193,8 @@ By processing messages in batches, the worker can significantly enhance throughp | Concurrency | Number of Go routines that process the messages from the Queue. Higher values are useful for slow I/O operations in the consumer's handler. | 100 | | Retrievers | Number of Go routines that retrieve messages from the Queue. Higher values are helpful for slow networks or when consumers are quicker. | 1 | | ErrorConfig | Defines the error threshold and interval for worker termination and error reporting function. | None | -| Consumer | The message consumer, either MessageConsumer or BatchConsumer. | None | +| Consumer | The message consumer, built with either `NewMessageConsumer` or `NewBatchConsumer`. This is a required configuration. | None | +| DeleterConfig | Size and timeout of the buffer that batches messages before they are deleted from the queue. | 10, 500ms | ## License From 91b4549b96a6b6534f0526835b8a6bcdfd343e63 Mon Sep 17 00:00:00 2001 From: Lee Graham-Merrett Date: Tue, 22 Sep 2026 11:31:05 +0100 Subject: [PATCH 2/2] feat: update minimum go version & ci upgrades --- .github/workflows/lint.yml | 15 +++++++-------- .github/workflows/release.yml | 4 ++-- .github/workflows/test.yml | 4 ++-- go.mod | 2 +- 4 files changed, 12 insertions(+), 13 deletions(-) diff --git a/.github/workflows/lint.yml b/.github/workflows/lint.yml index 1c045b6..f6b68b2 100644 --- a/.github/workflows/lint.yml +++ b/.github/workflows/lint.yml @@ -15,16 +15,15 @@ jobs: name: lint runs-on: ubuntu-latest steps: - - uses: actions/checkout@v3 - - uses: actions/setup-go@v4 + - uses: actions/checkout@v5 + - uses: actions/setup-go@v5 with: go-version-file: 'go.mod' cache: false - name: golangci-lint - uses: golangci/golangci-lint-action@v3 + uses: golangci/golangci-lint-action@v8 with: - version: v1.55 - skip-pkg-cache: true + version: v2.13 govulncheck: runs-on: ubuntu-latest @@ -38,14 +37,14 @@ jobs: commitlint: runs-on: ubuntu-latest steps: - - uses: actions/checkout@v3 + - uses: actions/checkout@v5 with: fetch-depth: 0 - name: Install NodeJS - uses: actions/setup-node@v4 + uses: actions/setup-node@v5 with: - node-version: 18 + node-version: 24 - name: Install commitlint and convention run: | diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 9aa5f48..770ef3b 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -17,11 +17,11 @@ jobs: pull-requests: write # to be able to comment on released pull requests steps: - name: Checkout - uses: actions/checkout@v3 + uses: actions/checkout@v5 with: fetch-depth: 0 - name: Setup Node.js - uses: actions/setup-node@v3 + uses: actions/setup-node@v5 with: node-version: "lts/*" - name: Install dependencies diff --git a/.github/workflows/test.yml b/.github/workflows/test.yml index 21744eb..98f4278 100644 --- a/.github/workflows/test.yml +++ b/.github/workflows/test.yml @@ -10,8 +10,8 @@ jobs: tests: runs-on: ubuntu-latest steps: - - uses: actions/checkout@v3 - - uses: actions/setup-go@v4 + - uses: actions/checkout@v5 + - uses: actions/setup-go@v5 with: go-version-file: 'go.mod' cache: false diff --git a/go.mod b/go.mod index b9949ea..b20235f 100644 --- a/go.mod +++ b/go.mod @@ -1,6 +1,6 @@ module github.com/Pod-Point/go-queue-worker -go 1.21 +go 1.26 require ( github.com/aws/aws-sdk-go-v2/service/sqs v1.34.5