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
34 changes: 28 additions & 6 deletions src/selfhost/monitored-work.ts
Original file line number Diff line number Diff line change
Expand Up @@ -86,12 +86,34 @@ export async function drainOrbRelayWithMonitor(args: {
result: events.length > 0 ? "events" : "empty",
});
for (const ev of events) {
const result = await args.enqueue(
args.env,
ev.deliveryId,
ev.eventName,
ev.rawBody,
);
// #audit-orb-relay-enqueue-isolation: an enqueue can throw uncaught (e.g. a D1/Postgres write failure
// inside recordWebhookEvent, not just the anticipated failures enqueueWebhookByEnv already returns as a
// string result) -- that must not abort the REST of this batch, or every event after the failing one
// is silently never attempted this tick. Isolate per event and treat a throw exactly like the existing
// non-throwing "enqueue_failed" result: don't ack (the relay redelivers it next drain) and keep going.
let result: EnqueueWebhookResult;
try {
result = await args.enqueue(
args.env,
ev.deliveryId,
ev.eventName,
ev.rawBody,
);
} catch (error) {
incr("gittensory_orb_webhook_total", {
event: orbRelayMetricEvent(ev.eventName),
result: "enqueue_failed",
});
console.error(
JSON.stringify({
level: "error",
event: "orb_relay_enqueue_threw",
eventName: ev.eventName,
error: error instanceof Error ? error.message : String(error),
}),
);
continue;
}
incr("gittensory_orb_webhook_total", {
event: orbRelayMetricEvent(ev.eventName),
result,
Expand Down
50 changes: 50 additions & 0 deletions test/unit/selfhost-monitored-work.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,56 @@ describe("self-host monitored recurring work", () => {
expect(metrics).toContain('gittensory_orb_webhook_total{event="check_suite",result="duplicate"} 1');
});

it("REGRESSION (#audit-orb-relay-enqueue-isolation): an enqueue that throws for one event does not abort the rest of the batch", async () => {
const state: OrbRelayDrainState = { pendingAck: [], lastDrainAtMs: null };
const drain = vi.fn().mockResolvedValue([
{ deliveryId: "ok-1", eventName: "pull_request", rawBody: "{}" },
{ deliveryId: "throws-2", eventName: "issues", rawBody: "{}" },
{ deliveryId: "ok-3", eventName: "check_suite", rawBody: "{}" },
]);
const enqueue = vi
.fn()
.mockResolvedValueOnce("queued")
.mockRejectedValueOnce(new Error("D1 write error"))
.mockResolvedValueOnce("queued");
const errors = vi.spyOn(console, "error").mockImplementation(() => undefined);

await drainOrbRelayWithMonitor({
state,
relayEnv: {},
env: {} as Env,
drain,
enqueue,
});

// All 3 events were attempted -- ok-3 was still reached even though throws-2 (the 2nd) rejected.
expect(enqueue).toHaveBeenCalledTimes(3);
expect(enqueue).toHaveBeenNthCalledWith(3, {}, "ok-3", "check_suite", "{}");
// throws-2 is NOT acked (the relay redelivers it next drain), but both successful events are.
expect(state.pendingAck).toEqual(["ok-1", "ok-3"]);
const logged = errors.mock.calls.map((c) => String(c[0])).find((line) => line.includes("orb_relay_enqueue_threw"));
expect(logged).toBeDefined();
expect(JSON.parse(logged!)).toMatchObject({ level: "error", event: "orb_relay_enqueue_threw", eventName: "issues", error: "D1 write error" });
const metrics = await renderMetrics();
expect(metrics).toContain('gittensory_orb_webhook_total{event="pull_request",result="queued"} 1');
expect(metrics).toContain('gittensory_orb_webhook_total{event="issues",result="enqueue_failed"} 1');
expect(metrics).toContain('gittensory_orb_webhook_total{event="check_suite",result="queued"} 1');
errors.mockRestore();
});

it("logs a non-Error enqueue rejection by stringifying it (the false ternary arm)", async () => {
const state: OrbRelayDrainState = { pendingAck: [], lastDrainAtMs: null };
const drain = vi.fn().mockResolvedValue([{ deliveryId: "throws-1", eventName: "pull_request", rawBody: "{}" }]);
const enqueue = vi.fn().mockRejectedValueOnce("not an Error instance");
const errors = vi.spyOn(console, "error").mockImplementation(() => undefined);

await drainOrbRelayWithMonitor({ state, relayEnv: {}, env: {} as Env, drain, enqueue });

const logged = errors.mock.calls.map((c) => String(c[0])).find((line) => line.includes("orb_relay_enqueue_threw"));
expect(JSON.parse(logged!)).toMatchObject({ error: "not an Error instance" });
errors.mockRestore();
});

it("clears previous Orb relay acks and stays quiet when the broker has no events", async () => {
const state: OrbRelayDrainState = { pendingAck: ["previous-delivery"], lastDrainAtMs: null };
const drain = vi.fn().mockResolvedValue([]);
Expand Down
Loading