diff --git a/CHANGELOG.md b/CHANGELOG.md index 9c29d53..ed09695 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,23 @@ # Changelog +## Unreleased + +- Make the failure-recovery demo's serialization proof count messages rather + than executions. A worker that loses its lease mid-operation leaves the + replacement to execute the same message again, which is the at-least-once + contract, so the proof failed on slower machines for behaviour it documents + elsewhere. Each serialization event now carries its attempt and process, so a + start pairs with its own finish instead of with whichever finish came next. + The proof asserts that every start has its own finish, that exactly the two + sent messages ran, and that the surviving attempt of each message never + overlaps another message's surviving attempt; a superseded attempt may + overlap anything, because it keeps running until it notices the lost lease and + its write is fenced out. The committed state check is unchanged, and the demo + reports the executions it saw. `assertSerializedExecution` moved into its own + module with unit coverage for the clean, retried, superseded-overlap, + still-running-replacement, unexplained-overlap, boundary, unfinished, + unmatched-finish, double-start, restart-after-finish, and lost-message cases. + ## 0.13.2 - 2026-08-17 - Accept a `key` on `schedule`, naming a reminder for the item it is waiting diff --git a/examples/failure-recovery/actor.ts b/examples/failure-recovery/actor.ts index a6235b1..c372732 100644 --- a/examples/failure-recovery/actor.ts +++ b/examples/failure-recovery/actor.ts @@ -25,16 +25,14 @@ export class RecoveryCounter extends Actor { async serialize({ controlDirectory }: { controlDirectory: string }): Promise { const message = this.currentMessage if (!message) throw new Error("serialize requires a durable message") - await appendFile( - join(controlDirectory, "serialization.jsonl"), - `${JSON.stringify({ event: "start", messageId: message.id, at: Date.now() })}\n`, - ) + // The attempt and the process identify the execution, so a start pairs with + // its own finish even when a superseded attempt outlives its replacement. + const execution = { messageId: message.id, attempt: message.attempt, processId: process.pid } + const path = join(controlDirectory, "serialization.jsonl") + await appendFile(path, `${JSON.stringify({ event: "start", ...execution, at: Date.now() })}\n`) await new Promise((resolve) => setTimeout(resolve, 100)) this.count += 1 - await appendFile( - join(controlDirectory, "serialization.jsonl"), - `${JSON.stringify({ event: "finish", messageId: message.id, at: Date.now() })}\n`, - ) + await appendFile(path, `${JSON.stringify({ event: "finish", ...execution, at: Date.now() })}\n`) return this.count } } diff --git a/examples/failure-recovery/demo.ts b/examples/failure-recovery/demo.ts index 9ab5edd..ba0f3a7 100644 --- a/examples/failure-recovery/demo.ts +++ b/examples/failure-recovery/demo.ts @@ -8,6 +8,11 @@ import { fork, type ChildProcess } from "node:child_process" import { createRuntime, type ActorReference, type MessageReference } from "solid-objects" import { sqlite } from "solid-objects/database/sqlite" import { RecoveryCounter } from "./actor.ts" +import { + assertSerializedExecution, + parseSerializationEvent, + type SerializationProof, +} from "./serialization.ts" interface WorkerMessage { event: string @@ -15,12 +20,6 @@ interface WorkerMessage { processed?: number } -interface SerializationEvent { - event: "start" | "finish" - messageId: string - at: number -} - interface ExternalEffectEvent { messageId: string attempt: number @@ -59,7 +58,7 @@ try { assert.equal(existsSync(directory), false) -async function proveSerialization(): Promise<{ finalState: number; overlap: false }> { +async function proveSerialization(): Promise { const controlDirectory = join(directory, "serialization") await mkdir(controlDirectory) const reference = runtime.ref(RecoveryCounter, "serialized") @@ -74,15 +73,10 @@ async function proveSerialization(): Promise<{ finalState: number; overlap: fals join(controlDirectory, "serialization.jsonl"), parseSerializationEvent, ) - assert.equal(events.length, 4) - const starts = events.filter((event) => event.event === "start") - const finishes = events.filter((event) => event.event === "finish") - assert.equal(starts.length, 2) - assert.equal(finishes.length, 2) - assert(Number(starts[1]?.at) >= Number(finishes[0]?.at)) + const proof = assertSerializedExecution(events, { messageCount: 2 }) const snapshot = await reference.snapshot() assert.equal(snapshot.count, 2) - return { finalState: snapshot.count, overlap: false } + return { ...proof, finalState: snapshot.count } } async function proveCrashRecovery(): Promise<{ @@ -200,18 +194,6 @@ async function readJsonLines( return (await readFile(path, "utf8")).trim().split("\n").filter(Boolean).map(parse) } -function parseSerializationEvent(line: string): SerializationEvent { - const event = JSON.parse(line) as Partial - if ( - (event.event !== "start" && event.event !== "finish") || - typeof event.messageId !== "string" || - typeof event.at !== "number" - ) { - throw new TypeError("invalid serialization event") - } - return { event: event.event, messageId: event.messageId, at: event.at } -} - function parseExternalEffectEvent(line: string): ExternalEffectEvent { const event = JSON.parse(line) as Partial if ( diff --git a/examples/failure-recovery/serialization.ts b/examples/failure-recovery/serialization.ts new file mode 100644 index 0000000..61cdd64 --- /dev/null +++ b/examples/failure-recovery/serialization.ts @@ -0,0 +1,133 @@ +import assert from "node:assert/strict" + +export interface SerializationEvent { + event: "start" | "finish" + messageId: string + attempt: number + processId: number + at: number +} + +export interface SerializationProof { + executions: number + retried: boolean + supersededOverlap: boolean +} + +interface Execution { + messageId: string + attempt: number + processId: number + startedAt: number + finishedAt: number +} + +export function parseSerializationEvent(line: string): SerializationEvent { + const event = JSON.parse(line) as Partial + if ( + (event.event !== "start" && event.event !== "finish") || + typeof event.messageId !== "string" || + typeof event.attempt !== "number" || + typeof event.processId !== "number" || + typeof event.at !== "number" + ) { + throw new TypeError("invalid serialization event") + } + return { + event: event.event, + messageId: event.messageId, + attempt: event.attempt, + processId: event.processId, + at: event.at, + } +} + +// One identity commits one state transition at a time. The control file is +// written outside the transaction, so it records execution attempts rather than +// commits: a worker that loses its lease keeps running until it notices, and its +// replacement executes the same message under a higher attempt. The superseded +// attempt may therefore overlap anything, because its write is fenced out and +// the committed state is what proves it. +// +// Each event carries its attempt and process, so a start pairs with its own +// finish rather than with whichever finish arrived next. Without that, a +// superseded attempt finishing late reads as its replacement finishing, and a +// second message could then overlap a replacement that is still running. +export function assertSerializedExecution( + events: readonly SerializationEvent[], + options: { messageCount: number }, +): SerializationProof { + const executions = pairExecutions(events) + + const messageIds = new Set(executions.map((execution) => execution.messageId)) + assert.equal( + messageIds.size, + options.messageCount, + `expected ${options.messageCount} messages to run, saw ${messageIds.size}`, + ) + + const survivingAttempt = new Map() + for (const execution of executions) { + const highest = survivingAttempt.get(execution.messageId) ?? 0 + if (execution.attempt > highest) survivingAttempt.set(execution.messageId, execution.attempt) + } + const surviving = executions.filter( + (execution) => survivingAttempt.get(execution.messageId) === execution.attempt, + ) + + for (const [index, execution] of surviving.entries()) { + for (const other of surviving.slice(index + 1)) { + assert( + !overlaps(execution, other), + `${describe(execution)} and ${describe(other)} overlap, and neither was superseded`, + ) + } + } + + const supersededOverlap = executions.some((execution) => + executions.some((other) => other !== execution && overlaps(execution, other)), + ) + + return { + executions: executions.length, + retried: executions.length > options.messageCount, + supersededOverlap, + } +} + +function pairExecutions(events: readonly SerializationEvent[]): Execution[] { + const started = new Map() + const executions: Execution[] = [] + + for (const event of [...events].sort((left, right) => left.at - right.at)) { + const key = `${event.messageId}#${event.attempt}#${event.processId}` + if (event.event === "start") { + assert(!started.has(key), `${describe(event)} started twice`) + started.set(key, event) + continue + } + const start = started.get(key) + assert(start !== undefined, `${describe(event)} finished with no matching start`) + started.delete(key) + executions.push({ + messageId: event.messageId, + attempt: event.attempt, + processId: event.processId, + startedAt: start.at, + finishedAt: event.at, + }) + } + + const unfinished = [...started.values()].map(describe) + assert.equal(unfinished.length, 0, `${unfinished.join(", ")} never wrote a finish`) + + return executions +} + +function overlaps(left: Execution, right: Execution): boolean { + return left.startedAt < right.finishedAt && right.startedAt < left.finishedAt +} + +function describe(execution: { messageId: string; attempt: number }): string { + return `${execution.messageId} attempt ${execution.attempt}` +} diff --git a/test/failure-recovery-serialization.test.ts b/test/failure-recovery-serialization.test.ts new file mode 100644 index 0000000..0982178 --- /dev/null +++ b/test/failure-recovery-serialization.test.ts @@ -0,0 +1,148 @@ +import { describe, expect, it } from "vitest" +import { + assertSerializedExecution, + type SerializationEvent, +} from "../examples/failure-recovery/serialization.js" + +function execution(options: { + messageId: string + attempt: number + startedAt: number + finishedAt: number + processId?: number +}): SerializationEvent[] { + const { messageId, attempt, startedAt, finishedAt, processId = attempt } = options + return [ + { event: "start", messageId, attempt, processId, at: startedAt }, + { event: "finish", messageId, attempt, processId, at: finishedAt }, + ] +} + +describe("serialization proof", () => { + it("accepts two executions that do not overlap", () => { + const events = [ + ...execution({ messageId: "a", attempt: 1, startedAt: 10, finishedAt: 20 }), + ...execution({ messageId: "b", attempt: 1, startedAt: 30, finishedAt: 40 }), + ] + + expect(assertSerializedExecution(events, { messageCount: 2 })).toEqual({ + executions: 2, + retried: false, + supersededOverlap: false, + }) + }) + + // A worker can lose its lease mid-operation, and the replacement executes the + // same message again under a higher attempt. That is the at-least-once + // contract, so the proof counts messages rather than executions. + it("accepts a message that executes twice after a lost lease", () => { + const events = [ + ...execution({ messageId: "a", attempt: 1, startedAt: 10, finishedAt: 20 }), + ...execution({ messageId: "a", attempt: 2, startedAt: 30, finishedAt: 40 }), + ...execution({ messageId: "b", attempt: 1, startedAt: 50, finishedAt: 60 }), + ] + + expect(assertSerializedExecution(events, { messageCount: 2 })).toEqual({ + executions: 3, + retried: true, + supersededOverlap: false, + }) + }) + + // The attempt that lost the lease is the one that was slow, so it is still + // running when its replacement starts. Its write is fenced out, and the + // committed state is what proves that, so the log may interleave here. + it("accepts a superseded attempt that outlives the start of its replacement", () => { + const events = [ + ...execution({ messageId: "a", attempt: 1, startedAt: 10, finishedAt: 30 }), + ...execution({ messageId: "a", attempt: 2, startedAt: 20, finishedAt: 40 }), + ...execution({ messageId: "b", attempt: 1, startedAt: 50, finishedAt: 60 }), + ] + + expect(assertSerializedExecution(events, { messageCount: 2 })).toEqual({ + executions: 3, + retried: true, + supersededOverlap: true, + }) + }) + + // The superseded attempt finishing late must not be read as its replacement + // finishing. The replacement is still running, so a second message that + // starts here is a real serialization failure. + it("rejects a second message that overlaps a still-running replacement", () => { + const events = [ + ...execution({ messageId: "a", attempt: 1, startedAt: 10, finishedAt: 35 }), + ...execution({ messageId: "a", attempt: 2, startedAt: 20, finishedAt: 60 }), + ...execution({ messageId: "b", attempt: 1, startedAt: 40, finishedAt: 50 }), + ] + + expect(() => assertSerializedExecution(events, { messageCount: 2 })).toThrow(/overlap/) + }) + + it("rejects executions that overlap when no message ran twice", () => { + const events = [ + ...execution({ messageId: "a", attempt: 1, startedAt: 10, finishedAt: 20 }), + ...execution({ messageId: "b", attempt: 1, startedAt: 15, finishedAt: 25 }), + ] + + expect(() => assertSerializedExecution(events, { messageCount: 2 })).toThrow(/overlap/) + }) + + it("accepts one execution that ends exactly as the next begins", () => { + const events = [ + ...execution({ messageId: "a", attempt: 1, startedAt: 10, finishedAt: 20 }), + ...execution({ messageId: "b", attempt: 1, startedAt: 20, finishedAt: 30 }), + ] + + expect(() => assertSerializedExecution(events, { messageCount: 2 })).not.toThrow() + }) + + it("rejects an execution that never finished", () => { + const events: SerializationEvent[] = [ + ...execution({ messageId: "a", attempt: 1, startedAt: 10, finishedAt: 20 }), + { event: "start", messageId: "b", attempt: 1, processId: 1, at: 30 }, + ] + + expect(() => assertSerializedExecution(events, { messageCount: 2 })).toThrow(/finish/) + }) + + it("rejects a finish with no start", () => { + const events: SerializationEvent[] = [ + ...execution({ messageId: "a", attempt: 1, startedAt: 10, finishedAt: 20 }), + { event: "finish", messageId: "b", attempt: 1, processId: 1, at: 30 }, + ] + + expect(() => assertSerializedExecution(events, { messageCount: 2 })).toThrow(/start/) + }) + + it("rejects one attempt that started twice", () => { + const events: SerializationEvent[] = [ + { event: "start", messageId: "a", attempt: 1, processId: 1, at: 10 }, + { event: "start", messageId: "a", attempt: 1, processId: 1, at: 15 }, + { event: "finish", messageId: "a", attempt: 1, processId: 1, at: 20 }, + { event: "finish", messageId: "a", attempt: 1, processId: 1, at: 25 }, + ...execution({ messageId: "b", attempt: 1, startedAt: 30, finishedAt: 40 }), + ] + + expect(() => assertSerializedExecution(events, { messageCount: 2 })).toThrow(/twice/) + }) + + it("rejects an attempt that starts again after it finished", () => { + const events: SerializationEvent[] = [ + ...execution({ messageId: "a", attempt: 1, startedAt: 10, finishedAt: 20 }), + { event: "start", messageId: "a", attempt: 1, processId: 1, at: 25 }, + ...execution({ messageId: "b", attempt: 1, startedAt: 30, finishedAt: 40 }), + ] + + expect(() => assertSerializedExecution(events, { messageCount: 2 })).toThrow(/finish/) + }) + + it("rejects a run that lost one of the messages", () => { + const events = [ + ...execution({ messageId: "a", attempt: 1, startedAt: 10, finishedAt: 20 }), + ...execution({ messageId: "a", attempt: 2, startedAt: 30, finishedAt: 40 }), + ] + + expect(() => assertSerializedExecution(events, { messageCount: 2 })).toThrow(/message/) + }) +})