diff --git a/bun.lock b/bun.lock index 61acacdc1..96da4beb8 100644 --- a/bun.lock +++ b/bun.lock @@ -1432,6 +1432,7 @@ "name": "@corbits/webhook-triggers", "version": "0.0.1", "dependencies": { + "@corbits/error-sink": "workspace:*", "@corbits/folded-runs": "workspace:*", "@intx/db": "workspace:*", "@intx/hub-api": "workspace:*", @@ -1477,8 +1478,8 @@ "name": "@corbits/workflow-deploy-source", "version": "0.0.1", "dependencies": { + "@corbits/error-sink": "workspace:*", "@intx/hub-sessions": "workspace:*", - "@intx/log": "0.3.0", "@intx/types": "workspace:*", "arktype": "catalog:", "drizzle-orm": "catalog:", diff --git a/packages/webhook-triggers/package.json b/packages/webhook-triggers/package.json index 03a9bc8ec..d3e8234f7 100644 --- a/packages/webhook-triggers/package.json +++ b/packages/webhook-triggers/package.json @@ -14,6 +14,7 @@ "test": "bun test" }, "dependencies": { + "@corbits/error-sink": "workspace:*", "@corbits/folded-runs": "workspace:*", "@intx/db": "workspace:*", "@intx/hub-api": "workspace:*", diff --git a/packages/webhook-triggers/src/launch.ts b/packages/webhook-triggers/src/launch.ts index 17942230d..ae847b70a 100644 --- a/packages/webhook-triggers/src/launch.ts +++ b/packages/webhook-triggers/src/launch.ts @@ -16,7 +16,9 @@ // `createWebhookIngressRoutes` — that would both hide the run (no // `store.recordFired` call) and, if the sender's webhook client retries // the same delivery on a 5xx, mint a second, duplicate run for one -// event. On exhausted retries this only logs, naming the run. +// event. On exhausted retries this only reports the failure through +// `@corbits/error-sink`, naming the run. +import { reportError } from "@corbits/error-sink"; import { and, eq } from "drizzle-orm"; import { domainOf, @@ -32,14 +34,11 @@ import { import type { DB } from "@intx/db"; import { tenant as tenantTable, workflowDefinition } from "@intx/db/schema"; import { generateId } from "@intx/hub-common"; -import { getLogger } from "@intx/log"; import { formatRunAddress } from "@intx/types"; import { renderInputTemplate } from "./mapping"; import type { WebhookTriggerRow } from "./schema"; -const log = getLogger(["webhook-triggers", "launch"]); - export type LaunchWebhookTriggerDeps = FoldedRunsDeps & { db: DB["db"]; cryptoProviderCache: CryptoProviderCache; @@ -161,11 +160,16 @@ export async function launchWebhookTrigger( cryptoProvider, }); if (!result.ok) { - const reason = - result.error instanceof Error - ? result.error.message - : String(result.error); - log.error`run ${instanceId} launched from webhook trigger "${trigger.id}" but its input failed to deliver after ${result.attempts} attempts: ${reason}`; + reportError(result.error, { + operation: "webhookTriggers.launch.deliverInput", + tenantId: trigger.tenantId, + agentId: triggerAddress, + extra: { + instanceId, + triggerId: trigger.id, + attempts: result.attempts, + }, + }); } return { instanceId, triggerAddress }; diff --git a/packages/webhook-triggers/test/launch.test.ts b/packages/webhook-triggers/test/launch.test.ts index 9a7194ebc..25c9393f6 100644 --- a/packages/webhook-triggers/test/launch.test.ts +++ b/packages/webhook-triggers/test/launch.test.ts @@ -5,7 +5,9 @@ // past this function (or `createWebhookIngressRoutes` would reject an // already-launched delivery, and a retried webhook client would then // mint a duplicate run for the same event). -import { describe, expect, mock, test } from "bun:test"; +import { afterEach, beforeEach, describe, expect, mock, test } from "bun:test"; + +let reportErrorCalls: unknown[] = []; const actualFoldedRuns = await import("@corbits/folded-runs"); @@ -35,6 +37,20 @@ mock.module("@corbits/folded-runs", () => ({ }, })); +beforeEach(async () => { + reportErrorCalls = []; + await mock.module("@corbits/error-sink", () => ({ + reportError: (...args: unknown[]) => { + reportErrorCalls.push(args); + return "ref_test"; + }, + })); +}); + +afterEach(() => { + mock.restore(); +}); + const { launchWebhookTrigger } = await import("../src/launch"); const DEFINITION_ROW = { @@ -95,9 +111,11 @@ describe("launchWebhookTrigger", () => { test("still returns the launched run when input delivery fails after every retry", async () => { launchFoldedRunCalls = []; sendFoldedMailWithRetryCalls = []; + reportErrorCalls = []; + const deliveryError = new Error("sidecar unreachable"); sendFoldedMailWithRetryResult = { ok: false, - error: new Error("sidecar unreachable"), + error: deliveryError, attempts: 3, }; @@ -111,6 +129,53 @@ describe("launchWebhookTrigger", () => { expect(sendFoldedMailWithRetryCalls).toHaveLength(1); }); + test("reports the exhausted delivery failure with the run's context", async () => { + launchFoldedRunCalls = []; + sendFoldedMailWithRetryCalls = []; + reportErrorCalls = []; + const deliveryError = new Error("sidecar unreachable"); + sendFoldedMailWithRetryResult = { + ok: false, + error: deliveryError, + attempts: 3, + }; + + const result = await launchWebhookTrigger(baseDeps(), TRIGGER, { + status: "ok", + }); + + expect(reportErrorCalls).toHaveLength(1); + const [cause, context] = reportErrorCalls[0] as [ + unknown, + { + operation: string; + tenantId: string; + agentId: string; + extra: Record; + }, + ]; + expect(cause).toBe(deliveryError); + expect(context.operation).toBe("webhookTriggers.launch.deliverInput"); + expect(context.tenantId).toBe(TRIGGER.tenantId); + expect(context.agentId).toBe(result.triggerAddress); + expect(context.extra).toEqual({ + instanceId: result.instanceId, + triggerId: TRIGGER.id, + attempts: 3, + }); + }); + + test("does not report anything when delivery succeeds", async () => { + launchFoldedRunCalls = []; + sendFoldedMailWithRetryCalls = []; + reportErrorCalls = []; + sendFoldedMailWithRetryResult = { ok: true, mail: { id: "m_1" } }; + + await launchWebhookTrigger(baseDeps(), TRIGGER, { status: "ok" }); + + expect(reportErrorCalls).toHaveLength(0); + }); + test("returns the launched run normally when delivery succeeds", async () => { launchFoldedRunCalls = []; sendFoldedMailWithRetryCalls = []; diff --git a/packages/workflow-deploy-source/package.json b/packages/workflow-deploy-source/package.json index 3c6711787..aa3eec059 100644 --- a/packages/workflow-deploy-source/package.json +++ b/packages/workflow-deploy-source/package.json @@ -14,8 +14,8 @@ "test": "bun test" }, "dependencies": { + "@corbits/error-sink": "workspace:*", "@intx/hub-sessions": "workspace:*", - "@intx/log": "0.3.0", "@intx/types": "workspace:*", "arktype": "catalog:", "drizzle-orm": "catalog:", diff --git a/packages/workflow-deploy-source/src/record-on-deploy.ts b/packages/workflow-deploy-source/src/record-on-deploy.ts index b31775e3d..4cdd35218 100644 --- a/packages/workflow-deploy-source/src/record-on-deploy.ts +++ b/packages/workflow-deploy-source/src/record-on-deploy.ts @@ -13,7 +13,7 @@ // (vendor/intx/db/src/schema/workflow-run-launch-spec.ts), written from // vendored `workflow-allocation-service.ts`, which this package does not // touch. -import { getLogger } from "@intx/log"; +import { reportError } from "@corbits/error-sink"; import type { AdoptingWorkflowDeployer, DeployAdoptedWorkflowFromSourceParams, @@ -23,8 +23,6 @@ import type { import type { WorkflowDeploySourceStore } from "./store"; -const logger = getLogger(["workflow-deploy-source", "record-on-deploy"]); - export type DeployWorkflowDeployer = SessionService & AdoptingWorkflowDeployer; function recordFromDeployParams( @@ -62,27 +60,29 @@ export function withDeploySourceRecording( ...sessionService, async deployWorkflowFromSource(params) { const result = await sessionService.deployWorkflowFromSource(params); - await recordOrLog(store, recordFromDeployParams(params)); + await recordOrReport(store, recordFromDeployParams(params)); return result; }, async deployAdoptedWorkflowFromSource(params) { const result = await sessionService.deployAdoptedWorkflowFromSource(params); - await recordOrLog(store, recordFromDeployParams(params)); + await recordOrReport(store, recordFromDeployParams(params)); return result; }, }; } -async function recordOrLog( +async function recordOrReport( store: WorkflowDeploySourceStore, entry: ReturnType, ): Promise { try { await store.record(entry); } catch (cause) { - logger.error`failed to record deploy source for anchor run ${entry.anchorRunId}: ${ - cause instanceof Error ? cause.message : String(cause) - }`; + reportError(cause, { + operation: "workflowDeploySource.record", + tenantId: entry.tenantId, + extra: { anchorRunId: entry.anchorRunId }, + }); } } diff --git a/packages/workflow-deploy-source/test/record-on-deploy.test.ts b/packages/workflow-deploy-source/test/record-on-deploy.test.ts index 5d8d54159..1f176ab2b 100644 --- a/packages/workflow-deploy-source/test/record-on-deploy.test.ts +++ b/packages/workflow-deploy-source/test/record-on-deploy.test.ts @@ -2,14 +2,28 @@ // successful deploy, forwards the underlying result untouched, and never // lets a recording failure fail the deploy call itself -- no DB, no // SessionService, both are hand-rolled minimal doubles. -import { describe, expect, test } from "bun:test"; +import { afterEach, beforeEach, describe, expect, mock, test } from "bun:test"; import type { AdoptingWorkflowDeployer, DeployWorkflowDefinitionResult, SessionService, } from "@intx/hub-sessions"; -import { withDeploySourceRecording } from "../src/record-on-deploy"; +let reportErrorCalls: unknown[] = []; +beforeEach(async () => { + reportErrorCalls = []; + await mock.module("@corbits/error-sink", () => ({ + reportError: (...args: unknown[]) => { + reportErrorCalls.push(args); + return "ref_test"; + }, + })); +}); +afterEach(() => { + mock.restore(); +}); + +const { withDeploySourceRecording } = await import("../src/record-on-deploy"); import type { WorkflowDeploySourceRecord, WorkflowDeploySourceStore, @@ -119,6 +133,44 @@ describe("withDeploySourceRecording", () => { expect(result.anchorRunId).toBe("run_1"); }); + test("a recording failure is reported with the run's context", async () => { + reportErrorCalls = []; + const recordError = new Error("db unreachable"); + const failingStore: WorkflowDeploySourceStore = { + record: async () => { + throw recordError; + }, + get: async () => null, + }; + const wrapped = withDeploySourceRecording(fakeDeployer(), failingStore); + + await wrapped.deployWorkflowFromSource(baseParams); + + expect(reportErrorCalls).toHaveLength(1); + const [cause, context] = reportErrorCalls[0] as [ + unknown, + { + operation: string; + tenantId: string; + extra: Record; + }, + ]; + expect(cause).toBe(recordError); + expect(context.operation).toBe("workflowDeploySource.record"); + expect(context.tenantId).toBe("tenant_1"); + expect(context.extra).toEqual({ anchorRunId: "run_1" }); + }); + + test("does not report anything when recording succeeds", async () => { + reportErrorCalls = []; + const store = recordingStore(); + const wrapped = withDeploySourceRecording(fakeDeployer(), store); + + await wrapped.deployWorkflowFromSource(baseParams); + + expect(reportErrorCalls).toHaveLength(0); + }); + test("other SessionService methods pass through unwrapped", async () => { const store = recordingStore(); const underlying = fakeDeployer();