diff --git a/internal/pkg/pipeline/task/sqs/README.md b/internal/pkg/pipeline/task/sqs/README.md index 85c3599..37f291c 100644 --- a/internal/pkg/pipeline/task/sqs/README.md +++ b/internal/pkg/pipeline/task/sqs/README.md @@ -18,12 +18,13 @@ The task automatically determines its mode based on the presence of input/output | `name` | string | - | Task name for identification | | `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 | +| `concurrency` | int | `10` | Number of concurrent `DeleteMessage` calls | | `max_messages` | int | `10` | Maximum number of messages to receive per batch | | `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` | | `message_group_id` | string | - | Message group ID for FIFO queues | +| `delivery` | string | `at-least-once` | Read mode: `at-least-once` deletes after downstream finishes; `at-most-once` deletes on receive. | | `task_concurrency` | int | `1` | Number of competing-consumer workers for this task | | `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 | @@ -73,10 +74,12 @@ tasks: ## Message Acknowledgment When reading from a queue, a message's receipt is deleted only once every downstream task -has finished with the record produced from it. A task that returns an error while holding a -record leaves the receipt alone, so SQS redelivers the message after the visibility timeout -rather than losing it. Delivery is therefore at-least-once: a pipeline may see a message -more than once. +has finished with the record produced from it (`delivery: at-least-once`, the default). A task +that returns an error while holding a record leaves the receipt alone, so SQS redelivers the +message after the visibility timeout rather than losing it. Delivery is therefore at-least-once: +a pipeline may see a message more than once. `delivery: at-most-once` forwards the record +without an ack and deletes on the same `concurrency` pool, so FIFO groups are not blocked by +downstream work; a crash can lose the message. That covers failures, not drops. A task configured to skip a bad record counts it as finished, so the receipt is deleted and the message does not come back — and `jq`'s diff --git a/internal/pkg/pipeline/task/sqs/sqs.go b/internal/pkg/pipeline/task/sqs/sqs.go index 0834306..97b0f08 100644 --- a/internal/pkg/pipeline/task/sqs/sqs.go +++ b/internal/pkg/pipeline/task/sqs/sqs.go @@ -28,6 +28,13 @@ const ( defaultRegion = "us-west-2" ) +type deliveryMode string + +const ( + deliveryAtMostOnce deliveryMode = "at-most-once" + deliveryAtLeastOnce deliveryMode = "at-least-once" +) + var ( awsRegionRegex = regexp.MustCompile(`^[a-z]{2}-[a-z]+-\d+$`) ctx = context.Background() @@ -41,12 +48,13 @@ var ( type sqs struct { task.ServerBase `yaml:",inline" json:",inline"` - 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"` - 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 + 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"` + 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 + Delivery deliveryMode `yaml:"delivery,omitempty" json:"delivery,omitempty"` client *qs.Client tracker *ack.Tracker @@ -70,6 +78,16 @@ func (s *sqs) Init() error { return fmt.Errorf("queue_url is required") } + switch s.Delivery { + case "", deliveryAtLeastOnce: + s.Delivery = deliveryAtLeastOnce + case deliveryAtMostOnce: + default: + return fmt.Errorf("invalid delivery mode %q: must be %q or %q", s.Delivery, deliveryAtMostOnce, deliveryAtLeastOnce) + } + + s.tracker = ack.NewTracker(s.Concurrency) + region := s.extractRegionFromQueueURL() awsConfig, err := config.LoadDefaultConfig(ctx, config.WithRegion(region)) if err != nil { @@ -77,7 +95,6 @@ func (s *sqs) Init() error { } s.client = qs.NewFromConfig(awsConfig) - s.tracker = ack.NewTracker(s.Concurrency) return nil } @@ -170,28 +187,38 @@ func (s *sqs) getMessages(ctx context.Context, output chan<- *record.Record) err for _, m := range receiveMessageOutput.Messages { - // nothing to forward to, so there's no downstream ack to wait - // for: delete the receipt right away. if output == nil { s.deleteMessage(m.MessageId, m.ReceiptHandle) continue } + if s.Delivery == deliveryAtMostOnce { + s.SendData(ctx, []byte(*m.Body), output) + del := ack.New() + del.AddBranch(1) + del.Done() + s.enqueueDelete(m, del) + continue + } + 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, - }) + s.enqueueDelete(m, msgAck) } } } } +func (s *sqs) enqueueDelete(m types.Message, a *ack.Ack) { + s.outstanding.Add(1) + s.tracker.Track(a, &messageAck{ + sqs: s, + messageId: m.MessageId, + receiptHandle: m.ReceiptHandle, + }) +} + // messageAck acknowledges one received message on behalf of ack.Tracker. type messageAck struct { sqs *sqs