Skip to content

Commit 6c8e00d

Browse files
extocijuliusmarmingeclaude
committed
perf(server): stop decoding unrelated events during startup (#12846)
Co-authored-by: Julius Marminge <51714798+juliusmarminge@users.noreply.github.com> Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
1 parent 93fa784 commit 6c8e00d

2 files changed

Lines changed: 258 additions & 21 deletions

File tree

Lines changed: 167 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,167 @@
1+
import * as NodeServices from "@effect/platform-node/NodeServices";
2+
import { assert, it } from "@effect/vitest";
3+
import {
4+
CommandId,
5+
EventId,
6+
ProjectId,
7+
ProviderInstanceId,
8+
ThreadId,
9+
type OrchestrationEvent,
10+
} from "@t3tools/contracts";
11+
import * as Effect from "effect/Effect";
12+
import * as FileSystem from "effect/FileSystem";
13+
import * as Layer from "effect/Layer";
14+
import * as Path from "effect/Path";
15+
16+
import { ServerConfig } from "../../config.ts";
17+
import { OrchestrationEventStoreLive } from "../../persistence/Layers/OrchestrationEventStore.ts";
18+
import { SqlitePersistenceMemory } from "../../persistence/Layers/Sqlite.ts";
19+
import { OrchestrationEventStore } from "../../persistence/Services/OrchestrationEventStore.ts";
20+
import { ProjectionStateRepository } from "../../persistence/Services/ProjectionState.ts";
21+
import { OrchestrationProjectionPipeline } from "../Services/ProjectionPipeline.ts";
22+
import { OrchestrationProjectionPipelineLive } from "./ProjectionPipeline.ts";
23+
24+
const TestLayer = OrchestrationProjectionPipelineLive.pipe(
25+
Layer.provideMerge(OrchestrationEventStoreLive),
26+
Layer.provideMerge(ServerConfig.layerTest(process.cwd(), { prefix: "t3-projection-cleanup-" })),
27+
Layer.provideMerge(SqlitePersistenceMemory),
28+
Layer.provideMerge(NodeServices.layer),
29+
);
30+
31+
const CLEANUP_PROJECTOR = "projection.attachment-cleanup";
32+
33+
const exists = (filePath: string) =>
34+
Effect.gen(function* () {
35+
const fileSystem = yield* FileSystem.FileSystem;
36+
return (yield* Effect.result(fileSystem.stat(filePath)))._tag === "Success";
37+
});
38+
39+
it.layer(TestLayer)("OrchestrationProjectionPipeline attachment cleanup bootstrap", (it) => {
40+
it.effect(
41+
"advances the cleanup cursor to the latest event, dedupes reverts, and holds on failure",
42+
() =>
43+
Effect.gen(function* () {
44+
const projectionPipeline = yield* OrchestrationProjectionPipeline;
45+
const eventStore = yield* OrchestrationEventStore;
46+
const projectionState = yield* ProjectionStateRepository;
47+
const fileSystem = yield* FileSystem.FileSystem;
48+
const path = yield* Path.Path;
49+
const { attachmentsDir } = yield* ServerConfig;
50+
const now = "2026-01-01T00:00:00.000Z";
51+
const projectId = ProjectId.make("project-cleanup");
52+
const revertedThreadId = ThreadId.make("thread-cleanup-revert");
53+
const deletedThreadId = ThreadId.make("thread-cleanup-delete");
54+
const revertedAttachmentPath = path.join(
55+
attachmentsDir,
56+
"thread-cleanup-revert-00000000-0000-4000-8000-000000000001.png",
57+
);
58+
const deletedAttachmentPath = path.join(
59+
attachmentsDir,
60+
"thread-cleanup-delete-00000000-0000-4000-8000-000000000001.png",
61+
);
62+
const cleanupCursor = projectionState
63+
.listAll()
64+
.pipe(
65+
Effect.map(
66+
(states) =>
67+
states.find((state) => state.projector === CLEANUP_PROJECTOR)?.lastAppliedSequence,
68+
),
69+
);
70+
71+
let eventOrdinal = 0;
72+
const append = (
73+
event: Pick<OrchestrationEvent, "type" | "aggregateKind" | "aggregateId" | "payload">,
74+
) => {
75+
eventOrdinal += 1;
76+
return eventStore.append({
77+
...event,
78+
eventId: EventId.make(`evt-cleanup-${eventOrdinal}`),
79+
occurredAt: now,
80+
commandId: CommandId.make(`cmd-cleanup-${eventOrdinal}`),
81+
causationEventId: null,
82+
correlationId: CommandId.make(`cmd-cleanup-${eventOrdinal}`),
83+
metadata: {},
84+
} as Omit<OrchestrationEvent, "sequence">);
85+
};
86+
const appendThreadCreated = (threadId: ThreadId) =>
87+
append({
88+
type: "thread.created",
89+
aggregateKind: "thread",
90+
aggregateId: threadId,
91+
payload: {
92+
threadId,
93+
projectId,
94+
title: "Cleanup thread",
95+
modelSelection: {
96+
instanceId: ProviderInstanceId.make("codex"),
97+
model: "gpt-5-codex",
98+
},
99+
runtimeMode: "full-access",
100+
branch: null,
101+
worktreePath: null,
102+
createdAt: now,
103+
updatedAt: now,
104+
},
105+
});
106+
const appendReverted = (turnCount: number) =>
107+
append({
108+
type: "thread.reverted",
109+
aggregateKind: "thread",
110+
aggregateId: revertedThreadId,
111+
payload: { threadId: revertedThreadId, turnCount },
112+
});
113+
114+
yield* append({
115+
type: "project.created",
116+
aggregateKind: "project",
117+
aggregateId: projectId,
118+
payload: {
119+
projectId,
120+
title: "Cleanup project",
121+
workspaceRoot: "/tmp/project-cleanup",
122+
defaultModelSelection: null,
123+
scripts: [],
124+
createdAt: now,
125+
updatedAt: now,
126+
},
127+
});
128+
yield* appendThreadCreated(revertedThreadId);
129+
yield* appendThreadCreated(deletedThreadId);
130+
yield* fileSystem.makeDirectory(attachmentsDir, { recursive: true });
131+
yield* fileSystem.writeFileString(revertedAttachmentPath, "stale");
132+
// Two reverts for one thread dedupe into a single prune.
133+
yield* appendReverted(2);
134+
yield* appendReverted(1);
135+
// The latest event is neither a revert nor a delete, yet the cursor must land on it.
136+
const latest = yield* append({
137+
type: "thread.archived",
138+
aggregateKind: "thread",
139+
aggregateId: deletedThreadId,
140+
payload: { threadId: deletedThreadId, archivedAt: now, updatedAt: now },
141+
});
142+
143+
yield* projectionPipeline.bootstrap;
144+
assert.isFalse(yield* exists(revertedAttachmentPath));
145+
assert.equal(yield* cleanupCursor, latest.sequence);
146+
147+
// A cleanup that cannot remove its files leaves the cursor behind for the next bootstrap.
148+
yield* fileSystem.makeDirectory(deletedAttachmentPath);
149+
yield* fileSystem.writeFileString(path.join(deletedAttachmentPath, "keep.txt"), "keep");
150+
yield* append({
151+
type: "thread.deleted",
152+
aggregateKind: "thread",
153+
aggregateId: deletedThreadId,
154+
payload: { threadId: deletedThreadId, deletedAt: now },
155+
});
156+
yield* projectionPipeline.bootstrap;
157+
assert.isTrue(yield* exists(deletedAttachmentPath));
158+
assert.equal(yield* cleanupCursor, latest.sequence);
159+
160+
yield* fileSystem.remove(deletedAttachmentPath, { recursive: true });
161+
yield* fileSystem.writeFileString(deletedAttachmentPath, "retry");
162+
yield* projectionPipeline.bootstrap;
163+
assert.isFalse(yield* exists(deletedAttachmentPath));
164+
assert.equal(yield* cleanupCursor, latest.sequence + 1);
165+
}),
166+
);
167+
});

‎apps/server/src/orchestration/Layers/ProjectionPipeline.ts‎

Lines changed: 91 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
11
import {
22
ApprovalRequestId,
3+
IsoDateTime,
34
isImportedAgentSessionMessageId,
5+
NonNegativeInt,
46
UserInputAttachmentAnswerPayload,
57
type ChatAttachment,
68
ThreadId,
@@ -18,12 +20,17 @@ import * as Path from "effect/Path";
1820
import * as Schema from "effect/Schema";
1921
import * as Stream from "effect/Stream";
2022
import * as SqlClient from "effect/unstable/sql/SqlClient";
23+
import * as SqlSchema from "effect/unstable/sql/SqlSchema";
2124
import {
2225
legacyThreadPullRequestKey,
2326
threadPullRequestKeysEqual,
2427
} from "@t3tools/shared/threadPullRequests";
2528

26-
import { toPersistenceSqlError, type ProjectionRepositoryError } from "../../persistence/Errors.ts";
29+
import {
30+
toPersistenceDecodeError,
31+
toPersistenceSqlError,
32+
type ProjectionRepositoryError,
33+
} from "../../persistence/Errors.ts";
2734
import { OrchestrationEventStore } from "../../persistence/Services/OrchestrationEventStore.ts";
2835
import { ProjectionPendingApprovalRepository } from "../../persistence/Services/ProjectionPendingApprovals.ts";
2936
import { ProjectionProjectRepository } from "../../persistence/Services/ProjectionProjects.ts";
@@ -117,6 +124,13 @@ interface AttachmentSideEffects {
117124
readonly prunedThreadRelativePaths: Map<string, Set<string>>;
118125
}
119126

127+
const AttachmentCleanupReplayRow = Schema.Struct({
128+
sequence: NonNegativeInt,
129+
occurredAt: IsoDateTime,
130+
cleanupType: Schema.NullOr(Schema.Literals(["thread.reverted", "thread.deleted"])),
131+
threadId: Schema.NullOr(ThreadId),
132+
});
133+
120134
const materializeAttachmentsForProjection = Effect.fn("materializeAttachmentsForProjection")(
121135
(input: { readonly attachments: ReadonlyArray<ChatAttachment> }) =>
122136
Effect.succeed(input.attachments.length === 0 ? [] : input.attachments),
@@ -498,6 +512,49 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
498512
const path = yield* Path.Path;
499513
const serverConfig = yield* ServerConfig;
500514

515+
// Attachment cleanup has a separate retry cursor because its filesystem work
516+
// runs after projection transactions commit. Normal runtime events therefore
517+
// accumulate behind that cursor. Bootstrap only needs the latest event for
518+
// the new cursor plus revert/delete metadata; reading full event payloads here
519+
// made restart cost grow with every event since the previous restart.
520+
const listAttachmentCleanupReplayRows = SqlSchema.findAll({
521+
Request: Schema.Struct({ sequenceExclusive: NonNegativeInt }),
522+
Result: AttachmentCleanupReplayRow,
523+
execute: ({ sequenceExclusive }) => sql`
524+
SELECT
525+
sequence,
526+
occurred_at AS "occurredAt",
527+
CASE
528+
WHEN aggregate_kind = 'thread'
529+
AND event_type IN ('thread.reverted', 'thread.deleted')
530+
THEN event_type
531+
ELSE NULL
532+
END AS "cleanupType",
533+
CASE
534+
WHEN aggregate_kind = 'thread'
535+
AND event_type IN ('thread.reverted', 'thread.deleted')
536+
THEN stream_id
537+
ELSE NULL
538+
END AS "threadId"
539+
-- Force the sequence rowid range. With this OR predicate SQLite can
540+
-- otherwise choose the aggregate index and scan every thread event.
541+
FROM orchestration_events NOT INDEXED
542+
WHERE sequence > ${sequenceExclusive}
543+
AND (
544+
sequence = (
545+
SELECT MAX(sequence)
546+
FROM orchestration_events
547+
WHERE sequence > ${sequenceExclusive}
548+
)
549+
OR (
550+
aggregate_kind = 'thread'
551+
AND event_type IN ('thread.reverted', 'thread.deleted')
552+
)
553+
)
554+
ORDER BY sequence ASC
555+
`,
556+
});
557+
501558
const applyProjectsProjection: ProjectorDefinition["apply"] = Effect.fn(
502559
"applyProjectsProjection",
503560
)(function* (event, _attachmentSideEffects) {
@@ -1962,7 +2019,10 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
19622019
];
19632020

19642021
const applyAttachmentSideEffects = Effect.fn("applyAttachmentSideEffects")(
1965-
function* (event: OrchestrationEvent, sideEffects: AttachmentSideEffects) {
2022+
function* (
2023+
event: Pick<OrchestrationEvent, "sequence" | "type">,
2024+
sideEffects: AttachmentSideEffects,
2025+
) {
19662026
if (
19672027
sideEffects.deletedThreadIds.size === 0 &&
19682028
sideEffects.prunedThreadRelativePaths.size === 0
@@ -2124,27 +2184,37 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
21242184

21252185
// Cleanup has its own cursor so retries never have to replay committed text.
21262186
// All message and activity references are current before any files are removed.
2127-
const pendingCleanup = new Map<string, OrchestrationEvent>();
2128-
let lastEvent: OrchestrationEvent | undefined;
2129-
yield* Stream.runForEach(
2130-
eventStore.readFromSequence(cleanupStart, Number.MAX_SAFE_INTEGER),
2131-
(event) =>
2132-
Effect.sync(() => {
2133-
lastEvent = event;
2134-
if (event.type === "thread.reverted" || event.type === "thread.deleted") {
2135-
pendingCleanup.set(`${event.type}:${event.payload.threadId}`, event);
2136-
}
2137-
}),
2187+
const pendingCleanup = new Map<
2188+
string,
2189+
Schema.Schema.Type<typeof AttachmentCleanupReplayRow>
2190+
>();
2191+
const cleanupReplayRows = yield* listAttachmentCleanupReplayRows({
2192+
sequenceExclusive: cleanupStart,
2193+
}).pipe(
2194+
Effect.mapError((cause) =>
2195+
Schema.isSchemaError(cause)
2196+
? toPersistenceDecodeError("ProjectionPipeline.bootstrap:decodeCleanupRows")(cause)
2197+
: toPersistenceSqlError("ProjectionPipeline.bootstrap:listCleanupRows")(cause),
2198+
),
21382199
);
2200+
const lastEvent = cleanupReplayRows.at(-1);
2201+
for (const row of cleanupReplayRows) {
2202+
if (row.cleanupType !== null && row.threadId !== null) {
2203+
pendingCleanup.set(`${row.cleanupType}:${row.threadId}`, row);
2204+
}
2205+
}
21392206
for (const event of pendingCleanup.values()) {
2140-
if (event.type !== "thread.reverted" && event.type !== "thread.deleted") continue;
2141-
const threadId = event.payload.threadId;
2142-
const cleaned = yield* applyAttachmentSideEffects(event, {
2143-
deletedThreadIds: new Set(event.type === "thread.deleted" ? [threadId] : []),
2144-
prunedThreadRelativePaths: new Map(
2145-
event.type === "thread.reverted" ? [[threadId, new Set<string>()]] : [],
2146-
),
2147-
});
2207+
if (event.cleanupType === null || event.threadId === null) continue;
2208+
const threadId = event.threadId;
2209+
const cleaned = yield* applyAttachmentSideEffects(
2210+
{ sequence: event.sequence, type: event.cleanupType },
2211+
{
2212+
deletedThreadIds: new Set(event.cleanupType === "thread.deleted" ? [threadId] : []),
2213+
prunedThreadRelativePaths: new Map(
2214+
event.cleanupType === "thread.reverted" ? [[threadId, new Set<string>()]] : [],
2215+
),
2216+
},
2217+
);
21482218
// Leave the cleanup cursor behind this event so the next bootstrap retries it.
21492219
if (!cleaned) return;
21502220
}

0 commit comments

Comments
 (0)