Skip to content
File

Blob: src/worker/dispatch/queue/consumer.ts

typescript63 lines
1import { RunQueueMessage } from "@/worker/contracts";
2 
3import { logger } from "@/worker/dispatch/shared/run-execution-context";
4import { executeDispatchedRun } from "./execute-dispatched-run";
5 
6const processQueueMessage = async (message: Message<unknown>, env: Env): Promise<void> => {
7 const payload = RunQueueMessage.assertDecode(message.body);
8 
9 const result = await executeDispatchedRun(env, {
10 projectId: payload.projectId,
11 runId: payload.runId,
12 });
13 
14 switch (result.kind) {
15 case "project_missing":
16 logger.warn("stale_queue_delivery", {
17 queueMessageId: message.id,
18 projectId: payload.projectId,
19 runId: payload.runId,
20 reason: "project_missing",
21 });
22 break;
23 
24 case "recovered":
25 logger.info("recovered_active_queue_delivery", {
26 queueMessageId: message.id,
27 projectId: payload.projectId,
28 runId: payload.runId,
29 });
30 break;
31 
32 case "stale":
33 logger.info("stale_queue_delivery", {
34 queueMessageId: message.id,
35 projectId: payload.projectId,
36 runId: payload.runId,
37 reason: result.reason,
38 });
39 break;
40 
41 case "executed":
42 break;
43 }
44 
45 message.ack();
46};
47 
48export const handleQueueBatch = async (batch: MessageBatch<RunQueueMessage>, env: Env): Promise<void> => {
49 for (const message of batch.messages) {
50 logger.info("queue_message_received", { queueMessageId: message.id });
51 
52 try {
53 await processQueueMessage(message, env);
54 } catch (error) {
55 logger.error("queue_message_failed", {
56 queueMessageId: message.id,
57 error: error instanceof Error ? error.message : String(error),
58 });
59 message.retry();
60 }
61 }
62};