Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion bun.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions packages/webhook-triggers/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
"test": "bun test"
},
"dependencies": {
"@corbits/error-sink": "workspace:*",
"@corbits/folded-runs": "workspace:*",
"@intx/db": "workspace:*",
"@intx/hub-api": "workspace:*",
Expand Down
22 changes: 13 additions & 9 deletions packages/webhook-triggers/src/launch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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;
Expand Down Expand Up @@ -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 };
Expand Down
69 changes: 67 additions & 2 deletions packages/webhook-triggers/test/launch.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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");

Expand Down Expand Up @@ -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 = {
Expand Down Expand Up @@ -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,
};

Expand All @@ -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<string, unknown>;
},
];
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 = [];
Expand Down
2 changes: 1 addition & 1 deletion packages/workflow-deploy-source/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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:",
Expand Down
18 changes: 9 additions & 9 deletions packages/workflow-deploy-source/src/record-on-deploy.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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(
Expand Down Expand Up @@ -62,27 +60,29 @@ export function withDeploySourceRecording<T extends DeployWorkflowDeployer>(
...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<typeof recordFromDeployParams>,
): Promise<void> {
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 },
});
}
}
56 changes: 54 additions & 2 deletions packages/workflow-deploy-source/test/record-on-deploy.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<string, unknown>;
},
];
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();
Expand Down
Loading