From 72877c8cf909519dd6e1d593ddcf4ce6bd00721a Mon Sep 17 00:00:00 2001 From: ischindl Date: Wed, 19 Aug 2026 07:11:47 +0200 Subject: [PATCH] fix(RUFU-074): idle backoff + jitter for the task-deleted outbox consumer (#3471) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit **Problem:** Each dashboard/engine project consumer polled `task_deleted` outbox on a fixed 5s setInterval, so ~44 per-project consumers thundered together on the same cadence — an idle DB query storm and CPU hot-spot even when projects were paused/idle. **Fix:** The outbox consumer reschedules itself from each poll outcome: an idle poll (zero events) grows the next delay by `TASK_DELETED_OUTBOX_BACKOFF_STEP_MS` toward `MAX_POLL_MS`, with ±20% jitter so the consumers de-synchronize; a poll that delivered events resets to the fast base. A paused/idle project drains its outbox and backoff alone drops the DB load. **Includes:** regression test (bounded jitter + idle growth), performance changeset, solution doc, deploy handoff script. ## Summary by CodeRabbit - **Performance** - Reduced unnecessary idle polling by gradually increasing the polling interval, up to 60 seconds, with bounded timing variation. - Restored the faster 5-second polling cadence when new events, waits, or transient errors occur. - Preserved event ordering, delivery guarantees, acknowledgements, and independent behavior across concurrent consumers. - **Documentation** - Added guidance on polling behavior, deployment verification, and monitoring targets. - **Tests** - Added coverage for backoff growth, jitter limits, event bursts, concurrent consumers, error handling, retries, and clean shutdown. --- .changeset/rufu-074-performance.md | 7 + .../task-deleted-outbox-idle-backoff.md | 121 ++++++ ...sk-deleted-outbox-consumer-backoff.test.ts | 410 ++++++++++++++++++ .../task-deleted-outbox-consumer.ts | 137 +++++- 4 files changed, 660 insertions(+), 15 deletions(-) create mode 100644 .changeset/rufu-074-performance.md create mode 100644 docs/solutions/performance/task-deleted-outbox-idle-backoff.md create mode 100644 packages/core/src/__tests__/task-deleted-outbox-consumer-backoff.test.ts diff --git a/.changeset/rufu-074-performance.md b/.changeset/rufu-074-performance.md new file mode 100644 index 0000000000..794de35618 --- /dev/null +++ b/.changeset/rufu-074-performance.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Back off idle task-lifecycle outbox consumers to a 60s cadence so paused/idle projects stop the 98% CPU poll storm. +category: performance +dev: TaskDeletedOutboxConsumer now self-reschedules with a tri-state poll outcome (active/idle/waiting) and ±20% jitter: only a genuinely idle poll (empty outbox) grows the next delay by 10s per idle poll toward a 60s cap; a poll that delivers events ("active") or a non-idle wait ("waiting" — retry-backoff window, lease contention, fencing, poll errors, shutdown races) resets to the fast 5s base, so transient failures recover at 5s cadence instead of an error streak masquerading as an idle streak. This targets a drop in task_lifecycle_consumer_cursors idx_scan from ~26/s toward <5/s and CPU from ~98% toward <50% when projects are paused/idle (the ~44 per-project dashboard+engine consumers no longer thunder on a fixed 5s interval), while cursor fencing, lease advance, per-event ordering, and at-least-once delivery are unchanged — backoff only changes when poll() runs, never the poll/dispatch/ack logic. A new event mid-backoff resets the cadence to 5s, bounding delivery latency. \ No newline at end of file diff --git a/docs/solutions/performance/task-deleted-outbox-idle-backoff.md b/docs/solutions/performance/task-deleted-outbox-idle-backoff.md new file mode 100644 index 0000000000..76302fd28d --- /dev/null +++ b/docs/solutions/performance/task-deleted-outbox-idle-backoff.md @@ -0,0 +1,121 @@ +--- +category: performance +module: packages/core/src/task-store/task-deleted-outbox-consumer.ts +date: 2026-08-13 +problem_type: performance +severity: high +applies_when: + - "Seeing a setInterval(5s) outbox poll storm peg CPU across many per-project consumers" + - "A per-project pause does NOT stop a task-store-level poller because the poller is bound to the task store, not the engine pause" + - "task_lifecycle_consumer_cursors idx_scan growing nonstop (~26/s) on idle/paused projects with a zero cursor delta" +component: task-store +tags: + - performance + - poll-storm + - idle-backoff + - task-lifecycle-consumer + - outbox + - fnxc-crossprocessdeleteobservation + - fnxc-tasklifecycleconsumeridlebackoff +related_components: + - task_store + - lifecycle_outbox + - dashboard + - engine +--- + +# Task-lifecycle outbox idle backoff (the RUFU-074 poll storm) + +## Symptom + +Production CPU stayed at **70–98%** and the health API took **0.77–2.0s** even after every project was +paused. `pg_stat_user_tables` showed `project.task_lifecycle_consumer_cursors.idx_scan` growing +**~26/s nonstop** with a zero `last_acked` delta while paused. A `cpuprofile` showed +`onStreamRead@?:166` (the stream/DB polling) dominant plus ~76% in the error/promise machinery. + +Root cause (measured): **44 lifecycle consumers** run nonstop in the single node process — 22× +`dashboard` + ~20× `engine` (verified in `project.task_lifecycle_consumer_registrations`). Each +registers via `createTaskStoreForBackend` with a **fixed 5s `setInterval`** poll. These pollers live +at the **task-store level** (per InProcessRuntime per project + dashboard store cache), NOT the +per-project engine pause — so **project pause never stops them**. A paused project stops writing +lifecycle events, but the consumers kept polling the same empty outbox every 5s, re-reading the cursor +and re-running Drizzle SQL string compilation forever. + +## The chosen fix: idle backoff + jitter inside the outbox consumer (Option 1) + +We deliberately chose **idle backoff with jitter inside `TaskDeletedOutboxConsumer`** (the operator's +Option 1) over the alternatives: + +- **Option 2** (evict/stop dashboard project stores for paused projects) — risks breaking live + SSE/real-time cross-process observation and adds coupling to the engine pause system. +- **Option 3** (deregister the engine consumer on project pause) — same coupling concern. +- **Option 4** (gate dashboard store creation on `project.status active` + only the open board + project) — changes store lifecycle semantics for no benefit. + +A paused/idle project stops writing lifecycle events, so its outbox naturally goes empty and **idle +backoff alone** yields the required CPU/idx_scan drop without touching delivery semantics. The fix +lives entirely at task-store level, so it covers both the `dashboard` and `engine` consumers for every +project automatically. + +### The rescheduling loop + +The consumer no longer uses a fixed `setInterval(5s)`. It now **self-reschedules with a `setTimeout`**, +feeding each poll's tri-state outcome (`active` / `idle` / `waiting`) back into the next delay: + +- **Idle poll** (the outbox was genuinely empty): `idlePollsSinceEvent` increments, and the next + delay grows by `TASK_DELETED_OUTBOX_BACKOFF_STEP_MS` (10s) per idle poll, toward + `TASK_DELETED_OUTBOX_MAX_POLL_MS` (60s): + `5s → 15s → 25s → ... → 60s` (capped). This is the ONLY outcome that extends the backoff. +- **Active poll** (delivered ≥1 event): `idlePollsSinceEvent = 0`, next delay resets to the fast + `TASK_DELETED_OUTBOX_POLL_MS` (5s) base — so a new event mid-backoff is delivered on the next poll after the cadence resets to 5s (bounded by that next poll; not immediate, but not stuck at the 60s cap). +- **Waiting poll** (retry-backoff window `cursor.retryBackoffUntil`, lease contention, fencing, + poll errors mapped by `pollSafely`, shutdown races): `idlePollsSinceEvent = 0` as well. Review + feedback (2026-08-18): retry and error outcomes must not be classified as idle polls — they do + not mean the outbox is empty, so they must not extend the backoff. Resetting to the fast base + keeps a transient failure on a 5s recovery cadence instead of an error streak growing the delay + to the 60s cap and hiding recovery behind it. +- **Jitter:** `applyPollJitter` applies ±20% multiplicative jitter so the ~44 consumers de-synchronize + and don't thunder on the same cadence; the delay never falls below the fast 5s base. + +### Contract preserved + +The `FNXC:CrossProcessDeleteObservation` contract is preserved exactly. Backoff only changes **when** +`poll()` runs, never the `poll`/`dispatch`/`ack` logic: + +- Cursor fencing via `advanceCursorWithFence` / fenced `acknowledgeTaskLifecycleEvent(fencingToken)`. +- Lease advance via `renewTaskLifecycleLease`. +- Per-event ordering (in-order sequence advance). +- At-least-once delivery in the crash window (dispatch before durable receipt/ack; `stop()` disarms the + rescheduled timer so not orphaned poll survives shutdown). + +## Verification + +Regression test `packages/core/src/__tests__/task-deleted-outbox-consumer-backoff.test.ts` (in-memory +fakes + fake timers, no real DB, no waits) proves: + +1. An idle `dashboard` consumer polls far below the fixed-5s rate and reaches the 60s cap, with + bounded jitter. +2. An active `engine` consumer keeps the fast 5s cadence while delivering events (fenced acks). +3. An idle engine consumer backs off **independently** from a concurrent active dashboard consumer. +4. A burst arriving **mid-backoff** is delivered in order with fenced acks (`fencingToken` honored, + sequence `1n,2n,3n`) and resets the cadence to the fast 5s base. +5. `stop()` disarms the timer (no orphaned post-stop poll); a manual post-stop `poll()` reports + `"waiting"` and reads nothing. +6. The deterministic growth curve `5s→15s→...→60s` holds and never falls below base. +7. A transient outbox-reader **error streak** is classified `"waiting"`, not `"idle"`: the cadence + stays on the fast 5s base through the failures and the next poll after recovery delivers the + pending event. +8. A cursor inside its **per-event retry-backoff window** takes the `"waiting"` early return before + reading the outbox and stays on the fast 5s base (the outbox read is skipped while the window + holds). + +## Symptom-verification (operator deploy) + +The production :4040 daemon is the Fusion host this agent runs inside; restarting it crosses the +shutdown boundary and the embedded PG global dir is privilege-fenced from agent sessions. Use the +handoff script `scripts/deploy-rufu-074.mjs` (operator-run) to measure before/after: + +- **Before baseline** (old build): `task_lifecycle_consumer_cursors.idx_scan` growth, daemon CPU, + health latency. +- **After** (new build): same three measures. +- **Targets:** idx_scan < 5/s, CPU < 50%, health < 0.5s. \ No newline at end of file diff --git a/packages/core/src/__tests__/task-deleted-outbox-consumer-backoff.test.ts b/packages/core/src/__tests__/task-deleted-outbox-consumer-backoff.test.ts new file mode 100644 index 0000000000..7c0fe37d3d --- /dev/null +++ b/packages/core/src/__tests__/task-deleted-outbox-consumer-backoff.test.ts @@ -0,0 +1,410 @@ +/* +FNXC:TaskLifecycleConsumerIdleBackoff 2026-08-13-06:41 (RUFU-074): +In-memory regression test (no real DB, fake timers, vi.mock fakes per docs/testing.md) proving that +a TaskDeletedOutboxConsumer with an idle outbox backs off its poll interval toward the 60s cap +instead of polling a fixed 5s forever, while burst events arriving mid-backoff are still delivered +with fenced, in-order cursor advance that resets the cadence to the fast 5s base. The invariant: +idle/active/waiting feedback from each poll drives the rescheduling delay; fencing/ordering/ +at-least-once delivery logic is untouched (backoff only changes when poll() runs). Review fix +2026-08-18-00:55: retry and error outcomes are "waiting", not "idle" — they must not extend the +backoff toward the cap, so transient failures recover on the fast 5s cadence. +*/ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { + TASK_DELETED_OUTBOX_BACKOFF_STEP_MS, + TASK_DELETED_OUTBOX_MAX_POLL_MS, + TASK_DELETED_OUTBOX_POLL_MS, + TASK_DELETED_OUTBOX_POLL_JITTER_RATIO, + TaskDeletedOutboxConsumer, + applyPollJitter, +} from "../task-store/task-deleted-outbox-consumer.js"; +import { buildConsumerId } from "../task-store/task-lifecycle-consumer-identity.js"; +import type { TaskStore } from "../store.js"; + +const mocks = vi.hoisted(() => ({ + registerTaskLifecycleConsumer: vi.fn(), + setTaskLifecycleConsumerActive: vi.fn(), + acquireTaskLifecycleLease: vi.fn(), + readTaskLifecycleConsumerCursor: vi.fn(), + readTaskLifecycleEventBounds: vi.fn(), + listTaskLifecycleEvents: vi.fn(), + hasTaskLifecycleConsumerReceipt: vi.fn(), + acknowledgeTaskLifecycleEvent: vi.fn(), + advanceTaskLifecycleConsumerCursor: vi.fn(), + renewTaskLifecycleLease: vi.fn(), + setTaskLifecycleConsumerRetry: vi.fn(), + parkTaskLifecycleConsumerDeadLetter: vi.fn(), + releaseTaskLifecycleLease: vi.fn(), +})); + +vi.mock("../postgres/data-layer.js", () => ({ + recordRunAuditEvent: vi.fn().mockResolvedValue(undefined), +})); + +vi.mock("../task-store/task-lifecycle-consumer-registry.js", () => ({ + acknowledgeTaskLifecycleEvent: mocks.acknowledgeTaskLifecycleEvent, + acquireTaskLifecycleLease: mocks.acquireTaskLifecycleLease, + advanceTaskLifecycleConsumerCursor: mocks.advanceTaskLifecycleConsumerCursor, + hasTaskLifecycleConsumerReceipt: mocks.hasTaskLifecycleConsumerReceipt, + listTaskLifecycleEvents: mocks.listTaskLifecycleEvents, + registerTaskLifecycleConsumer: mocks.registerTaskLifecycleConsumer, + releaseTaskLifecycleLease: mocks.releaseTaskLifecycleLease, + readTaskLifecycleConsumerCursor: mocks.readTaskLifecycleConsumerCursor, + readTaskLifecycleEventBounds: mocks.readTaskLifecycleEventBounds, + renewTaskLifecycleLease: mocks.renewTaskLifecycleLease, + setTaskLifecycleConsumerActive: mocks.setTaskLifecycleConsumerActive, + setTaskLifecycleConsumerRetry: mocks.setTaskLifecycleConsumerRetry, + parkTaskLifecycleConsumerDeadLetter: mocks.parkTaskLifecycleConsumerDeadLetter, +})); + +type LifecycleEvent = { + eventId: string; + seq: bigint; + eventType: string; + taskId: string; + occurredAt: string; + createdAt: string; + payload: Record; +}; + +function makeCursor(lastAckedSeq = 0n) { + return { + lastAckedSeq, + retryAttempts: 0, + retryBackoffUntil: null, + leaseToken: "lease", + fencingToken: 41n, + leaseExpiresAt: "2026-01-01T00:00:60.000Z", + // Recent ack timestamp so needsReconciliation does not trigger the 30-day fallback (which + // requires a real layer.db the fake does not provide). + updatedAt: new Date().toISOString(), + }; +} + +function taskDeletedEvent(seq: bigint, taskId: string, eventId?: string): LifecycleEvent { + return { + eventId: eventId ?? `evt-${seq}`, + seq, + eventType: "task:deleted", + taskId, + occurredAt: "2026-01-02T00:00:00.000Z", + createdAt: "2026-01-02T00:00:00.000Z", + payload: { + taskId, + previousColumn: "done", + previousStatus: null, + deletedAt: "2026-01-02T00:00:00.000Z", + allowResurrection: false, + githubIssueAction: null, + closureContext: null, + deletedBy: null, + }, + }; +} + +/** Fake TaskStore carrying asyncLayer + consumerId + taskCache for the consumer's seam. */ +function makeStore(consumerId: string, taskCache: Map = new Map()) { + return { + asyncLayer: { projectId: `project-${consumerId}` }, + consumerId, + taskCache, + emitObservedTaskDeleted: vi.fn(), + } as unknown as TaskStore; +} + +/** Records the fake-clock time of every outbox read (one per completed poll). */ +function scriptDefaults(events: () => LifecycleEvent[]) { + const pollTimes: number[] = []; + mocks.acquireTaskLifecycleLease.mockImplementation(async () => ({ + token: "lease-token", + fencingToken: 41n, + expiresAt: "2026-01-01T00:00:60.000Z", + })); + mocks.readTaskLifecycleConsumerCursor.mockImplementation(async () => makeCursor()); + mocks.readTaskLifecycleEventBounds.mockImplementation(async () => ({ + oldestSeq: null, + oldestOccurredAt: null, + headSeq: 0n, + })); + mocks.listTaskLifecycleEvents.mockImplementation(async () => { + pollTimes.push(Date.now()); + return events(); + }); + mocks.hasTaskLifecycleConsumerReceipt.mockImplementation(async () => false); + mocks.acknowledgeTaskLifecycleEvent.mockImplementation(async () => true); + mocks.renewTaskLifecycleLease.mockImplementation(async () => true); + mocks.registerTaskLifecycleConsumer.mockImplementation(async () => undefined); + mocks.setTaskLifecycleConsumerActive.mockImplementation(async () => undefined); + mocks.releaseTaskLifecycleLease.mockImplementation(async () => undefined); + return { pollTimes }; +} + +function collect(times: number[]): number[] { + const sorted = [...times].sort((a, b) => a - b); + const gaps: number[] = []; + for (let i = 1; i < sorted.length; i++) gaps.push(sorted[i] - sorted[i - 1]); + return gaps; +} + +beforeEach(() => { + vi.useFakeTimers(); + // Fail fast if an unmocked registry call leaks through to DB code. + vi.clearAllMocks(); +}); + +afterEach(() => { + vi.useRealTimers(); +}); + +describe("TaskDeletedOutboxConsumer idle backoff", () => { + it("backs an idle dashboard consumer way below the fixed-5s poll rate and reaches the 60s cap", async () => { + const { pollTimes } = scriptDefaults(() => []); + const store = makeStore(buildConsumerId("dashboard")); + const consumer = new TaskDeletedOutboxConsumer(store); + await consumer.start(); + const startPolls = pollTimes.length; + + // Fixed-5s poller would fire ~12 times in 60s; a backed-off idle consumer must not. + await vi.advanceTimersByTimeAsync(60_000); + const reads = pollTimes.length - startPolls; + expect(reads).toBeLessThan(6); + + // Let the cadence plateau so the deterministic delay provably hits the cap. + await vi.advanceTimersByTimeAsync(300_000); + const computeNextDelay = (consumer as unknown as { computeNextPollDelayMs(): number }).computeNextPollDelayMs; + expect(computeNextDelay.call(consumer)).toBe(TASK_DELETED_OUTBOX_MAX_POLL_MS); + + // De-synchronization jitter keeps the capped delay bounded and never below the fast base. + for (let i = 0; i < 200; i++) { + const jittered = applyPollJitter(TASK_DELETED_OUTBOX_MAX_POLL_MS); + expect(jittered).toBeGreaterThanOrEqual(TASK_DELETED_OUTBOX_POLL_MS); + expect(jittered).toBeLessThanOrEqual(TASK_DELETED_OUTBOX_MAX_POLL_MS * (1 + TASK_DELETED_OUTBOX_POLL_JITTER_RATIO)); + } + + await consumer.stop(); + }); + + it("keeps an active engine consumer on the fast 5s cadence while it delivers events", async () => { + const consumerId = buildConsumerId("engine"); + const store = makeStore( + consumerId, + new Map([["FN-STAY", { id: "FN-STAY", title: "stay" }]]), + ); + const { pollTimes } = scriptDefaults(() => [taskDeletedEvent(10n, "FN-STAY")]); + const consumer = new TaskDeletedOutboxConsumer(store); + await consumer.start(); + + await vi.advanceTimersByTimeAsync(60_000); + const gaps = collect(pollTimes); + // Active cadence stays at the jittered 5s base (4000-6000ms), with no spurious backoff. + expect(gaps).not.toHaveLength(0); + expect(Math.min(...gaps)).toBeGreaterThanOrEqual(TASK_DELETED_OUTBOX_POLL_MS * (1 - TASK_DELETED_OUTBOX_POLL_JITTER_RATIO)); + expect(Math.max(...gaps)).toBeLessThanOrEqual(TASK_DELETED_OUTBOX_POLL_MS * (1 + TASK_DELETED_OUTBOX_POLL_JITTER_RATIO)); + // Events were actually delivered through the fenced ack seam. + expect(store.emitObservedTaskDeleted).toHaveBeenCalled(); + expect(mocks.acknowledgeTaskLifecycleEvent).toHaveBeenCalled(); + + await consumer.stop(); + }); + + it("backs off an idle engine consumer independently from a concurrent active dashboard consumer", async () => { + // Both consumers share the module-level listTaskLifecycleEvents mock, so the reader dispatches + // on the per-store layer.projectId (the only cross-consumer discriminator the registry receives) + // rather than arming two competing global mockImplementations. + const idleTimes: number[] = []; + const activeTimes: number[] = []; + scriptDefaults(() => []); + const idleStore = makeStore(buildConsumerId("engine"), new Map([["FN-ACTIVE", { id: "FN-ACTIVE", title: "a" }]])); + const activeStore = makeStore(buildConsumerId("dashboard"), new Map([["FN-ACTIVE", { id: "FN-ACTIVE", title: "a" }]])); + + let seq = 100n; + mocks.listTaskLifecycleEvents.mockImplementation(async (layer: { projectId: string }) => { + if (layer.projectId === idleStore.asyncLayer!.projectId) { + idleTimes.push(Date.now()); + return []; // idle engine consumer -> empty outbox every poll + } + activeTimes.push(Date.now()); + seq += 1n; + return [taskDeletedEvent(seq, "FN-ACTIVE")]; // active dashboard consumer -> fresh event + }); + + const idleConsumer = new TaskDeletedOutboxConsumer(idleStore); + await idleConsumer.start(); + // Give the idle consumer a head start so backoff is already engaged. + await vi.advanceTimersByTimeAsync(80_000); + + const activeConsumer = new TaskDeletedOutboxConsumer(activeStore); + await activeConsumer.start(); + await vi.advanceTimersByTimeAsync(60_000); + const idleGaps = collect(idleTimes); + const activeGaps = collect(activeTimes); + // No cross-talk: the idle engine consumer backs off well beyond a 5s cadence... + expect(idleGaps).toContainEqual(expect.any(Number)); + expect(Math.max(...idleGaps)).toBeGreaterThanOrEqual(TASK_DELETED_OUTBOX_POLL_MS * 2); + // ...while the active dashboard consumer keeps polling at the fast base. + expect(activeGaps).not.toHaveLength(0); + expect(Math.max(...activeGaps)).toBeLessThanOrEqual(TASK_DELETED_OUTBOX_POLL_MS * (1 + TASK_DELETED_OUTBOX_POLL_JITTER_RATIO)); + + await activeConsumer.stop(); + await idleConsumer.stop(); + }); + + it("delivers a burst arriving mid-backoff in order with fenced acks and resets to the fast cadence", async () => { + const burst = [ + taskDeletedEvent(1n, "FN-A"), + taskDeletedEvent(2n, "FN-B"), + taskDeletedEvent(3n, "FN-C"), + ]; + const store = makeStore( + buildConsumerId("dashboard"), + new Map([ + ["FN-A", { id: "FN-A", title: "a" }], + ["FN-B", { id: "FN-B", title: "b" }], + ["FN-C", { id: "FN-C", title: "c" }], + ]), + ); + // Start idle, then arm a one-time burst once the consumer has backed off. The stateful reader + // returns the full burst exactly once (one poll) and then the empty outbox, so the burst is + // consumed by a single mid-backoff poll rather than re-delivered every poll in the window. + let burstLeft = 0; + const readEvents = (): LifecycleEvent[] => { + if (burstLeft > 0) { + burstLeft -= 1; + return burst; + } + return []; + }; + const { pollTimes } = scriptDefaults(() => readEvents()); + const consumer = new TaskDeletedOutboxConsumer(store); + await consumer.start(); + await vi.advanceTimersByTimeAsync(60_000); // engaged deeper backoff + const beforeBurstPolls = pollTimes.length; + + const preBurstAcks = mocks.acknowledgeTaskLifecycleEvent.mock.calls.length; + burstLeft = 1; // burst arrives mid-backoff; next delayed poll consumes all of it + await vi.advanceTimersByTimeAsync(120_000); + + // The burst poll delivered all N events in order. + const deliveredIds = store.emitObservedTaskDeleted.mock.calls.map((c) => c[0].id); + expect(deliveredIds).toEqual(["FN-A", "FN-B", "FN-C"]); + const burstAcks = mocks.acknowledgeTaskLifecycleEvent.mock.calls.slice(preBurstAcks); + expect(burstAcks).toHaveLength(3); + // acknowledgeTaskLifecycleEvent(layer, { consumerId, eventId, seq, priorSeq, fencingToken }): the + // fenced-ack payload is the second argument, so the ack assertions read call[1]. + const ackedSeqs = burstAcks.map((c) => c[1].seq); + expect(ackedSeqs).toEqual([1n, 2n, 3n]); // cursor advanced in order + for (const call of burstAcks) { + expect(call[1].fencingToken).toBe(41n); // every fenced ack honors the acquired lease token + } + + // The burst poll was observed (outbox read happened), proving mid-backoff delivery. + const deliveredSince = store.emitObservedTaskDeleted.mock.calls.length; + expect(deliveredSince).toBe(3); + expect(pollTimes.length).toBeGreaterThan(beforeBurstPolls); + + // Cadence resets to the fast base after the burst (idlePollsSinceEvent cleared). + const burstAt = pollTimes[beforeBurstPolls]; + const nextRead = pollTimes[beforeBurstPolls + 1]; + const resetGap = nextRead - burstAt; + expect(resetGap).toBeGreaterThanOrEqual(TASK_DELETED_OUTBOX_POLL_MS * (1 - TASK_DELETED_OUTBOX_POLL_JITTER_RATIO)); + expect(resetGap).toBeLessThanOrEqual(TASK_DELETED_OUTBOX_POLL_MS * (1 + TASK_DELETED_OUTBOX_POLL_JITTER_RATIO)); + + await consumer.stop(); + }); + + it("disarms the rescheduled timer on stop so shutdown cannot leave an orphaned poll", async () => { + const { pollTimes } = scriptDefaults(() => []); + const store = makeStore(buildConsumerId("dashboard")); + const consumer = new TaskDeletedOutboxConsumer(store); + await consumer.start(); + + await consumer.stop(); + const before = pollTimes.length; + await vi.advanceTimersByTimeAsync(10 * TASK_DELETED_OUTBOX_MAX_POLL_MS); + // No further polls fire after stop(); a manual poll() reports "waiting" (not running) and does + // not read the outbox. + expect(pollTimes.length).toBe(before); + const result = await consumer.poll(); + expect(result).toBe("waiting"); + expect(pollTimes.length).toBe(before); + + // A second stop() must also be safe (no throw, no timer double-clear crash). + await expect(consumer.stop()).resolves.toBeUndefined(); + }); + + it("keeps the backoff step constant growing toward the cap and never below base", () => { + // The deterministic growth curve: 5s -> 15s -> 25s -> ... -> 60s. + expect(TASK_DELETED_OUTBOX_BACKOFF_STEP_MS).toBe(10_000); + const consumer = new TaskDeletedOutboxConsumer(makeStore(buildConsumerId("engine"))); + const computeNextDelay = (consumer as unknown as { computeNextPollDelayMs(): number }).computeNextPollDelayMs; + // Drives the private idle counter through repeated simulated idle polls via the private seam. + const state = consumer as unknown as { idlePollsSinceEvent: number }; + for (let idle = 1; idle <= 10; idle++) { + state.idlePollsSinceEvent = idle; + expect(computeNextDelay.call(consumer)).toBe( + Math.min(TASK_DELETED_OUTBOX_POLL_MS + idle * TASK_DELETED_OUTBOX_BACKOFF_STEP_MS, TASK_DELETED_OUTBOX_MAX_POLL_MS), + ); + } + expect(computeNextDelay.call(consumer)).toBe(TASK_DELETED_OUTBOX_MAX_POLL_MS); + }); + + it("does not classify poll errors as idle: the error streak stays on the fast 5s base and recovers quickly", async () => { + const store = makeStore(buildConsumerId("engine"), new Map([[ + "FN-ERR", { id: "FN-ERR", title: "err" }, + ]])); + const { pollTimes } = scriptDefaults(() => [taskDeletedEvent(5n, "FN-ERR")]); + const consumer = new TaskDeletedOutboxConsumer(store); + await consumer.start(); + + // Simulate a transient DB failure on the next polls (the outbox reader records the poll time, + // then throws; pollSafely maps the error to "waiting"). + mocks.listTaskLifecycleEvents.mockImplementation(async () => { + pollTimes.push(Date.now()); + throw new Error("pg transient failure"); + }); + await vi.advanceTimersByTimeAsync(20_000); + const errorGaps = collect(pollTimes); + // Error polls are "waiting", never "idle": the cadence must stay at the jittered 5s base + // (4000-6000ms) — an idle streak of the same length would have grown it to 35s+. + expect(errorGaps.length).toBeGreaterThan(0); + expect(Math.max(...errorGaps)).toBeLessThanOrEqual(TASK_DELETED_OUTBOX_POLL_MS * (1 + TASK_DELETED_OUTBOX_POLL_JITTER_RATIO)); + + // The failure clears: the very next poll (still on the 5s base, not hidden behind a grown + // cap) delivers the pending event. (At-least-once: the mock receipt stays false, so the event + // is redelivered.) + mocks.listTaskLifecycleEvents.mockImplementation(async () => { pollTimes.push(Date.now()); return [taskDeletedEvent(5n, "FN-ERR")]; }); + const beforeRecovery = pollTimes.length; + const callsBeforeRecovery = store.emitObservedTaskDeleted.mock.calls.length; + await vi.advanceTimersByTimeAsync(TASK_DELETED_OUTBOX_POLL_MS * 2 * (1 + TASK_DELETED_OUTBOX_POLL_JITTER_RATIO)); + expect(pollTimes.length).toBeGreaterThan(beforeRecovery); + expect(store.emitObservedTaskDeleted.mock.calls.length).toBeGreaterThan(callsBeforeRecovery); + + await consumer.stop(); + }); + + it("does not classify the per-event retry-backoff window as idle", async () => { + const store = makeStore(buildConsumerId("engine")); + scriptDefaults(() => []); + // The cursor is inside its per-event retry window: the poll must take the "waiting" early + // return BEFORE reading the outbox. Record the cursor reads (one per waiting poll) instead of + // outbox reads, which the early return never reaches. + const cursorReads: number[] = []; + mocks.readTaskLifecycleConsumerCursor.mockImplementation(async () => { + cursorReads.push(Date.now()); + return { ...makeCursor(), retryBackoffUntil: new Date(Date.now() + 300_000).toISOString() }; + }); + const consumer = new TaskDeletedOutboxConsumer(store); + await consumer.start(); + + await vi.advanceTimersByTimeAsync(20_000); + const gaps = collect(cursorReads); + // Waiting polls stay on the fast 5s base — the retry window is a scheduled wait, not an idle + // streak growing the delay toward the 60s cap. + expect(gaps.length).toBeGreaterThan(0); + expect(Math.max(...gaps)).toBeLessThanOrEqual(TASK_DELETED_OUTBOX_POLL_MS * (1 + TASK_DELETED_OUTBOX_POLL_JITTER_RATIO)); + // The outbox was never read during the window (the early return precedes it). + expect(mocks.listTaskLifecycleEvents).not.toHaveBeenCalled(); + + await consumer.stop(); + }); +}); \ No newline at end of file diff --git a/packages/core/src/task-store/task-deleted-outbox-consumer.ts b/packages/core/src/task-store/task-deleted-outbox-consumer.ts index 8075fa5e75..ad6808b746 100644 --- a/packages/core/src/task-store/task-deleted-outbox-consumer.ts +++ b/packages/core/src/task-store/task-deleted-outbox-consumer.ts @@ -22,12 +22,26 @@ import { } from "./task-lifecycle-consumer-registry.js"; export const TASK_DELETED_OUTBOX_POLL_MS = 5_000; +export const TASK_DELETED_OUTBOX_MAX_POLL_MS = 60_000; +export const TASK_DELETED_OUTBOX_BACKOFF_STEP_MS = 10_000; +export const TASK_DELETED_OUTBOX_POLL_JITTER_RATIO = 0.2; export const TASK_DELETED_OUTBOX_LEASE_MS = 15_000; export const TASK_DELETED_OUTBOX_BATCH_SIZE = 100; export const TASK_DELETED_OUTBOX_RETENTION_DAYS = 30; const outboxConsumerLog = createLogger("task-deleted-outbox-consumer"); +/** + * Apply ±`TASK_DELETED_OUTBOX_POLL_JITTER_RATIO` multiplicative jitter to a poll delay so the ~44 + * per-project dashboard/engine consumers de-synchronize instead of firing on the same cadence. + * Imported by the regression test to assert its bounded, non-negative range. The delay never falls + * below the fast base so jitter cannot speed an idle consumer back into the storm. + */ +export function applyPollJitter(delayMs: number, ratio = TASK_DELETED_OUTBOX_POLL_JITTER_RATIO): number { + const delta = delayMs * ratio * (Math.random() * 2 - 1); + return Math.max(TASK_DELETED_OUTBOX_POLL_MS, Math.round(delayMs + delta)); +} + type OutboxEventForValidation = { eventId: string; eventType: string; @@ -66,24 +80,62 @@ function assertTaskDeletedOutboxEvent(event: OutboxEventForValidation): void { * before the durable receipt/cursor acknowledgement, intentionally yielding at-least-once * observed notifications in the crash window; observed dispatch has no writer-owned effects. */ + +/** + * Tri-state poll outcome that drives the idle backoff. Only "idle" (the outbox was genuinely + * empty) extends the delay toward the 60s cap. "active" means at least one event was + * delivered/processed. "waiting" covers retry-backoff windows, lease contention, fencing, + * errors, and shutdown races: none of those mean the outbox is idle, so none may extend the + * backoff. Resetting to the fast base on "waiting" keeps transient failures on a 5s recovery + * cadence instead of letting an error streak masquerade as an idle streak and hide recovery + * behind the 60s cap. + */ +export type TaskDeletedOutboxPollOutcome = "active" | "idle" | "waiting"; + export class TaskDeletedOutboxConsumer { - private pollTimer: ReturnType | null = null; + private pollTimer: ReturnType | null = null; private renewalTimer: ReturnType | null = null; private running = false; private lease: TaskLifecycleLease | null = null; + private idlePollsSinceEvent = 0; constructor(private readonly store: TaskStore) {} + /* + FNXC:TaskLifecycleConsumerIdleBackoff 2026-08-13-06:41: + The outbox consumer reschedules itself instead of polling a fixed 5s setInterval forever. Each + poll feeds back its outcome: an idle poll (the outbox returned zero events) grows the next delay + by TASK_DELETED_OUTBOX_BACKOFF_STEP_MS toward TASK_DELETED_OUTBOX_MAX_POLL_MS, with ±20% jitter + so the ~44 per-project dashboard/engine consumers de-synchronize instead of thundering together; + a poll that delivered events resets to the fast TASK_DELETED_OUTBOX_POLL_MS base. A paused/idle + project stops writing lifecycle events, so its outbox drains and backoff alone drops the DB + idx_scan/CPU storm on task_lifecycle_consumer_cursors without touching delivery semantics. Any + new event mid-backoff resets the cadence to 5s, bounding delivery latency. Cursor fencing, lease + advance, per-event ordering, and at-least-once delivery are unchanged — backoff only changes when + poll() runs, never the poll/dispatch/ack logic. + + FNXC:TaskLifecycleConsumerIdleBackoff 2026-08-18-00:55 (RUFU-074 review fix): + Review feedback: retry and error outcomes must not be classified as idle polls. The outcome is + now tri-state (TaskDeletedOutboxPollOutcome): only a genuinely empty outbox reports "idle" and + extends the backoff. Retry-backoff windows (cursor.retryBackoffUntil), lease contention, fencing, + and poll errors report "waiting" and reset the idle streak to the fast 5s base, so a transient + failure recovers at 5s cadence instead of an error streak growing the delay to the 60s cap and + hiding recovery behind it. + */ async start(): Promise { if (this.running || !this.store.asyncLayer || !this.store.consumerId) return; this.running = true; - await this.pollSafely(); - this.pollTimer = setInterval(() => void this.pollSafely(), TASK_DELETED_OUTBOX_POLL_MS); + /* Drain any backlog with the immediate poll, then feed its outcome into the first scheduled + delay so a fresh idle consumer backs off from the very first timer (not after one wasted 5s + tick). */ + const outcome = await this.pollSafely(); + this.recordPollOutcome(outcome); + this.scheduleNextPoll(); } async stop(): Promise { this.running = false; - if (this.pollTimer) clearInterval(this.pollTimer); + if (this.pollTimer) clearTimeout(this.pollTimer); if (this.renewalTimer) clearInterval(this.renewalTimer); this.pollTimer = null; this.renewalTimer = null; @@ -106,18 +158,28 @@ export class TaskDeletedOutboxConsumer { } } - private async pollSafely(): Promise { + private async pollSafely(): Promise { try { - await this.poll(); + return await this.poll(); } catch (error) { outboxConsumerLog.warn("Task:deleted outbox poll failed; delivery will retry", error); + return "waiting"; } } - async poll(): Promise { + /** + * Idle/active/waiting backoff contract: returns "active" when the poll delivered/processed at + * least one lifecycle event, "idle" when it found the outbox genuinely empty, and "waiting" when + * it took a non-delivering early return that does NOT mean the outbox is idle (retry-backoff + * window, lease contention, fencing, errors, shutdown races). The rescheduling loop extends the + * next poll delay toward the 60s cap on "idle" only and resets to the fast 5s base on "active" + * and "waiting", so transient failures recover quickly instead of being hidden behind an idle + * cap grown out of an error streak. + */ + async poll(): Promise { const layer = this.store.asyncLayer; const consumerId = this.store.consumerId; - if (!this.running || !layer || !consumerId) return; + if (!this.running || !layer || !consumerId) return "waiting"; await registerTaskLifecycleConsumer(layer, consumerId); /* FNXC:CrossProcessDeleteObservation 2026-08-01-12:14: @@ -126,7 +188,7 @@ export class TaskDeletedOutboxConsumer { */ if (!this.running) { await setTaskLifecycleConsumerActive(layer, consumerId, false); - return; + return "waiting"; } const now = new Date(); const acquired = await acquireTaskLifecycleLease( @@ -136,27 +198,29 @@ export class TaskDeletedOutboxConsumer { new Date(now.getTime() + TASK_DELETED_OUTBOX_LEASE_MS).toISOString(), now.toISOString(), ); - if (!acquired) return; + if (!acquired) return "waiting"; if (!this.running) { await releaseTaskLifecycleLease(layer, consumerId, acquired); - return; + return "waiting"; } this.lease = acquired; this.startRenewal(acquired); try { const cursor = await readTaskLifecycleConsumerCursor(layer, consumerId); - if (!cursor) return; - if (cursor.retryBackoffUntil && Date.parse(cursor.retryBackoffUntil) > Date.now()) return; + if (!cursor) return "waiting"; + // In the per-event retry window: the cursor already scheduled when the failed event may be + // retried. This is a retry wait, not an idle outbox — the next probe stays on the fast base. + if (cursor.retryBackoffUntil && Date.parse(cursor.retryBackoffUntil) > Date.now()) return "waiting"; const reconciliationReason = await this.needsReconciliation(cursor.lastAckedSeq, cursor.updatedAt); if (reconciliationReason) { const reconciled = await this.reconcile(cursor.lastAckedSeq, acquired, reconciliationReason); if (!reconciled) { await this.recordLeaseFenced(acquired, 0); - return; + return "waiting"; } } const currentCursor = await readTaskLifecycleConsumerCursor(layer, consumerId); - if (!currentCursor) return; + if (!currentCursor) return "waiting"; const events = await listTaskLifecycleEvents(layer, currentCursor.lastAckedSeq, TASK_DELETED_OUTBOX_BATCH_SIZE); let priorSeq = currentCursor.lastAckedSeq; let dispatchedCount = 0; @@ -228,6 +292,7 @@ export class TaskDeletedOutboxConsumer { }); } if (this.running) await setTaskLifecycleConsumerActive(layer, consumerId, true); + return events.length > 0 ? "active" : "idle"; } finally { if (this.renewalTimer) clearInterval(this.renewalTimer); this.renewalTimer = null; @@ -257,6 +322,48 @@ export class TaskDeletedOutboxConsumer { return null; } + /** + * Reschedule the next poll with a setTimeout whose delay reflects the previous poll's outcome, + * replacing the old fixed-interval setInterval so an idle consumer stops thundering against the + * DB. pollSafely always resolves (it maps errors to "waiting"), so the reschedule chain never + * stalls. + */ + private scheduleNextPoll(): void { + if (!this.running) return; + if (this.pollTimer) clearTimeout(this.pollTimer); + this.pollTimer = setTimeout(() => { + this.pollTimer = null; + void this.pollSafely().then((outcome) => { + this.recordPollOutcome(outcome); + this.scheduleNextPoll(); + }); + }, this.nextPollDelayMs()); + } + + /** + * FNXC:TaskLifecycleConsumerIdleBackoff 2026-08-18-00:55 (RUFU-074 review fix): + * Only a genuinely idle poll (empty outbox) extends the backoff. Active and waiting outcomes + * reset the idle streak so the next poll lands on the fast 5s base — an error or retry streak + * must never masquerade as an idle streak and push the cadence toward the 60s cap. + */ + private recordPollOutcome(outcome: TaskDeletedOutboxPollOutcome): void { + if (outcome === "idle") this.idlePollsSinceEvent += 1; + else this.idlePollsSinceEvent = 0; + } + + /** Jittered delay for the upcoming poll, derived from the accumulated idle-poll count. */ + private nextPollDelayMs(): number { + return applyPollJitter(this.computeNextPollDelayMs()); + } + + /** Deterministic idle-backoff delay (no jitter): 5s -> 15s -> 25s -> ... -> 60s cap. */ + private computeNextPollDelayMs(): number { + if (this.idlePollsSinceEvent <= 0) return TASK_DELETED_OUTBOX_POLL_MS; + const grown = TASK_DELETED_OUTBOX_POLL_MS + + this.idlePollsSinceEvent * TASK_DELETED_OUTBOX_BACKOFF_STEP_MS; + return Math.min(grown, TASK_DELETED_OUTBOX_MAX_POLL_MS); + } + private async reconcile( priorSeq: bigint, lease: TaskLifecycleLease,