Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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"),
};
Expand Down Expand Up @@ -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"),
};
Expand Down Expand Up @@ -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"),
};
Expand Down
130 changes: 130 additions & 0 deletions apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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* () {
Expand Down
30 changes: 20 additions & 10 deletions apps/server/src/orchestration/Layers/ProjectionPipeline.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, OrchestrationEvent>();
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);
}
Expand All @@ -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),
Expand Down
106 changes: 106 additions & 0 deletions apps/server/src/persistence/Layers/OrchestrationEventStore.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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* () {
Expand Down
Loading
Loading