diff --git a/apps/coordinator/src/server.ts b/apps/coordinator/src/server.ts index 4cc6e9a7d..82372736e 100644 --- a/apps/coordinator/src/server.ts +++ b/apps/coordinator/src/server.ts @@ -399,7 +399,7 @@ export class TriggerServer { return true; } - const success = await this.#serverRPC.send("COMPLETE_REQUEST", { + const success = await this.#serverRPC.send("RESOLVE_REQUEST", { id: data.id, output: data.output, meta: { diff --git a/packages/internal-bridge/src/schemas/host.ts b/packages/internal-bridge/src/schemas/host.ts index 1b300a16c..f173de5a4 100644 --- a/packages/internal-bridge/src/schemas/host.ts +++ b/packages/internal-bridge/src/schemas/host.ts @@ -18,7 +18,7 @@ export const HostRPCSchema = { }), response: z.void().nullable(), }, - COMPLETE_REQUEST: { + RESOLVE_REQUEST: { request: z.object({ id: z.string(), output: JsonSchema.default({}), diff --git a/packages/internal-platform/src/messages/zodSubscriber.ts b/packages/internal-platform/src/messages/zodSubscriber.ts index d59f37eab..7d7340c6f 100644 --- a/packages/internal-platform/src/messages/zodSubscriber.ts +++ b/packages/internal-platform/src/messages/zodSubscriber.ts @@ -115,7 +115,8 @@ export class ZodSubscriber { this.#logger.error("[ZodSubscriber] Error handling message", e); } - consumer.negativeAcknowledge(msg); + // TODO: Add support for dead letter queue + await consumer.acknowledge(msg); } } diff --git a/packages/trigger-sdk/src/client.ts b/packages/trigger-sdk/src/client.ts index c828920f3..1a567ad16 100644 --- a/packages/trigger-sdk/src/client.ts +++ b/packages/trigger-sdk/src/client.ts @@ -109,7 +109,7 @@ export class TriggerClient { sender: ServerRPCSchema, receiver: HostRPCSchema, handlers: { - COMPLETE_REQUEST: async (data) => { + RESOLVE_REQUEST: async (data) => { const requestCallbacks = this.#responseCompleteCallbacks.get(data.id); if (!requestCallbacks) {