From d67e88734594ee9a8e07e39449a6f4419525ce18 Mon Sep 17 00:00:00 2001 From: Divyanshu Tiwari Date: Tue, 15 Sep 2026 19:40:37 +0530 Subject: [PATCH 1/3] feat(sqs): stop read mode after max_records Match Kafka: cap forwarded messages per worker and size the next ReceiveMessage to remaining so leftover receipts are not held. --- internal/pkg/pipeline/task/sqs/README.md | 7 ++-- internal/pkg/pipeline/task/sqs/sqs.go | 44 +++++++++++++++++------- 2 files changed, 36 insertions(+), 15 deletions(-) diff --git a/internal/pkg/pipeline/task/sqs/README.md b/internal/pkg/pipeline/task/sqs/README.md index 85c3599..0c0b063 100644 --- a/internal/pkg/pipeline/task/sqs/README.md +++ b/internal/pkg/pipeline/task/sqs/README.md @@ -19,7 +19,8 @@ The task automatically determines its mode based on the presence of input/output | `type` | string | `sqs` | Must be "sqs" | | `queue_url` | string | - | SQS queue URL (required) | | `concurrency` | int | `10` | Number of concurrent workers that acknowledge (delete) fully-processed messages | -| `max_messages` | int | `10` | Maximum number of messages to receive per batch | +| `max_messages` | int | `10` | Per poll: max messages per `ReceiveMessage` (AWS cap 10) | +| `max_records` | int | `0` (unlimited) | Read-mode cap on records forwarded downstream (per worker). `0` = unlimited. | | `wait_time_seconds` | int | `10` | Long polling wait time in seconds | | `exit_on_empty` | bool | `false` | Exit when a receive returns no messages. FIFO: empty poll is not drain while this task holds receipts or the queue still has visible, in-flight, or delayed messages. | | `end_after` | duration | - | Stop polling after this much time (read mode); e.g. `5m` | @@ -28,8 +29,8 @@ The task automatically determines its mode based on the presence of input/output | `context` | map | - | JQ expressions whose results are stored on each record for downstream tasks | | `fail_on_error` | bool | `false` | Whether to stop the pipeline if this task encounters an error | -In read mode the task polls until the queue drains (`exit_on_empty`) or `end_after` elapses; -with neither set it polls indefinitely. +In read mode the task polls until the queue drains (`exit_on_empty`), `end_after` elapses, or `max_records` is reached; +with none of those set it polls indefinitely. ## Example Configurations diff --git a/internal/pkg/pipeline/task/sqs/sqs.go b/internal/pkg/pipeline/task/sqs/sqs.go index 0834306..90a1d0f 100644 --- a/internal/pkg/pipeline/task/sqs/sqs.go +++ b/internal/pkg/pipeline/task/sqs/sqs.go @@ -44,11 +44,13 @@ type sqs struct { QueueURL string `yaml:"queue_url" json:"queue_url" validate:"required"` Concurrency int `yaml:"concurrency,omitempty" json:"concurrency,omitempty"` MaxMessages int32 `yaml:"max_messages,omitempty" json:"max_messages,omitempty"` + MaxRecords int `yaml:"max_records,omitempty" json:"max_records,omitempty" validate:"omitempty,gte=0"` WaitTimeSeconds int `yaml:"wait_time_seconds,omitempty" json:"wait_time_seconds,omitempty"` ExitOnEmpty bool `yaml:"exit_on_empty,omitempty" json:"exit_on_empty,omitempty"` MessageGroupId string `yaml:"message_group_id,omitempty" json:"message_group_id,omitempty"` // used for FIFO queues client *qs.Client + receive func(context.Context, *qs.ReceiveMessageInput, ...func(*qs.Options)) (*qs.ReceiveMessageOutput, error) // test fake; nil uses client tracker *ack.Tracker outstanding atomic.Int32 } @@ -136,6 +138,12 @@ func (s *sqs) getMessages(ctx context.Context, output chan<- *record.Record) err defer cancel() } + recv := s.receive + if recv == nil { + recv = s.client.ReceiveMessage + } + + recordsRead := 0 for { select { case <-ctx.Done(): @@ -143,9 +151,16 @@ func (s *sqs) getMessages(ctx context.Context, output chan<- *record.Record) err return nil default: - receiveMessageOutput, err := s.client.ReceiveMessage(ctx, &qs.ReceiveMessageInput{ + batch := s.MaxMessages + if s.MaxRecords > 0 { + if remaining := s.MaxRecords - recordsRead; remaining < int(batch) { + batch = int32(remaining) + } + } + + receiveMessageOutput, err := recv(ctx, &qs.ReceiveMessageInput{ QueueUrl: &s.QueueURL, - MaxNumberOfMessages: s.MaxMessages, + MaxNumberOfMessages: batch, WaitTimeSeconds: int32(s.WaitTimeSeconds), }) @@ -174,18 +189,23 @@ func (s *sqs) getMessages(ctx context.Context, output chan<- *record.Record) err // for: delete the receipt right away. if output == nil { s.deleteMessage(m.MessageId, m.ReceiptHandle) - continue + } else { + msgAck := ack.New() + s.SendData(ack.WithContext(ctx, msgAck), []byte(*m.Body), output) + + s.outstanding.Add(1) + s.tracker.Track(msgAck, &messageAck{ + sqs: s, + messageId: m.MessageId, + receiptHandle: m.ReceiptHandle, + }) } - msgAck := ack.New() - s.SendData(ack.WithContext(ctx, msgAck), []byte(*m.Body), output) - - s.outstanding.Add(1) - s.tracker.Track(msgAck, &messageAck{ - sqs: s, - messageId: m.MessageId, - receiptHandle: m.ReceiptHandle, - }) + recordsRead++ + if s.MaxRecords > 0 && recordsRead >= s.MaxRecords { + fmt.Printf("SQS max_records (%d) reached for queue %s, stopping reader\n", s.MaxRecords, s.QueueURL) + return nil + } } } } From 8a5c7f9cc2ffa131195cfa3be93533c33a0e727f Mon Sep 17 00:00:00 2001 From: Divyanshu Tiwari Date: Tue, 15 Sep 2026 19:49:03 +0530 Subject: [PATCH 2/3] fix(sqs): keep max_messages as the ReceiveMessage batch size max_records only stops forwarding, matching Kafka. Do not shrink the poll. --- internal/pkg/pipeline/task/sqs/sqs.go | 9 +-------- 1 file changed, 1 insertion(+), 8 deletions(-) diff --git a/internal/pkg/pipeline/task/sqs/sqs.go b/internal/pkg/pipeline/task/sqs/sqs.go index 90a1d0f..8354436 100644 --- a/internal/pkg/pipeline/task/sqs/sqs.go +++ b/internal/pkg/pipeline/task/sqs/sqs.go @@ -151,16 +151,9 @@ func (s *sqs) getMessages(ctx context.Context, output chan<- *record.Record) err return nil default: - batch := s.MaxMessages - if s.MaxRecords > 0 { - if remaining := s.MaxRecords - recordsRead; remaining < int(batch) { - batch = int32(remaining) - } - } - receiveMessageOutput, err := recv(ctx, &qs.ReceiveMessageInput{ QueueUrl: &s.QueueURL, - MaxNumberOfMessages: batch, + MaxNumberOfMessages: s.MaxMessages, WaitTimeSeconds: int32(s.WaitTimeSeconds), }) From cc0359f77b4463ab0d8779064634b92b1dd33fcf Mon Sep 17 00:00:00 2001 From: Divyanshu Tiwari Date: Tue, 15 Sep 2026 19:55:15 +0530 Subject: [PATCH 3/3] refactor(sqs): drop unused ReceiveMessage test hook --- internal/pkg/pipeline/task/sqs/sqs.go | 8 +------- 1 file changed, 1 insertion(+), 7 deletions(-) diff --git a/internal/pkg/pipeline/task/sqs/sqs.go b/internal/pkg/pipeline/task/sqs/sqs.go index 8354436..7b12b6d 100644 --- a/internal/pkg/pipeline/task/sqs/sqs.go +++ b/internal/pkg/pipeline/task/sqs/sqs.go @@ -50,7 +50,6 @@ type sqs struct { MessageGroupId string `yaml:"message_group_id,omitempty" json:"message_group_id,omitempty"` // used for FIFO queues client *qs.Client - receive func(context.Context, *qs.ReceiveMessageInput, ...func(*qs.Options)) (*qs.ReceiveMessageOutput, error) // test fake; nil uses client tracker *ack.Tracker outstanding atomic.Int32 } @@ -138,11 +137,6 @@ func (s *sqs) getMessages(ctx context.Context, output chan<- *record.Record) err defer cancel() } - recv := s.receive - if recv == nil { - recv = s.client.ReceiveMessage - } - recordsRead := 0 for { select { @@ -151,7 +145,7 @@ func (s *sqs) getMessages(ctx context.Context, output chan<- *record.Record) err return nil default: - receiveMessageOutput, err := recv(ctx, &qs.ReceiveMessageInput{ + receiveMessageOutput, err := s.client.ReceiveMessage(ctx, &qs.ReceiveMessageInput{ QueueUrl: &s.QueueURL, MaxNumberOfMessages: s.MaxMessages, WaitTimeSeconds: int32(s.WaitTimeSeconds),