From b28d7dbb3e1e05b952108681e4ab45ac717c60cd Mon Sep 17 00:00:00 2001 From: Kieran Klukas Date: Sat, 15 Aug 2026 08:39:00 -0400 Subject: [PATCH] fix: harden actor lifecycle recovery --- src/actor.ts | 66 ++++++++++++++++++++++++++---------- src/store.ts | 82 ++++++++++++++++++++++++++++++++++++++++++--- tests/actor.test.ts | 81 ++++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 207 insertions(+), 22 deletions(-) diff --git a/src/actor.ts b/src/actor.ts index eb78c28..f20b937 100644 --- a/src/actor.ts +++ b/src/actor.ts @@ -13,6 +13,7 @@ import { Event, type EventData, type EventName, + type MessageEndData, makeId, type ReasoningDeltaData, type ReasoningSignatureData, @@ -171,7 +172,7 @@ export type RunStep = * One single-writer conversation. Owns the durable event log (via Store), the * set of live SSE subscribers, the cancel flag, and the in-flight run. * - * Every `persist*` helper assigns the next monotonic seq, writes the event + * Every `persist*` helper allocates a monotonic seq in SQLite, writes the event * durably (synchronous, before fan-out), then emits it to all subscribers. * Deltas are batched (BATCH_MAX_DELTAS per persist) with an in-memory tail; * the tail is flushed to durable storage on a full batch, keeping resume @@ -186,23 +187,16 @@ export class ConversationActor { /** Aborts the in-flight provider stream; set only while a run is streaming. */ private currentAbort: (() => void) | null = null; private currentRunId: string | null = null; - private seq = 0; lastActivity = Date.now(); constructor(conversationId: string, store: Store) { this.conversationId = conversationId; this.store = store; - this.seq = this.store.lastSeq(conversationId); - } - - private nextSeq(): number { - return ++this.seq; } private persist(eventName: EventName, data: EventData): WireEvent { - const seq = this.nextSeq(); + const seq = this.store.appendAndBump(this.conversationId, eventName, data); const id = makeId(this.conversationId, seq); - this.store.appendAndBump(this.conversationId, seq, eventName, data); const e: WireEvent = { id, event: eventName, data }; this.emit(e); this.lastActivity = Date.now(); @@ -288,6 +282,9 @@ export class ConversationActor { * leak into a later run. */ requestCancel(runId?: string): void { + // This is deliberately durable as well as in-memory: the HTTP process may + // not be the process that owns the provider stream. + this.store.requestCancel(this.conversationId); if (runId !== undefined) { this.cancelRunId = runId; } else if (this.currentRunId !== null) { @@ -300,6 +297,9 @@ export class ConversationActor { // model can sit mid-token for seconds). The loop's own isCancelled check // still covers the between-steps case. if (this.currentAbort && this.isCancelled()) this.currentAbort(); + // A local in-flight run has already observed this request synchronously. + // Consume the durable copy so it cannot poison the following queued run. + if (this.currentRunId !== null) this.store.takeCancelRequest(this.conversationId); } isCancelled(): boolean { @@ -340,11 +340,13 @@ export class ConversationActor { const callIds = new Set(); const resultIds = new Set(); const signedMsgs = new Set(); + const completedMsgs = new Set(); for (const e of this.store.replay(this.conversationId, 0)) { if (e.event === Event.ToolCall) callIds.add((e.data as ToolCallData).toolCallId); else if (e.event === Event.ToolResult) resultIds.add((e.data as ToolResultData).toolCallId); else if (e.event === Event.ReasoningSig) signedMsgs.add((e.data as ReasoningSignatureData).messageId); + else if (e.event === Event.MsgEnd) completedMsgs.add((e.data as MessageEndData).messageId); } const paired = new Set([...callIds].filter((id) => resultIds.has(id))); @@ -424,13 +426,19 @@ export class ConversationActor { if (d.content.length > 0) userParts.push({ type: "text", text: d.content }); for (const a of d.attachments ?? []) userParts.push(await this.renderAttachment(a, opts)); } else if (e.event === Event.TextDelta) { + const d = e.data as TextDeltaData; + // A lease reclaim starts a fresh provider request. Deltas from an actor + // that died before message-end are therefore only a partial attempt, + // not prior assistant context to prepend to that retry. + if (!completedMsgs.has(d.messageId)) continue; flushUser(); if (toolResults.length > 0) flushTools(); // text after a result: new step - asstText += (e.data as TextDeltaData).delta; + asstText += d.delta; } else if (e.event === Event.ReasoningDelta) { + const d = e.data as ReasoningDeltaData; + if (!completedMsgs.has(d.messageId)) continue; flushUser(); if (toolResults.length > 0) flushTools(); // reasoning after a result: new step - const d = e.data as ReasoningDeltaData; // Accumulate only for turns that carry a signature — everything else is // dropped from history, so there's no need to build the string. if (signedMsgs.has(d.messageId)) asstReasoning += d.delta; @@ -509,7 +517,6 @@ export class ConversationActor { /** Replay the durable log strictly after `afterSeq`, oldest first. */ replay(afterSeq: number): WireEvent[] { - if (afterSeq >= this.seq) return []; return this.store.replay(this.conversationId, Math.max(0, afterSeq)).map((r) => ({ id: makeId(this.conversationId, r.seq), event: r.event as EventName, @@ -568,7 +575,8 @@ export class ConversationActor { ): Promise { this.currentRunId = runId; // Honor a cancel that arrived before this run was claimed. - if (this.cancelNext) { + const durableCancel = this.store.takeCancelRequest(this.conversationId); + if (this.cancelNext || durableCancel) { this.cancelRunId = runId; this.cancelNext = false; } @@ -590,11 +598,10 @@ export class ConversationActor { // long thinking stretch before the first answer token can't be reaped. const flush = (eventName: EventName, text: string): void => { if (text.length === 0) return; - const seq = this.nextSeq(); - const id = makeId(this.conversationId, seq); const truncated = truncateUtf8(text); const data = { runId, threadId: this.conversationId, messageId, delta: truncated }; - this.store.appendAndBump(this.conversationId, seq, eventName, data); + const seq = this.store.appendAndBump(this.conversationId, eventName, data); + const id = makeId(this.conversationId, seq); this.emit({ id, event: eventName, data }); this.lastActivity = Date.now(); onProgress?.(seq); @@ -611,9 +618,23 @@ export class ConversationActor { }; const tick = setInterval(() => { + if (this.store.isConversationDeleted(this.conversationId)) { + this.cancelRunId = runId; + abortController.abort(); + return; + } flushReasoning(); flushDelta(); }, BATCH_FLUSH_MS); + // The HTTP server and the worker can be separate processes. Poll the + // durable request while a provider is quiet so a cancel reaches whichever + // process owns this run, not only the actor that received the HTTP call. + const cancelTick = setInterval(() => { + if (this.store.takeCancelRequest(this.conversationId)) { + this.cancelRunId = runId; + abortController.abort(); + } + }, 250); let usage: TokenUsage | undefined; // Reasoning duration: wall-clock from the start of generation to the first @@ -625,6 +646,11 @@ export class ConversationActor { let firstTextAt = 0; try { for await (const step of steps(abortController.signal)) { + if (this.store.isConversationDeleted(this.conversationId)) { + this.cancelRunId = runId; + abortController.abort(); + break; + } if (this.isCancelled()) { abortController.abort(); break; @@ -698,7 +724,7 @@ export class ConversationActor { } catch (err) { // Aborting mid-token (via requestCancel) surfaces here as a rejection — // that's a clean stop, not an error. Anything else is a real failure. - if (!this.isCancelled()) { + if (!this.isCancelled() && !this.store.isConversationDeleted(this.conversationId)) { errored = true; // Also log server-side: the RunErr event reaches the client, but the // terminal is where you're watching, and a stack is more diagnostic. @@ -711,9 +737,14 @@ export class ConversationActor { } } finally { clearInterval(tick); + clearInterval(cancelTick); this.currentAbort = null; } + if (this.store.isConversationDeleted(this.conversationId)) { + this.currentRunId = null; + return; + } flushReasoning(); flushDelta(); const finish = errored ? "error" : this.isCancelled() ? "aborted" : "stop"; @@ -786,6 +817,7 @@ export class ConversationActor { phase: string; data?: unknown; }): void { + if (this.store.isConversationDeleted(this.conversationId)) return; this.persist(Event.ToolProgress, { threadId: this.conversationId, ...p }); } } diff --git a/src/store.ts b/src/store.ts index 4b5e3c9..1e4bd2f 100644 --- a/src/store.ts +++ b/src/store.ts @@ -310,6 +310,21 @@ CREATE TABLE IF NOT EXISTS jobs ( ); CREATE INDEX IF NOT EXISTS idx_jobs_status ON jobs (status); +-- A deletion is a durable tombstone, not merely the absence of a row. An +-- in-flight worker may still hold an actor after DELETE; the tombstone prevents +-- that actor from recreating the conversation while it unwinds. +CREATE TABLE IF NOT EXISTS deleted_conversations ( + id TEXT PRIMARY KEY, + deleted_at INTEGER NOT NULL +); + +-- Cancellation must cross the server/worker process boundary. One pending row +-- cancels the current run, or the next one when it has not yet been claimed. +CREATE TABLE IF NOT EXISTS cancel_requests ( + conversation_id TEXT PRIMARY KEY, + created_at INTEGER NOT NULL +); + CREATE TABLE IF NOT EXISTS model_settings ( model_ref TEXT PRIMARY KEY, visible INTEGER NOT NULL DEFAULT 0, @@ -509,6 +524,10 @@ export class Store { private findOrphanBlobsStmt: ReturnType; private deleteBlobStmt: ReturnType; private addBlobRefStmt: ReturnType; + private claimConversationOwnerStmt: ReturnType; + private isDeletedConversationStmt: ReturnType; + private requestCancelStmt: ReturnType; + private takeCancelStmt: ReturnType; constructor(databasePath: string = getConfig().server.dbPath) { // Ensure the parent directory exists so a fresh checkout (where `data/` is @@ -582,7 +601,8 @@ export class Store { ); this.upsertConversationStmt = this.db.prepare( `INSERT INTO conversations (id, created_at, last_seq) VALUES (?, ?, ?) - ON CONFLICT(id) DO UPDATE SET last_seq = excluded.last_seq`, + ON CONFLICT(id) DO UPDATE SET last_seq = conversations.last_seq + 1 + RETURNING last_seq`, ); this.enqueueStmt = this.db.prepare( @@ -719,6 +739,21 @@ export class Store { `INSERT INTO blob_refs (sha256, conversation_id) VALUES (?, ?) ON CONFLICT(sha256, conversation_id) DO NOTHING`, ); + this.claimConversationOwnerStmt = this.db.prepare( + `INSERT INTO conversations (id, created_at, last_seq, owner_sub) VALUES (?, ?, 0, ?) + ON CONFLICT(id) DO UPDATE SET owner_sub = COALESCE(conversations.owner_sub, excluded.owner_sub) + RETURNING owner_sub`, + ); + this.isDeletedConversationStmt = this.db.prepare( + "SELECT 1 FROM deleted_conversations WHERE id = ?", + ); + this.requestCancelStmt = this.db.prepare( + `INSERT INTO cancel_requests (conversation_id, created_at) VALUES (?, ?) + ON CONFLICT(conversation_id) DO UPDATE SET created_at = excluded.created_at`, + ); + this.takeCancelStmt = this.db.prepare( + "DELETE FROM cancel_requests WHERE conversation_id = ? RETURNING conversation_id", + ); } /** @@ -1130,6 +1165,9 @@ export class Store { this.db.prepare("DELETE FROM events WHERE conversation_id = ?").run(id); this.db.prepare("DELETE FROM jobs WHERE conversation_id = ?").run(id); this.db.prepare("DELETE FROM conversations WHERE id = ?").run(id); + this.db + .prepare("INSERT OR IGNORE INTO deleted_conversations (id, deleted_at) VALUES (?, ?)") + .run(id, Date.now()); // Of those candidates, the ones no conversation references anymore. const stillRef = this.db.prepare("SELECT 1 FROM blob_refs WHERE sha256 = ? LIMIT 1"); return candidates.filter((sha) => stillRef.get(sha) == null); @@ -1292,6 +1330,29 @@ export class Store { .run(sub, id); } + /** Atomically create or claim an unowned conversation for an authenticated user. */ + claimConversationOwner(id: string, sub: string): boolean { + const row = this.claimConversationOwnerStmt.get(id, Date.now(), sub) as + | { owner_sub: string } + | null; + return row?.owner_sub === sub; + } + + /** Whether this id was deleted and must never be revived by a stale actor. */ + isConversationDeleted(id: string): boolean { + return this.isDeletedConversationStmt.get(id) != null; + } + + /** Request cancellation durably so the process running the job can observe it. */ + requestCancel(conversationId: string): void { + this.requestCancelStmt.run(conversationId, Date.now()); + } + + /** Consume one cancellation request. A pre-claim cancel therefore reaches the next run. */ + takeCancelRequest(conversationId: string): boolean { + return this.takeCancelStmt.get(conversationId) != null; + } + /** The kloe user `sub` that owns a conversation, or undefined (auth-off / legacy). */ getConversationOwner(id: string): string | undefined { const row = this.db.query("SELECT owner_sub FROM conversations WHERE id = ?").get(id) as { @@ -1473,9 +1534,20 @@ export class Store { this.addBlobRefStmt.run(sha256, conversationId); } - /** Atomic append + seq advance in one transaction. */ - appendAndBump(conversationId: string, seq: number, eventName: string, data: unknown): void { - this.db.transaction(() => { + /** + * Atomic append + sequence allocation. The sequence is allocated in SQLite, + * not an actor-local counter, because the HTTP process may append a steer + * while a standalone worker is writing the assistant stream. + */ + appendAndBump(conversationId: string, eventName: string, data: unknown): number { + return this.db.transaction(() => { + if (this.isConversationDeleted(conversationId)) { + throw new Error(`conversation "${conversationId}" was deleted`); + } + const row = this.upsertConversationStmt.get(conversationId, Date.now(), 1) as { + last_seq: number; + }; + const seq = row.last_seq; const id = `${conversationId}:${seq}`; this.insertEventStmt.run( id, @@ -1485,7 +1557,7 @@ export class Store { JSON.stringify(data), Date.now(), ); - this.upsertConversationStmt.run(conversationId, Date.now(), seq); + return seq; })(); } diff --git a/tests/actor.test.ts b/tests/actor.test.ts index 39ca337..b928a3c 100644 --- a/tests/actor.test.ts +++ b/tests/actor.test.ts @@ -38,6 +38,15 @@ test("event ids are monotonic across a conversation", async () => { } }); +test("separate actors allocate one shared sequence without collisions", () => { + const first = new ConversationActor("t-shared-seq", store); + const second = new ConversationActor("t-shared-seq", store); + first.appendUser("from server", "r1"); + second.appendUser("from worker", "r2"); + + expect(store.replay("t-shared-seq", 0).map((e) => e.seq)).toEqual([1, 2]); +}); + test("a usage step is stamped onto the message-end event", async () => { const a = new ConversationActor("t-usage", store); const events: WireEvent[] = []; @@ -199,6 +208,78 @@ test("history drops an assistant turn that produced no text (stopped before firs expect(await a.history()).toEqual([{ role: "user", content: "hi" }]); }); +test("history drops durable deltas from a run that crashed before message-end", async () => { + const conv = "t-history-crash"; + const a = new ConversationActor(conv, store); + a.appendUser("retry this", "r1"); + store.appendAndBump(conv, Event.RunStart, { threadId: conv, runId: "r1", messageId: "m1" }); + store.appendAndBump(conv, Event.MsgStart, { + threadId: conv, + runId: "r1", + messageId: "m1", + type: "text", + }); + store.appendAndBump(conv, Event.TextDelta, { + threadId: conv, + runId: "r1", + messageId: "m1", + delta: "partial answer", + }); + + expect(await a.history()).toEqual([{ role: "user", content: "retry this" }]); +}); + +test("a durable cancel reaches the process that claims the next run", async () => { + const conv = "t-durable-cancel"; + const a = new ConversationActor(conv, store); + store.requestCancel(conv); + const events: WireEvent[] = []; + a.follow({ push: (e) => events.push(e), closed: false }); + + await a.runText("r1", "m1", async function* (_signal) { + yield { kind: "text", chunk: "never" }; + }); + + expect(events.some((e) => e.event === Event.Cancelled)).toBe(true); + expect(events.some((e) => e.event === Event.TextDelta)).toBe(false); +}); + +test("a deleted conversation cannot be recreated by its stale actor", () => { + const conv = "t-deleted"; + const a = new ConversationActor(conv, store); + a.appendUser("before delete", "r1"); + store.deleteConversation(conv); + + expect(() => a.appendUser("late write", "r2")).toThrow(`conversation "${conv}" was deleted`); + expect(store.replay(conv, 0)).toEqual([]); +}); + +test("deletion stops an in-flight actor without writing a terminal event", async () => { + const conv = "t-delete-running"; + const a = new ConversationActor(conv, store); + a.appendUser("start", "r1"); + let reached: () => void = () => {}; + const started = new Promise((resolve) => (reached = resolve)); + const run = a.runText("r1", "m1", async function* (signal) { + yield { kind: "text", chunk: "partial" }; + reached(); + await new Promise((_resolve, reject) => { + signal.addEventListener("abort", () => reject(new DOMException("aborted", "AbortError"))); + }); + }); + + await started; + store.deleteConversation(conv); + await run; + expect(store.replay(conv, 0)).toEqual([]); +}); + +test("only the first authenticated writer can claim a new conversation", () => { + expect(store.claimConversationOwner("t-owner-race", "alice")).toBe(true); + expect(store.claimConversationOwner("t-owner-race", "bob")).toBe(false); + expect(store.getConversationOwner("t-owner-race")).toBe("alice"); +}); + test("cancel mid-token aborts the stream at once and finishes aborted, not error", async () => { const a = new ConversationActor("t-abort", store); const events: WireEvent[] = []; -- 2.51.2