A background job processor.
The work finished.
Did the job?
A team asks for a CSV report. A worker generates it, commits the result, and disappears before recording success. Follow what the next worker must know to recover.
One small system, developed through five decisions. Read the implementation, run the failure, and inspect the evidence behind each promise.
The case file
Export a monthly team report.
A small CSV contains resolved-ticket counts for two teams. Requests identify a workspace, month, and immutable fixture revision. A new request is a new export; retrying the same request keeps its identity.
- The user needs
- An accepted job they can find again, a useful failure reason, and one intended report.
- We control
- Job storage, worker claims, the report generator, and the output store.
- The first boundary
- Local teaching system, synthetic data, at most six retained jobs. No public API or authentication service is running on this page.
- 01 / AcceptedThe job is recorded.
A retry of submission can find the same intent.
- 02 / Output committedThe CSV exists.
Its bytes and receipt share a transaction.
- 03 / AcknowledgedThe job says succeeded.
This later write can be missing after a crash.
Follow the pressure
Give each new mechanism a reason to exist.
A request should not own the whole export.
Generating a tiny CSV while handling its request is a reasonable beginning. When the source becomes slow, the client’s connection becomes a poor record of whether the work happened. A timeout leaves the caller guessing.
Persist a job before returning its identity. Bind a submission key to immutable input. A repeated submission returns the original job; the same key with different input fails. In the implementation, a transaction makes that lookup and insert one decision.
export function enqueue(
store: Store,
requestId: string,
input: ReportInput,
fault: Fault = 'none',
policy: OutputPolicy = 'per-job'
): string {
if (!/^[a-z0-9-]{1,48}$/.test(requestId)) throw new Error('Invalid request ID');
renderReport(input); // Validate before accepting; the invalid fault simulates a later permanent failure.
return store.transaction((db) => {
const existing = db.jobs.find(
(job) => job.requestId === requestId && job.input.workspace === input.workspace
);
if (existing) {
if (existing.fingerprint !== fingerprint(input))
throw new Error('Request ID reused with different input');
record(db, existing, 'request-reused', 'Repeated submission returned the original job.');
return existing.id;
}
if (db.jobs.length >= limits.jobs)
throw new Error('Demo storage limit reached: reset to start another run');
const job: Job = {
id: `export-${db.jobs.length + 1}`,
requestId,
input: { ...input },
fingerprint: fingerprint(input),
fault,
policy,
status: 'queued',
attempts: 0,
token: 0,
leaseUntil: 0,
due: db.clock,
createdAt: db.clock,
cancelRequested: false,
lastError: ''
};
db.jobs.push(job);
record(db, job, 'accepted', 'Job committed before acceptance was returned.');
return job.id;
});
} What this buys: acceptance can be reconciled independently of execution. It introduces storage, retention, and admission limits to operate.
For an HTTP service, authenticate and authorize before this call, then return an accepted job reference only after commit. Those HTTP and identity boundaries are outside this runnable example. A workspace field in a function argument is not an authorization check.
A queue needs an owner for work in flight.
Starting every accepted export immediately lets arrivals determine resource use. Keep pending work in acceptance order and let a fixed number of slots claim eligible jobs. A delayed retry does not block a later eligible export.
Removing a job from a list would lose it if the worker died. Keep the record, mark it running, and give its claim an expiry and an ownership token. Active work renews its lease. Expiry makes abandoned work eligible for recovery and invalidates the old token.
export function claim(store: Store): { jobId: string; token: number } | undefined {
return store.transaction((db) => {
const job = db.jobs.find(
(candidate) =>
(candidate.status === 'queued' || candidate.status === 'retry') && candidate.due <= db.clock
);
if (!job) return undefined;
job.status = 'running';
job.attempts += 1;
job.token += 1;
job.leaseUntil = db.clock + limits.lease;
record(
db,
job,
'claimed',
`Attempt ${job.attempts} owns token ${job.token}; lease ends at tick ${job.leaseUntil}.`
);
return { jobId: job.id, token: job.token };
});
}
function owned(db: Database, jobId: string, token: number): Job | undefined {
return db.jobs.find(
(job) =>
job.id === jobId &&
job.status === 'running' &&
job.token === token &&
job.leaseUntil > db.clock
);
} What this buys: another worker can recover an abandoned claim. Every output and acknowledgement must still check current ownership; an expired lease cannot stop an old process by itself.
The lab’s slots enforce a limit within one processor. The queue is a bounded array scanned in acceptance order, which keeps the fixture readable. A larger store needs indexed eligibility queries and an explicit fairness policy.
Another attempt needs a reason and a budget.
A temporarily unavailable report source might recover. An unsupported source schema needs a fix. Distinguish those outcomes before deciding to retry. Keep the attempt count and next eligible time with the job so losing a worker does not erase the budget.
function retryOrFail(db: Database, job: Job, reason: string, permanent = false) {
job.lastError = reason;
job.leaseUntil = 0;
job.token += 1; // Any old worker loses permission to mutate this job.
if (permanent || job.attempts >= limits.attempts) {
// A final lost acknowledgement must not hide a committed result.
const resultExists = db.artifacts.some((artifact) => artifact.jobId === job.id);
job.status = resultExists ? 'succeeded' : 'failed';
record(
db,
job,
resultExists ? 'reconciled' : 'failed',
resultExists ? 'Attempt budget ended; reconciled the committed CSV receipt.' : reason
);
return;
}
// Reproducible jitter for this fixture; production uses independent jitter and a wall-clock budget.
const jitter = Number(job.id.split('-')[1]) % 2;
job.due = db.clock + Math.min(8, 2 ** job.attempts) + jitter;
job.status = 'retry';
record(db, job, 'retry-scheduled', `${reason} Next attempt is eligible at tick ${job.due}.`);
} What this buys: transient failures get a bounded recovery path. Permanent failures stop. If the last allowed attempt left an output receipt, reconciliation can recognize success without another execution.
This fixture permits three attempts and adds reproducible jitter to the delay. A deployed service needs independent jitter, real request deadlines, a total retry-age budget, and guidance from dependencies such as Retry-After. Its retry traffic consumes the same bounded worker capacity as new work.
The result needs an identity that survives an attempt.
Now put the crash between CSV commit and job acknowledgement. The job still says running, but executing the report again could create another result. A successful retry alone does not tell us whether we duplicated the effect.
Use one output key per workspace and job. Atomically check ownership, look up the receipt, and store the CSV with its immutable-input fingerprint. Retry with the same identity reuses the stored result. The lab lets you replace this with a new key per attempt to expose the duplicate.
export function publish(store: Store, jobId: string, token: number): boolean {
return store.transaction((db) => {
const job = owned(db, jobId, token);
if (!job || job.cancelRequested) return false;
const key =
job.policy === 'per-job'
? `${job.input.workspace}/${job.id}`
: `${job.input.workspace}/${job.id}/attempt-${job.attempts}`;
const existing = db.artifacts.find((artifact) => artifact.key === key);
if (existing) {
if (existing.fingerprint !== job.fingerprint) throw new Error('Output key conflict');
record(
db,
job,
'output-reused',
'The original CSV receipt was reused; no second artifact was created.'
);
} else {
db.artifacts.push({
key,
jobId,
fingerprint: job.fingerprint,
attempt: job.attempts,
committedAt: db.clock,
csv: renderReport(job.input)
});
record(db, job, 'output-committed', `CSV bytes and receipt committed together as ${key}.`);
}
return true;
});
}
export function acknowledge(store: Store, jobId: string, token: number): boolean {
return store.transaction((db) => {
const job = owned(db, jobId, token);
if (!job) return false;
if (!db.artifacts.some((artifact) => artifact.jobId === jobId))
throw new Error('Cannot acknowledge missing output');
job.status = 'succeeded';
job.leaseUntil = 0;
job.lastError = '';
record(db, job, 'acknowledged', 'Job marked succeeded after its output was committed.');
return true;
});
} What this buys: one logical stored CSV for this job while its receipt is retained. The handler may run again. This does not promise exactly-once execution or deduplication at an unrelated external service.
We deliberately keep output commit and acknowledgement separate to expose this failure window. In this small application both records share one database, so combining them into one transaction would also be a sensible simplification. Once an effect belongs to another system, its own identity and reconciliation contract become essential.
Stop is a request with a boundary.
Queued work can be cancelled before it starts. Running work records a stop request, then checks it at a cooperative checkpoint. The UI keeps “stop requested” distinct from “cancelled.” An abandoned cancelled claim is resolved when its lease expires.
When the CSV has already committed, cancellation reports that it is too late. Deleting an output would be a separate authorized operation. Pretending that cancellation reverses a committed effect would mislead the caller.
These lifecycle questions also occur in real concurrent code. Go’s pipeline and cancellation guide shows why a stage that stops consuming must let upstream work stop too. Our manual clock models checkpoints; it does not execute goroutines or cancel network requests.
What this buys: the caller can distinguish requested, observed, and too late. Worker shutdown still needs its own policy: stop claiming, drain for a deadline, then recover abandoned leases.
Predict → run → inspect → explain
Make the uncertain outcome visible.
Start with the default run. Let export-1 crash after output, inspect its CSV beside the running job, recover the worker, and finish. Then reset with one output key per attempt and compare the artifact count.
The report exists. The job still says running.
Follow export-1 through an uncertain outcome. The other exports show whether useful work can continue.
Active run: 2 worker slots · crash after CSV commit · one output key per job
Three exports are queued. Predict the output count, then run until the first worker crashes.
- Logical clock
- 0ticks
- Active slots
- 0/ 2
- Oldest pending job
- 0ticks since acceptance
- Export-1 outputs
- 0across 0 attempt(s)
| Job | State | Attempts | Next boundary | Action |
|---|---|---|---|---|
| export-12026-08 | queued | 0 / 3 | Waiting for a slot | |
| export-22026-07 | queued | 0 / 3 | Waiting for a slot | |
| export-32026-06 | queued | 0 / 3 | Waiting for a slot |
Committed CSVs
No output has committed yet.
Recent events
t=0 export-3 · attempt 0
acceptedJob committed before acceptance was returned.
request-3t=0 export-2 · attempt 0
acceptedJob committed before acceptance was returned.
request-2t=0 export-1 · attempt 0
acceptedJob committed before acceptance was returned.
request-1
Newest first · last 10 of up to 120 retained events.
A trace you can reproduce
Two attempts. One stored report.
This is the deterministic export-1 sequence with the default fault and one output key per job, without extra manual actions. These are logical ticks from the implementation, not measured latency.
| Tick | Event | What we know |
|---|---|---|
| 0 | Accepted | The request has a stored job identity. No CSV exists. |
| 1 | Attempt 1 claims | Token 1 owns a lease. The job is running. |
| 3 | CSV commits; worker crashes | One CSV exists. Job acknowledgement is missing. |
| 7 | Lease expires | The old token is invalidated. A retry is scheduled for tick 10. |
| 10 | Attempt 2 claims | Token 3 owns the new attempt; the original output key remains stable. |
| 12 | Output receipt reused | The stored CSV count remains one. |
| 13 | Acknowledged | The job now agrees with its committed result. |
Read identities by responsibility.
The submission key finds the original request. The job ID survives all attempts. Attempt numbers distinguish executions. An ownership token fences writes. The output key identifies the intended effect. Copying one ID into every field would erase these distinctions.
Observe what can become stuck.
Track oldest pending age, active slots, retry eligibility, terminal failures, and committed-but-unacknowledged results. The event list correlates this fixture’s work; it is a bounded diagnostic history. The page does not operate a monitoring backend or send alerts.
From browser state to a process boundary
Run the same processor against SQLite.
The browser keeps state in memory. The local runner uses a SQLite Store, commits the output, and exits the Node process with code 17 before acknowledging. A second process reads that database and recovers the original job.
cd src/lib/content/case-studies/background-job-processor
# Choose a fresh filename. The crash phase deliberately exits with code 17.
node --experimental-strip-types run.ts crash /tmp/heyrian-report-case.sqlite
# Start a new process against the same database.
node --experimental-strip-types run.ts recover /tmp/heyrian-report-case.sqlite
node --experimental-strip-types run.ts inspect /tmp/heyrian-report-case.sqlite Use a fresh database filename for each run; existing runs are preserved. The first command’s exit code 17 is the intended crash point. The recovery output should show succeeded, two attempts, and the same single artifact. The clocks still advance in deterministic steps in both processes.
Inspect the real transaction boundary
The SQLite adapter acquires its write transaction before reading the state, then commits the whole mutation or rolls it back. It stores the fixture in one JSON row. This favors inspection over throughput; it is not a normalized production queue schema.
transaction<T>(change: (draft: Database) => T): T {
this.db.exec('BEGIN IMMEDIATE');
try {
const draft = this.read();
const result = change(draft);
this.db
.prepare('UPDATE processor_state SET body = ? WHERE id = 1')
.run(JSON.stringify(draft));
this.db.exec('COMMIT');
return result;
} catch (error) {
this.db.exec('ROLLBACK');
throw error;
}
} SQLite’s transaction documentation describes BEGIN IMMEDIATE and write contention. This runner uses WAL and synchronous FULL. Its process-exit test does not establish resilience to disk failure or every operating-system crash.
Verified here
- Process exit after output commit, followed by recovery in a fresh process.
- Transaction rollback and reopening the database.
- Stable output identity versus the deliberate duplicate-output variant.
- Expired ownership, bounded attempts, slow-work heartbeats, cancellation, and admission bounds.
Execution scope
TypeScript in the browser; TypeScript with Node 22.21.1 and SQLite 3.50.4 for the local proof. Node’s SQLite module is experimental in that runtime. No Go implementation is included in this first case study.
The module is built into the verified Node runtime; no database server or additional package is needed. See the Node 22 SQLite API and TypeScript execution documentation.
yarn test:unit --run --project server src/lib/content/case-studies/background-job-processor Complete source · processor.ts
/** A deterministic report processor. Time and worker scheduling are explicit inputs. */
export type Fault = 'none' | 'slow' | 'transient' | 'invalid' | 'crash-after-output';
export type OutputPolicy = 'per-job' | 'per-attempt';
export type JobStatus = 'queued' | 'running' | 'retry' | 'succeeded' | 'failed' | 'cancelled';
export interface ReportInput {
workspace: string;
month: string;
}
export interface Job {
id: string;
requestId: string;
input: ReportInput;
fingerprint: string;
fault: Fault;
policy: OutputPolicy;
status: JobStatus;
attempts: number;
token: number;
leaseUntil: number;
due: number;
createdAt: number;
cancelRequested: boolean;
lastError: string;
}
export interface Artifact {
key: string;
jobId: string;
fingerprint: string;
attempt: number;
committedAt: number;
csv: string;
}
export interface ProcessorEvent {
tick: number;
jobId: string;
requestId: string;
attempt: number;
kind: string;
message: string;
}
export interface Database {
version: 1;
clock: number;
jobs: Job[];
artifacts: Artifact[];
events: ProcessorEvent[];
}
export const emptyDatabase = (): Database => ({
version: 1,
clock: 0,
jobs: [],
artifacts: [],
events: []
});
export interface Store {
read(): Database;
// A failed callback must commit nothing. Concurrent writers must serialize.
transaction<T>(change: (draft: Database) => T): T;
}
export class MemoryStore implements Store {
private state = emptyDatabase();
read() {
return structuredClone(this.state);
}
transaction<T>(change: (draft: Database) => T): T {
const draft = this.read();
const result = change(draft);
this.state = structuredClone(draft);
return result;
}
}
export const limits = { jobs: 6, attempts: 3, lease: 4, events: 120 } as const;
export const isTerminal = (job: Job) => ['succeeded', 'failed', 'cancelled'].includes(job.status);
export function record(db: Database, job: Job, kind: string, message: string) {
db.events.push({
tick: db.clock,
jobId: job.id,
requestId: job.requestId,
attempt: job.attempts,
kind,
message
});
db.events = db.events.slice(-limits.events);
}
export const fingerprint = (input: ReportInput) =>
JSON.stringify([input.workspace, input.month, 'fixture-v1']);
/** The report reads a fixed, synthetic data revision so a retry has identical meaning. */
export function renderReport(input: ReportInput): string {
if (!/^[a-z0-9-]{1,32}$/.test(input.workspace) || !/^2026-(0[1-9]|1[0-2])$/.test(input.month)) {
throw new Error('Invalid synthetic report input');
}
return `workspace,month,team,resolved\n${input.workspace},${input.month},support,42\n${input.workspace},${input.month},operations,17\n`;
}
export function enqueue(
store: Store,
requestId: string,
input: ReportInput,
fault: Fault = 'none',
policy: OutputPolicy = 'per-job'
): string {
if (!/^[a-z0-9-]{1,48}$/.test(requestId)) throw new Error('Invalid request ID');
renderReport(input); // Validate before accepting; the invalid fault simulates a later permanent failure.
return store.transaction((db) => {
const existing = db.jobs.find(
(job) => job.requestId === requestId && job.input.workspace === input.workspace
);
if (existing) {
if (existing.fingerprint !== fingerprint(input))
throw new Error('Request ID reused with different input');
record(db, existing, 'request-reused', 'Repeated submission returned the original job.');
return existing.id;
}
if (db.jobs.length >= limits.jobs)
throw new Error('Demo storage limit reached: reset to start another run');
const job: Job = {
id: `export-${db.jobs.length + 1}`,
requestId,
input: { ...input },
fingerprint: fingerprint(input),
fault,
policy,
status: 'queued',
attempts: 0,
token: 0,
leaseUntil: 0,
due: db.clock,
createdAt: db.clock,
cancelRequested: false,
lastError: ''
};
db.jobs.push(job);
record(db, job, 'accepted', 'Job committed before acceptance was returned.');
return job.id;
});
}
export function claim(store: Store): { jobId: string; token: number } | undefined {
return store.transaction((db) => {
const job = db.jobs.find(
(candidate) =>
(candidate.status === 'queued' || candidate.status === 'retry') && candidate.due <= db.clock
);
if (!job) return undefined;
job.status = 'running';
job.attempts += 1;
job.token += 1;
job.leaseUntil = db.clock + limits.lease;
record(
db,
job,
'claimed',
`Attempt ${job.attempts} owns token ${job.token}; lease ends at tick ${job.leaseUntil}.`
);
return { jobId: job.id, token: job.token };
});
}
function owned(db: Database, jobId: string, token: number): Job | undefined {
return db.jobs.find(
(job) =>
job.id === jobId &&
job.status === 'running' &&
job.token === token &&
job.leaseUntil > db.clock
);
}
export function publish(store: Store, jobId: string, token: number): boolean {
return store.transaction((db) => {
const job = owned(db, jobId, token);
if (!job || job.cancelRequested) return false;
const key =
job.policy === 'per-job'
? `${job.input.workspace}/${job.id}`
: `${job.input.workspace}/${job.id}/attempt-${job.attempts}`;
const existing = db.artifacts.find((artifact) => artifact.key === key);
if (existing) {
if (existing.fingerprint !== job.fingerprint) throw new Error('Output key conflict');
record(
db,
job,
'output-reused',
'The original CSV receipt was reused; no second artifact was created.'
);
} else {
db.artifacts.push({
key,
jobId,
fingerprint: job.fingerprint,
attempt: job.attempts,
committedAt: db.clock,
csv: renderReport(job.input)
});
record(db, job, 'output-committed', `CSV bytes and receipt committed together as ${key}.`);
}
return true;
});
}
export function acknowledge(store: Store, jobId: string, token: number): boolean {
return store.transaction((db) => {
const job = owned(db, jobId, token);
if (!job) return false;
if (!db.artifacts.some((artifact) => artifact.jobId === jobId))
throw new Error('Cannot acknowledge missing output');
job.status = 'succeeded';
job.leaseUntil = 0;
job.lastError = '';
record(db, job, 'acknowledged', 'Job marked succeeded after its output was committed.');
return true;
});
}
function retryOrFail(db: Database, job: Job, reason: string, permanent = false) {
job.lastError = reason;
job.leaseUntil = 0;
job.token += 1; // Any old worker loses permission to mutate this job.
if (permanent || job.attempts >= limits.attempts) {
// A final lost acknowledgement must not hide a committed result.
const resultExists = db.artifacts.some((artifact) => artifact.jobId === job.id);
job.status = resultExists ? 'succeeded' : 'failed';
record(
db,
job,
resultExists ? 'reconciled' : 'failed',
resultExists ? 'Attempt budget ended; reconciled the committed CSV receipt.' : reason
);
return;
}
// Reproducible jitter for this fixture; production uses independent jitter and a wall-clock budget.
const jitter = Number(job.id.split('-')[1]) % 2;
job.due = db.clock + Math.min(8, 2 ** job.attempts) + jitter;
job.status = 'retry';
record(db, job, 'retry-scheduled', `${reason} Next attempt is eligible at tick ${job.due}.`);
}
export function cancel(store: Store, jobId: string): string {
return store.transaction((db) => {
const job = db.jobs.find((candidate) => candidate.id === jobId);
if (!job) return 'Job not found.';
if (db.artifacts.some((artifact) => artifact.jobId === job.id))
return 'Too late to cancel: the CSV is already committed. Cancellation cannot undo it.';
if (isTerminal(job)) return `Job is already ${job.status}.`;
if (job.status === 'running') {
job.cancelRequested = true;
record(
db,
job,
'cancel-requested',
'Worker will observe cancellation at its next cooperative checkpoint.'
);
return 'Cancellation requested; work has not stopped yet.';
}
job.status = 'cancelled';
job.token += 1;
record(db, job, 'cancelled', 'Cancelled before the next attempt started.');
return 'Pending work cancelled.';
});
}
interface Work {
jobId: string;
token: number;
remaining: number;
phase: 'render' | 'ack';
}
export interface Worker {
id: string;
online: boolean;
work?: Work;
}
/** Worker slots are volatile. The Store owns job and output state across worker restarts. */
export class Processor {
private store: Store;
private slots: Worker[];
constructor(store: Store, concurrency = 2) {
if (!Number.isInteger(concurrency) || concurrency < 1 || concurrency > 3)
throw new Error('Choose 1–3 workers');
this.store = store;
this.slots = Array.from({ length: concurrency }, (_, index) => ({
id: `worker-${index + 1}`,
online: true
}));
}
workers(): Worker[] {
return structuredClone(this.slots);
}
recover(): void {
for (const slot of this.slots) slot.online = true;
}
crash(workerId: string): boolean {
const slot = this.slots.find((worker) => worker.id === workerId);
if (!slot?.online) return false;
if (slot.work) {
const jobId = slot.work.jobId;
this.store.transaction((db) => {
const job = db.jobs.find((item) => item.id === jobId)!;
record(
db,
job,
'worker-crashed',
`${slot.id} disappeared; its lease remains until expiry.`
);
});
}
slot.online = false;
slot.work = undefined;
return true;
}
advance(): void {
this.store.transaction((db) => {
db.clock += 1;
for (const job of db.jobs) {
if (job.status !== 'running' || job.leaseUntil > db.clock) continue;
if (job.cancelRequested && !db.artifacts.some((artifact) => artifact.jobId === job.id)) {
job.status = 'cancelled';
job.token += 1;
job.leaseUntil = 0;
record(
db,
job,
'cancelled',
'Expired attempt was cancelled before any output committed.'
);
} else retryOrFail(db, job, 'Worker lease expired.');
}
});
for (const slot of this.slots) {
if (!slot.online) continue;
const work = slot.work;
if (work) {
const job = this.store.transaction((db) => {
const active = owned(db, work.jobId, work.token);
if (!active) return undefined;
if (active.cancelRequested) {
active.status = 'cancelled';
active.token += 1;
active.leaseUntil = 0;
record(db, active, 'cancelled', 'Worker stopped at a checkpoint before publishing.');
return undefined;
}
active.leaseUntil = db.clock + limits.lease; // Heartbeat while this worker is making progress.
return structuredClone(active);
});
if (!job) {
slot.work = undefined;
continue;
}
if (work.phase === 'ack') {
acknowledge(this.store, work.jobId, work.token);
slot.work = undefined;
} else if (--work.remaining <= 0) {
if (job.fault === 'invalid' || (job.fault === 'transient' && job.attempts < 3)) {
this.store.transaction((db) => {
const active = owned(db, work.jobId, work.token);
if (active)
retryOrFail(
db,
active,
job.fault === 'invalid'
? 'Unsupported source schema.'
: 'Report source temporarily unavailable.',
job.fault === 'invalid'
);
});
slot.work = undefined;
} else if (publish(this.store, work.jobId, work.token)) {
work.phase = 'ack';
if (job.fault === 'crash-after-output' && job.attempts === 1) this.crash(slot.id);
} else slot.work = undefined;
}
}
if (!slot.work && slot.online) {
const next = claim(this.store);
if (next) {
const job = this.store.read().jobs.find((item) => item.id === next.jobId)!;
slot.work = { ...next, phase: 'render', remaining: job.fault === 'slow' ? 6 : 2 };
}
}
}
}
}
export function seed(
store: Store,
fault: Fault = 'crash-after-output',
policy: OutputPolicy = 'per-job'
) {
enqueue(store, 'request-1', { workspace: 'acme', month: '2026-08' }, fault, policy);
enqueue(store, 'request-2', { workspace: 'acme', month: '2026-07' });
enqueue(store, 'request-3', { workspace: 'acme', month: '2026-06' });
}
Complete source · sqlite-store.ts
import { DatabaseSync } from 'node:sqlite';
import { emptyDatabase, type Database, type Store } from './processor.ts';
/** A tiny local adapter. One JSON row keeps the teaching transaction boundary inspectable. */
export class SqliteStore implements Store {
private db: DatabaseSync;
constructor(filename: string) {
this.db = new DatabaseSync(filename);
this.db.exec(`
PRAGMA journal_mode = WAL;
PRAGMA synchronous = FULL;
PRAGMA busy_timeout = 3000;
CREATE TABLE IF NOT EXISTS processor_state (
id INTEGER PRIMARY KEY CHECK (id = 1),
body TEXT NOT NULL
);
`);
this.db
.prepare('INSERT OR IGNORE INTO processor_state (id, body) VALUES (1, ?)')
.run(JSON.stringify(emptyDatabase()));
}
read(): Database {
const row = this.db.prepare('SELECT body FROM processor_state WHERE id = 1').get()!;
const state = JSON.parse(String(row.body)) as Database;
if (state.version !== 1) throw new Error('Unsupported processor state version');
return state;
}
transaction<T>(change: (draft: Database) => T): T {
this.db.exec('BEGIN IMMEDIATE');
try {
const draft = this.read();
const result = change(draft);
this.db
.prepare('UPDATE processor_state SET body = ? WHERE id = 1')
.run(JSON.stringify(draft));
this.db.exec('COMMIT');
return result;
} catch (error) {
this.db.exec('ROLLBACK');
throw error;
}
}
close() {
this.db.close();
}
}
Complete source · run.ts
import { existsSync } from 'node:fs';
import { enqueue, isTerminal, Processor } from './processor.ts';
import { SqliteStore } from './sqlite-store.ts';
// Run with Node 22.21+; node:sqlite is experimental in the verified Node 22 runtime.
const [phase, filename] = process.argv.slice(2);
if (!filename || !['crash', 'recover', 'inspect'].includes(phase)) {
throw new Error(
'Usage: node --experimental-strip-types run.ts crash|recover|inspect /tmp/new-report-run.sqlite'
);
}
if (phase === 'crash' && existsSync(filename))
throw new Error('Choose a fresh database filename; existing runs are preserved');
if (phase !== 'crash' && !existsSync(filename)) throw new Error('Run the crash phase first');
const store = new SqliteStore(filename);
if (phase === 'crash') {
enqueue(store, 'request-1', { workspace: 'acme', month: '2026-08' }, 'crash-after-output');
}
const processor = new Processor(store, 1);
if (phase !== 'inspect') {
for (let step = 0; step < 80 && !store.read().jobs.every(isTerminal); step++) {
processor.advance();
if (phase === 'crash' && processor.workers().some((worker) => !worker.online)) {
console.log(JSON.stringify({ phase, state: store.read() }, null, 2));
// Deliberately exit without acknowledging the job or gracefully closing its database.
process.exit(17);
}
}
}
console.log(JSON.stringify({ phase, state: store.read() }, null, 2));
store.close();
The next pressures
Keep the guarantees attached to their boundaries.
What if the effect is outside this database?
This implementation can transact CSV bytes, their receipt, and ownership checks together. Uploading to object storage or sending an email moves the effect across a boundary. A database receipt written separately cannot make that external effect atomic. Use the receiver’s idempotency contract, immutable object identity and finalization, or explicit reconciliation; state any remaining duplicate risk.
Duplicate delivery is an ordinary queue concern. Amazon SQS documents at-least-once delivery for its standard queues. Adopting a queue product does not automatically deduplicate your handler’s effects.
What changes with real workers and more than one tenant?
The Store serializes mutations, but worker-slot limits are local to one Processor. Multiple processes need coordinated capacity if the limit is global. Use database time for leases, an indexed claim strategy, bounded task deadlines, and a policy for workers that continue after expiry. Cancellation must reach the actual I/O operation.
Authenticate submissions and authorize job status, cancellation, and CSV reads. Add per-workspace admission and fair scheduling so one tenant cannot monopolize the queue. The synthetic workspace field and six-record fixture cap do not implement these protections.
What changes when the input, schema, or retention policy evolves?
The report reads fixed fixture-v1 data. A real export needs to define whether it reads a captured source revision or current data on each attempt. Version job payloads and report recipes, migrate populated storage, and keep an explicit policy for old pending work.
Retain submission and output receipts for the supported retry window. Deleting a receipt removes evidence used for deduplication. This adapter accepts state version 1 and rejects other versions; a migration path is future work. Its 120-event ring is not an audit archive.
Ideas brought together
Name the job each idea is doing.
- Queue and deque
- Hold pending work until capacity is available; delayed retries need eligibility as well as order.
- Bounded parallelism
- Separate arrival volume from the number of tasks allowed to make progress.
- Cancellation propagation
- Carry intent to a checkpoint and distinguish a requested stop from an observed stop.
- Retry, backoff & idempotency
- Give uncertain work another attempt while preserving the identity of the intended effect.
The same leases, retries, and limit carry a report-sending job in Work queues and background jobs. If the storage boundary is the part you want to investigate next, follow Hexagonal / ports & adapters. For distinguishing retryable and permanent outcomes, revisit Kinds and sentinels.