fix(RUFU-074): idle backoff + jitter for the task-deleted outbox consumer (#3471)

**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.

<!-- This is an auto-generated comment: release notes by coderabbit.ai
-->
## 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.
<!-- end of auto-generated comment: release notes by coderabbit.ai -->
This commit is contained in:
ischindl
2026-08-19 07:11:47 +02:00
committed by GitHub
parent f195ff5b3d
commit 72877c8cf9
4 changed files with 660 additions and 15 deletions

View File

@@ -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.

View File

@@ -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.

View File

@@ -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<string, unknown>;
};
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<string, unknown> = 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();
});
});

View File

@@ -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<typeof setInterval> | null = null;
private pollTimer: ReturnType<typeof setTimeout> | null = null;
private renewalTimer: ReturnType<typeof setInterval> | 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<void> {
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<void> {
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<void> {
private async pollSafely(): Promise<TaskDeletedOutboxPollOutcome> {
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<void> {
/**
* 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<TaskDeletedOutboxPollOutcome> {
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,