From f93ac23a3d9f22b12703d73f1081c2363f258057 Mon Sep 17 00:00:00 2001 From: ThomasCrund <39323104+ThomasCrund@users.noreply.github.com> Date: Fri, 11 Sep 2026 12:15:20 +1000 Subject: [PATCH] fix(server): bootstrap attachment cleanup without decoding the event log The attachment-cleanup bootstrap added in #9871 replayed every event after its cursor and filtered for thread.deleted and thread.reverted in JS. Even after #10777 released consumed pages, a single 500-row page can hold over a gigabyte of serialized payload (historical session records carrying a 73 MB lastError), so the backend still exhausts its heap before the cursor moves and the desktop app crash-loops at startup on affected profiles. Cleanup now reads only deleted and reverted events through a new type-filtered event store query, so unrelated rows are never decoded. The head sequence is captured before replay and the cleanup cursor advances to it after a successful pass, including when no cleanup events exist. Ordering, the last-wins dedupe per thread, the recreate check, and the retry-on-failure cursor behavior are unchanged. Regression tests append a caught-up history containing an undecodable row of another type and assert bootstrap succeeds, removes the deleted thread's attachment, and lands the cursor on the head; the event store tests cover the typed reader's range and type filtering and the head query. Done with Claude Fable 5.1 in Claude Code. Co-Authored-By: Claude Fable 5.1 --- .../Layers/OrchestrationEngine.test.ts | 6 + .../Layers/ProjectionPipeline.test.ts | 130 ++++++++++++++++++ .../Layers/ProjectionPipeline.ts | 30 ++-- .../Layers/OrchestrationEventStore.test.ts | 106 ++++++++++++++ .../Layers/OrchestrationEventStore.ts | 99 +++++++++++++ .../Services/OrchestrationEventStore.ts | 23 ++++ 6 files changed, 384 insertions(+), 10 deletions(-) diff --git a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts index 6d062ebcab4b..b9dcf28696ed 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts @@ -353,6 +353,8 @@ describe("OrchestrationEngine", () => { }), ), hasEventAfter: () => Effect.succeed(false), + readEventsOfTypes: () => Stream.empty, + getHead: () => Effect.succeedNone, readAggregateRange: () => Stream.die("unused aggregate replay"), getAggregateReplayStats: () => Effect.die("unused aggregate replay stats"), }; @@ -1498,6 +1500,8 @@ describe("OrchestrationEngine", () => { return Stream.fromIterable(events); }, hasEventAfter: () => Effect.succeed(false), + readEventsOfTypes: () => Stream.empty, + getHead: () => Effect.succeedNone, readAggregateRange: () => Stream.die("unused aggregate replay"), getAggregateReplayStats: () => Effect.die("unused aggregate replay stats"), }; @@ -1738,6 +1742,8 @@ describe("OrchestrationEngine", () => { return Stream.fromIterable(events); }, hasEventAfter: () => Effect.succeed(false), + readEventsOfTypes: () => Stream.empty, + getHead: () => Effect.succeedNone, readAggregateRange: () => Stream.die("unused aggregate replay"), getAggregateReplayStats: () => Effect.die("unused aggregate replay stats"), }; diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts index a038a5662169..0f669efaf5bb 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts @@ -2068,6 +2068,136 @@ it.layer(Layer.fresh(makeProjectionPipelinePrefixedTestLayer("t3-projection-atta }, ); +it.layer(Layer.fresh(makeProjectionPipelinePrefixedTestLayer("t3-projection-attachments-typed-")))( + "OrchestrationProjectionPipeline", + (it) => { + it.effect("cleanup bootstrap never decodes unrelated history and advances to the head", () => + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const sql = yield* SqlClient.SqlClient; + const projectionPipeline = yield* OrchestrationProjectionPipeline; + const eventStore = yield* OrchestrationEventStore; + const projectionState = yield* ProjectionStateRepository; + const { attachmentsDir } = yield* ServerConfig; + const now = "2026-01-01T00:00:00.000Z"; + const projectId = ProjectId.make("project-typed"); + const threadId = ThreadId.make("thread-typed-gone"); + const attachmentPath = path.join( + attachmentsDir, + "thread-typed-gone-00000000-0000-4000-8000-000000000001.png", + ); + + yield* eventStore.append({ + type: "project.created", + eventId: EventId.make("evt-typed-project"), + aggregateKind: "project", + aggregateId: projectId, + occurredAt: now, + commandId: CommandId.make("cmd-typed-project"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-typed-project"), + metadata: {}, + payload: { + projectId, + title: "Typed", + workspaceRoot: "/tmp/project-typed", + defaultModelSelection: null, + scripts: [], + createdAt: now, + updatedAt: now, + }, + }); + yield* eventStore.append({ + type: "thread.created", + eventId: EventId.make("evt-typed-create"), + aggregateKind: "thread", + aggregateId: threadId, + occurredAt: now, + commandId: CommandId.make("cmd-typed-create"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-typed-create"), + metadata: {}, + payload: { + threadId, + projectId, + title: "Thread typed", + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + runtimeMode: "full-access", + branch: null, + worktreePath: null, + createdAt: now, + updatedAt: now, + }, + }); + yield* eventStore.append({ + type: "thread.deleted", + eventId: EventId.make("evt-typed-delete"), + aggregateKind: "thread", + aggregateId: threadId, + occurredAt: now, + commandId: CommandId.make("cmd-typed-delete"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-typed-delete"), + metadata: {}, + payload: { threadId, deletedAt: now }, + }); + // A caught-up profile whose history holds a row the schema cannot decode. + // It stands in for the oversized records a full replay would materialize: + // cleanup must reach the head without ever reading it. + const headAt = "2026-01-01T00:00:01.000Z"; + yield* sql` + INSERT INTO orchestration_events ( + event_id, aggregate_kind, stream_id, stream_version, event_type, occurred_at, + actor_kind, payload_json, metadata_json + ) VALUES ( + 'evt-typed-broken', 'thread', 'thread-typed-broken', 0, 'thread.message-sent', + ${headAt}, 'provider', '{', '{}' + ) + `; + const head = Option.getOrThrow(yield* eventStore.getHead()); + yield* projectionState.upsertMany( + Object.values(ORCHESTRATION_PROJECTOR_NAMES).map((projector) => ({ + projector, + lastAppliedSequence: head.sequence, + updatedAt: headAt, + })), + ); + yield* fileSystem.makeDirectory(attachmentsDir, { recursive: true }); + yield* fileSystem.writeFileString(attachmentPath, "gone"); + + yield* projectionPipeline.bootstrap; + + assert.isFalse(yield* exists(attachmentPath)); + const cleanupCursor = projectionState.getByProjector({ + projector: "projection.attachment-cleanup", + }); + assert.deepEqual( + yield* cleanupCursor, + Option.some({ + projector: "projection.attachment-cleanup", + lastAppliedSequence: head.sequence, + updatedAt: headAt, + }), + ); + + yield* projectionPipeline.bootstrap; + assert.deepEqual( + yield* cleanupCursor, + Option.some({ + projector: "projection.attachment-cleanup", + lastAppliedSequence: head.sequence, + updatedAt: headAt, + }), + ); + }), + ); + }, +); + it.layer(BaseTestLayer)("OrchestrationProjectionPipeline", (it) => { it.effect("replays a bootstrap backlog larger than the event store default limit", () => Effect.gen(function* () { diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts index e303e7323729..fa7859e981e0 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts @@ -2105,17 +2105,29 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti lastAppliedSequence: cleanupStart, updatedAt: cleanupState?.updatedAt ?? "1970-01-01T00:00:00.000Z", }); + // Nothing appends until the engine starts its worker, so the head captured here + // is what the projectors replay through. Capturing it first keeps the cleanup + // cursor at or behind every projector cursor. + const head = yield* eventStore.getHead(); yield* Effect.forEach(projectors, bootstrapProjector, { concurrency: 1, discard: true }); + if (Option.isNone(head)) { + return; + } // Cleanup has its own cursor so retries never have to replay committed text. // All message and activity references are current before any files are removed. + // Only deletions and reverts matter here, and they are filtered in SQL: a full + // replay would decode every historical payload, and one page of oversized + // records is enough to exhaust the heap before the cursor can move. const pendingCleanup = new Map(); - let lastEvent: OrchestrationEvent | undefined; yield* Stream.runForEach( - eventStore.readFromSequence(cleanupStart, Number.MAX_SAFE_INTEGER), + eventStore.readEventsOfTypes({ + types: ["thread.reverted", "thread.deleted"], + fromSequenceExclusive: cleanupStart, + toSequenceInclusive: head.value.sequence, + }), (event) => Effect.sync(() => { - lastEvent = event; if (event.type === "thread.reverted" || event.type === "thread.deleted") { pendingCleanup.set(`${event.type}:${event.payload.threadId}`, event); } @@ -2133,13 +2145,11 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti // Leave the cleanup cursor behind this event so the next bootstrap retries it. if (!cleaned) return; } - if (lastEvent) { - yield* projectionStateRepository.upsert({ - projector: cleanupProjector, - lastAppliedSequence: lastEvent.sequence, - updatedAt: lastEvent.occurredAt, - }); - } + yield* projectionStateRepository.upsert({ + projector: cleanupProjector, + lastAppliedSequence: head.value.sequence, + updatedAt: head.value.occurredAt, + }); }).pipe( Effect.provideService(FileSystem.FileSystem, fileSystem), Effect.provideService(Path.Path, path), diff --git a/apps/server/src/persistence/Layers/OrchestrationEventStore.test.ts b/apps/server/src/persistence/Layers/OrchestrationEventStore.test.ts index 2287dd10af45..e8c24daa6fa2 100644 --- a/apps/server/src/persistence/Layers/OrchestrationEventStore.test.ts +++ b/apps/server/src/persistence/Layers/OrchestrationEventStore.test.ts @@ -11,6 +11,7 @@ import { import { assert, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; import * as Schema from "effect/Schema"; import * as Stream from "effect/Stream"; import * as SqlClient from "effect/unstable/sql/SqlClient"; @@ -315,8 +316,113 @@ layer("OrchestrationEventStore", (it) => { assert.deepEqual(yield* Stream.runCollect(store.readFromSequence(0, -1)), []); }), ); + + it.effect("reads only the requested event types through the captured head", () => + Effect.gen(function* () { + const store = yield* OrchestrationEventStore; + const sql = yield* SqlClient.SqlClient; + const now = "2026-01-01T00:00:00.000Z"; + const deletedThreadId = ThreadId.make("typed-deleted"); + const revertedThreadId = ThreadId.make("typed-reverted"); + const base = { + aggregateKind: "thread" as const, + occurredAt: now, + commandId: null, + causationEventId: null, + correlationId: null, + metadata: {}, + }; + const before = yield* store.append(messageEvent(deletedThreadId, "typed-before")); + const deleted = yield* store.append({ + ...base, + type: "thread.deleted", + eventId: EventId.make("typed-deleted-1"), + aggregateId: deletedThreadId, + payload: { threadId: deletedThreadId, deletedAt: now }, + }); + // An undecodable payload of another type must never be touched by a typed read. + yield* sql` + INSERT INTO orchestration_events ( + event_id, aggregate_kind, stream_id, stream_version, event_type, occurred_at, + actor_kind, payload_json, metadata_json + ) VALUES ( + 'typed-broken', 'thread', 'typed-broken-thread', 0, 'thread.message-sent', + ${now}, 'provider', '{', '{}' + ) + `; + const reverted = yield* store.append({ + ...base, + type: "thread.reverted", + eventId: EventId.make("typed-reverted-1"), + aggregateId: revertedThreadId, + payload: { threadId: revertedThreadId, turnCount: 1 }, + }); + const headAt = "2026-01-01T00:00:01.000Z"; + const afterHead = yield* store.append({ + ...base, + type: "thread.deleted", + eventId: EventId.make("typed-deleted-2"), + aggregateId: revertedThreadId, + occurredAt: headAt, + payload: { threadId: revertedThreadId, deletedAt: headAt }, + }); + + assert.deepEqual( + yield* store.getHead(), + Option.some({ sequence: afterHead.sequence, occurredAt: headAt }), + ); + + const types = ["thread.deleted", "thread.reverted"] as const; + const replayed = yield* Stream.runCollect( + store.readEventsOfTypes({ + types, + fromSequenceExclusive: before.sequence, + toSequenceInclusive: reverted.sequence, + }), + ); + assert.deepEqual( + replayed.map((event) => [event.sequence, event.type]), + [ + [deleted.sequence, "thread.deleted"], + [reverted.sequence, "thread.reverted"], + ], + ); + assert.deepEqual( + yield* Stream.runCollect( + store.readEventsOfTypes({ + types, + fromSequenceExclusive: reverted.sequence, + toSequenceInclusive: reverted.sequence, + }), + ), + [], + ); + assert.deepEqual( + yield* Stream.runCollect( + store.readEventsOfTypes({ + types: [], + fromSequenceExclusive: 0, + toSequenceInclusive: afterHead.sequence, + }), + ), + [], + ); + // The broken row is still there for a full replay to trip over. + const fullReplay = yield* Stream.runCollect(store.readFromSequence(before.sequence)).pipe( + Effect.flip, + ); + assert.isTrue(isPersistenceDecodeError(fullReplay)); + }), + ); }); +it.effect("getHead is none for an empty store", () => + Effect.gen(function* () { + const store = yield* OrchestrationEventStore; + assert.deepEqual(yield* store.getHead(), Option.none()); + }).pipe(Effect.provide(OrchestrationEventStoreLive.pipe(Layer.provide(SqlitePersistenceMemory)))), +); + for (const reader of ["all", "aggregate"] as const) { it.effect(`releases consumed pages during ${reader} replay`, () => Effect.gen(function* () { diff --git a/apps/server/src/persistence/Layers/OrchestrationEventStore.ts b/apps/server/src/persistence/Layers/OrchestrationEventStore.ts index a13b4f77e39e..a5405bf915d0 100644 --- a/apps/server/src/persistence/Layers/OrchestrationEventStore.ts +++ b/apps/server/src/persistence/Layers/OrchestrationEventStore.ts @@ -72,6 +72,16 @@ const ReadFromSequenceRequestSchema = Schema.Struct({ sequenceExclusive: NonNegativeInt, limit: Schema.Number, }); +const ReadEventsOfTypesRequestSchema = Schema.Struct({ + types: Schema.Array(OrchestrationEventType), + fromSequenceExclusive: NonNegativeInt, + toSequenceInclusive: NonNegativeInt, + limit: Schema.Number, +}); +const EventHeadRowSchema = Schema.Struct({ + sequence: NonNegativeInt, + occurredAt: IsoDateTime, +}); const AggregateReplayRequestSchema = Schema.Struct({ aggregateKind: OrchestrationAggregateKind, aggregateId: Schema.String, @@ -228,6 +238,44 @@ const makeEventStore = Effect.gen(function* () { `, }); + const readEventRowsOfTypes = SqlSchema.findAll({ + Request: ReadEventsOfTypesRequestSchema, + Result: OrchestrationEventPersistedRowSchema, + execute: (request) => + sql` + SELECT + sequence, + event_id AS "eventId", + event_type AS "type", + aggregate_kind AS "aggregateKind", + stream_id AS "aggregateId", + occurred_at AS "occurredAt", + command_id AS "commandId", + causation_event_id AS "causationEventId", + correlation_id AS "correlationId", + payload_json AS "payload", + metadata_json AS "metadata" + FROM orchestration_events + WHERE sequence > ${request.fromSequenceExclusive} + AND sequence <= ${request.toSequenceInclusive} + AND ${sql.in("event_type", request.types)} + ORDER BY sequence ASC + LIMIT ${request.limit} + `, + }); + + const readEventHeadRow = SqlSchema.findOneOption({ + Request: Schema.Void, + Result: EventHeadRowSchema, + execute: () => + sql` + SELECT sequence, occurred_at AS "occurredAt" + FROM orchestration_events + ORDER BY sequence DESC + LIMIT 1 + `, + }); + const readAggregateReplayStats = SqlSchema.findOne({ Request: AggregateReplayRequestSchema, Result: AggregateReplayStatsRowSchema, @@ -323,6 +371,55 @@ const makeEventStore = Effect.gen(function* () { ); }; + const readEventsOfTypes: OrchestrationEventStoreShape["readEventsOfTypes"] = (input) => { + if (input.types.length === 0 || input.fromSequenceExclusive >= input.toSequenceInclusive) { + return Stream.empty; + } + return Stream.paginate(input.fromSequenceExclusive, (cursor) => + readEventRowsOfTypes({ + types: input.types, + fromSequenceExclusive: cursor, + toSequenceInclusive: input.toSequenceInclusive, + limit: READ_PAGE_SIZE, + }).pipe( + Effect.mapError( + toPersistenceSqlOrDecodeError( + "OrchestrationEventStore.readEventsOfTypes:query", + "OrchestrationEventStore.readEventsOfTypes:decodeRows", + ), + ), + Effect.flatMap((rows) => + Effect.forEach(rows, (row) => + decodeEvent(row).pipe( + Effect.mapError( + toPersistenceDecodeError("OrchestrationEventStore.readEventsOfTypes:rowToEvent"), + ), + ), + ), + ), + Effect.map((events) => { + const last = events.at(-1); + return [ + events, + last === undefined || events.length < READ_PAGE_SIZE + ? Option.none() + : Option.some(last.sequence), + ] as const; + }), + ), + ); + }; + + const getHead: OrchestrationEventStoreShape["getHead"] = () => + readEventHeadRow().pipe( + Effect.mapError( + toPersistenceSqlOrDecodeError( + "OrchestrationEventStore.getHead:query", + "OrchestrationEventStore.getHead:decodeRow", + ), + ), + ); + const findEventAfter = SqlSchema.findOneOption({ Request: HasEventAfterRequestSchema, Result: Schema.Struct({ sequence: Schema.Number }), @@ -414,6 +511,8 @@ const makeEventStore = Effect.gen(function* () { return { append, readFromSequence, + readEventsOfTypes, + getHead, readAggregateRange, getAggregateReplayStats, readAll: () => readFromSequence(0, Number.MAX_SAFE_INTEGER), diff --git a/apps/server/src/persistence/Services/OrchestrationEventStore.ts b/apps/server/src/persistence/Services/OrchestrationEventStore.ts index 3e131b068b01..8f504861e745 100644 --- a/apps/server/src/persistence/Services/OrchestrationEventStore.ts +++ b/apps/server/src/persistence/Services/OrchestrationEventStore.ts @@ -12,6 +12,7 @@ import { OrchestrationEvent } from "@t3tools/contracts"; import * as Context from "effect/Context"; import type * as Effect from "effect/Effect"; +import type * as Option from "effect/Option"; import type * as Stream from "effect/Stream"; import type { OrchestrationEventStoreError } from "../Errors.ts"; @@ -30,6 +31,11 @@ export interface OrchestrationAggregateReplayStats { readonly hasCreateEvent: boolean; } +export interface OrchestrationEventHead { + readonly sequence: number; + readonly occurredAt: string; +} + /** * OrchestrationEventStoreShape - Service API for orchestration event persistence. */ @@ -60,6 +66,23 @@ export interface OrchestrationEventStoreShape { limit?: number, ) => Stream.Stream; + /** + * Replay only the listed event types after a sequence, through a captured + * head. Rows of other types are filtered in SQL and never decoded, so a + * history full of oversized payloads costs nothing here. + */ + readonly readEventsOfTypes: (input: { + readonly types: ReadonlyArray; + readonly fromSequenceExclusive: number; + readonly toSequenceInclusive: number; + }) => Stream.Stream; + + /** The newest event's sequence and time, without decoding its payload. */ + readonly getHead: () => Effect.Effect< + Option.Option, + OrchestrationEventStoreError + >; + /** Read one aggregate through a captured global head, without decoding other streams. */ readonly readAggregateRange: ( input: OrchestrationAggregateReplayRange & { readonly limit?: number },