From 623ea134f69ebef72f0922576a40e2393f9543b5 Mon Sep 17 00:00:00 2001 From: sak0a Date: Sun, 30 Aug 2026 03:53:54 +0200 Subject: [PATCH] fix(server): replay every unapplied projection event Override the event-store default during bootstrap and cover backlogs larger than one thousand events. --- .../Layers/ProjectionPipeline.test.ts | 71 +++++++++++++++++++ .../Layers/ProjectionPipeline.ts | 34 ++++----- 2 files changed, 89 insertions(+), 16 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts index e23de9bdd..52f9203e8 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts @@ -1337,6 +1337,77 @@ it.layer( }); it.layer(BaseTestLayer)("OrchestrationProjectionPipeline", (it) => { + it.effect("replays a bootstrap backlog larger than the event store default limit", () => + Effect.gen(function* () { + const projectionPipeline = yield* OrchestrationProjectionPipeline; + const eventStore = yield* OrchestrationEventStore; + const sql = yield* SqlClient.SqlClient; + const now = "2026-01-01T00:00:00.000Z"; + const projectId = ProjectId.make("project-bootstrap-backlog"); + + const sequenceRows = yield* sql<{ readonly maxSequence: number | null }>` + SELECT MAX(sequence) AS "maxSequence" FROM orchestration_events + `; + const sequenceBeforeBacklog = sequenceRows[0]?.maxSequence ?? 0; + const appendedEvents = yield* Effect.forEach( + Array.from({ length: 1_001 }, (_, index) => index), + (index) => { + const eventId = EventId.make(`evt-bootstrap-backlog-${index}`); + const commandId = CommandId.make(`cmd-bootstrap-backlog-${index}`); + return eventStore.append({ + type: "project.created", + eventId, + aggregateKind: "project", + aggregateId: projectId, + occurredAt: now, + commandId, + causationEventId: null, + correlationId: CorrelationId.make(commandId), + metadata: {}, + payload: { + projectId, + title: `Bootstrap backlog ${index}`, + workspaceRoot: "/tmp/project-bootstrap-backlog", + defaultModelSelection: null, + scripts: [], + createdAt: now, + updatedAt: now, + }, + }); + }, + ); + const lastSequence = appendedEvents[appendedEvents.length - 1]!.sequence; + + yield* Effect.forEach( + Object.values(ORCHESTRATION_PROJECTOR_NAMES), + (projector) => { + const lastAppliedSequence = + projector === ORCHESTRATION_PROJECTOR_NAMES.projects + ? sequenceBeforeBacklog + : lastSequence; + return sql` + INSERT INTO projection_state (projector, last_applied_sequence, updated_at) + VALUES (${projector}, ${lastAppliedSequence}, ${now}) + ON CONFLICT (projector) + DO UPDATE SET + last_applied_sequence = excluded.last_applied_sequence, + updated_at = excluded.updated_at + `; + }, + { discard: true }, + ); + + yield* projectionPipeline.bootstrap; + + const stateRows = yield* sql<{ readonly lastAppliedSequence: number }>` + SELECT last_applied_sequence AS "lastAppliedSequence" + FROM projection_state + WHERE projector = ${ORCHESTRATION_PROJECTOR_NAMES.projects} + `; + assert.deepEqual(stateRows, [{ lastAppliedSequence: lastSequence }]); + }), + ); + it.effect("resumes from projector last_applied_sequence without replaying older events", () => Effect.gen(function* () { const projectionPipeline = yield* OrchestrationProjectionPipeline; diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts index 5b036b499..269829da8 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts @@ -2028,22 +2028,24 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti } const replayFromSequence = Math.min(...lastAppliedByProjector.values()); - yield* Stream.runForEach(eventStore.readFromSequence(replayFromSequence), (event) => - Effect.gen(function* () { - const projectorsToAdvance = projectors.filter( - (projector) => event.sequence > (lastAppliedByProjector.get(projector.name) ?? 0), - ); - if (projectorsToAdvance.length === 0) { - return; - } - const postCommit = yield* sql.withTransaction( - applyProjectorsForEvent(event, projectorsToAdvance), - ); - for (const projector of projectorsToAdvance) { - lastAppliedByProjector.set(projector.name, event.sequence); - } - yield* postCommit; - }), + yield* Stream.runForEach( + eventStore.readFromSequence(replayFromSequence, Number.MAX_SAFE_INTEGER), + (event) => + Effect.gen(function* () { + const projectorsToAdvance = projectors.filter( + (projector) => event.sequence > (lastAppliedByProjector.get(projector.name) ?? 0), + ); + if (projectorsToAdvance.length === 0) { + return; + } + const postCommit = yield* sql.withTransaction( + applyProjectorsForEvent(event, projectorsToAdvance), + ); + for (const projector of projectorsToAdvance) { + lastAppliedByProjector.set(projector.name, event.sequence); + } + yield* postCommit; + }), ); });