import { SQSClient, SendMessageCommand, ReceiveMessageCommand, DeleteMessageCommand, ChangeMessageVisibilityCommand, type SendMessageCommandInput, type ReceiveMessageCommandInput, type DeleteMessageCommandInput, type ChangeMessageVisibilityCommandInput, } from "@aws-sdk/client-sqs"; import { Queue } from "../core/queue.ts"; import type { JobStatus, JobMeta, QueueMessage, BaseJobOptions, WithDelay, } from "../interfaces/job.ts"; import type { QueueOptions } from "../interfaces/plugin.ts"; // Driver-specific job request interface export interface SqsJobRequest extends BaseJobOptions, WithDelay { /** Job payload */ payload: TPayload; // SQS supports delays (0-900 seconds max) but not priority ordering } interface SqsClient { send: SQSClient["send"]; } export class SqsQueue> extends Queue< TJobMap, SqsJobRequest > { #onFailure: "delete" | "leaveInQueue" constructor( private client: SqsClient, private queueUrl: string, options: QueueOptions & { onFailure: "delete" | "leaveInQueue" } ) { super(options); // SQS supports long polling via WaitTimeSeconds this.supportsLongPolling = true; this.#onFailure = options.onFailure; } protected async pushMessage(payload: unknown, meta: JobMeta): Promise { const messageAttributes: Record< string, { StringValue: string; DataType: string } > = {}; messageAttributes.name = { StringValue: meta.name, DataType: "String", }; if (meta.ttr) { messageAttributes.ttr = { StringValue: meta.ttr.toString(), DataType: "Number", }; } if (meta.priority) { messageAttributes.priority = { StringValue: meta.priority.toString(), DataType: "Number", }; } const command = new SendMessageCommand({ QueueUrl: this.queueUrl, MessageBody: JSON.stringify(payload), DelaySeconds: meta.delaySeconds || 0, MessageAttributes: messageAttributes, }); const result = await this.client.send(command); if (!result.MessageId) { throw new Error("Failed to send message - no MessageId returned"); } return result.MessageId; } protected async reserve(timeout: number): Promise { const command = new ReceiveMessageCommand({ QueueUrl: this.queueUrl, MaxNumberOfMessages: 1, WaitTimeSeconds: timeout, MessageAttributeNames: ["All"], }); const result = await this.client.send(command); if (!result.Messages || result.Messages.length === 0) { return null; } const message = result.Messages[0]; if ( !message || !message.Body || !message.MessageId || !message.ReceiptHandle ) { return null; } const payload = JSON.parse(message.Body); const meta: JobMeta = { name: message.MessageAttributes?.name?.StringValue || "", ttr: message.MessageAttributes?.ttr?.StringValue ? parseInt(message.MessageAttributes.ttr.StringValue) : undefined, priority: message.MessageAttributes?.priority?.StringValue ? parseInt(message.MessageAttributes.priority.StringValue) : undefined, receiptHandle: message.ReceiptHandle, }; if (meta.ttr) { const visibilityCommand = new ChangeMessageVisibilityCommand({ QueueUrl: this.queueUrl, ReceiptHandle: message.ReceiptHandle, VisibilityTimeout: meta.ttr, }); await this.client.send(visibilityCommand); } return { id: message.MessageId, payload, meta, }; } protected async completeJob(message: QueueMessage): Promise { if (!message.meta.receiptHandle) { throw new Error( "Cannot complete SQS message: receiptHandle is missing from metadata" ); } const deleteCommand = new DeleteMessageCommand({ QueueUrl: this.queueUrl, ReceiptHandle: message.meta.receiptHandle, }); await this.client.send(deleteCommand); } protected async failJob( message: QueueMessage, error: unknown ): Promise { if (this.#onFailure === "leaveInQueue") { if (!message.meta.receiptHandle) { throw new Error( "Cannot fail SQS message: receiptHandle is missing from metadata" ); } // For SQS, we delete failed messages too since SQS doesn't track failure states // Future enhancement could send to a dead letter queue const deleteCommand = new DeleteMessageCommand({ QueueUrl: this.queueUrl, ReceiptHandle: message.meta.receiptHandle, }); await this.client.send(deleteCommand); } } async status(id: string): Promise { throw new Error("SQS does not support status"); } }