assessment-model/src/app/utils/sqs.ts
2026-04-18 18:55:19 +00:00

85 lines
2.5 KiB
TypeScript

// utils/sqs.ts
import {
SQSClient,
SendMessageCommand,
GetQueueUrlCommand,
ListQueuesCommand,
SendMessageCommandOutput,
} from "@aws-sdk/client-sqs";
// If you prefer explicit creds via env, keep your current config;
// otherwise, this ctor will use the default credential chain (env vars, shared profile, role, etc.)
const sqsClient = new SQSClient({
region: process.env.SQS_AWS_REGION,
credentials: {
accessKeyId: process.env.SQS_AWS_ACCESS_KEY_ID as string,
secretAccessKey: process.env.SQS_AWS_SECRET_ACCESS_KEY as string,
},
});
const queueUrlCache = new Map<string, string>();
// Export if you want to reuse elsewhere
export async function getQueueUrl(queueName: string): Promise<string> {
const cached = queueUrlCache.get(queueName);
if (cached) return cached;
const resp = await sqsClient.send(
new GetQueueUrlCommand({ QueueName: queueName })
);
if (!resp.QueueUrl)
throw new Error(`Could not resolve SQS URL for queue: ${queueName}`);
queueUrlCache.set(queueName, resp.QueueUrl);
return resp.QueueUrl;
}
type SendOptions = {
queueName?: string; // defaults to env
groupId?: string; // for FIFO queues only
deduplicationId?: string; // for FIFO queues only
delaySeconds?: number; // 0-900
};
/**
* Send a message to SQS. Handles both standard and FIFO queues.
*/
export async function sendToQueue(
messageBody: unknown,
opts: SendOptions = {}
): Promise<SendMessageCommandOutput> {
const queueName =
opts.queueName ?? (process.env.AWS_SQS_QUEUE_NAME as string);
if (!queueName)
throw new Error("Missing AWS_SQS_QUEUE_NAME or sendToQueue opts.queueName");
const queueUrl = await getQueueUrl(queueName);
const params: any = {
QueueUrl: queueUrl,
MessageBody: JSON.stringify(messageBody),
};
// If it's a FIFO queue (ends with .fifo), include group/dedupe if provided
const isFifo = queueUrl.endsWith(".fifo");
if (isFifo) {
params.MessageGroupId = opts.groupId ?? "default-group";
if (opts.deduplicationId)
params.MessageDeduplicationId = opts.deduplicationId;
}
if (typeof opts.delaySeconds === "number") {
params.DelaySeconds = opts.delaySeconds;
}
return sqsClient.send(new SendMessageCommand(params));
}
/**
* List queues in the configured region.
* Optionally filter by name prefix.
*/
export async function listQueues(prefix?: string): Promise<string[]> {
const resp = await sqsClient.send(
new ListQueuesCommand(prefix ? { QueueNamePrefix: prefix } : {})
);
return resp.QueueUrls ?? [];
}