Skip to content

Commit 11acae1

Browse files
perf(server): port #13765 to V2 — background PR discovery reads only threads that can still settle
Main's periodic discovery and settlement sweeps stopped reading settled threads. V2's settlement sweep already reads candidates only (getSettlementCandidates filters settled, pinned, archived, and busy threads in SQL), but V2's PR discovery ran a full getShellSnapshot every minute: every thread, active and archived, settled or not, with every thread's run, item, and session subqueries. Discovery then dropped archived and settled threads unless they were in backfill. getShellSnapshot takes an `unsettledOnly` option (SQL and memory store), and the periodic discovery sweep reads active, unsettled threads unless a backfill pass is pending, as on main. A pending backfill entry for a thread with no branch is cleared so it cannot keep every pass on the full read. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
1 parent 1b10d08 commit 11acae1

5 files changed

Lines changed: 138 additions & 14 deletions

File tree

‎apps/server/src/orchestration-v2/Orchestrator.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -268,6 +268,8 @@ export interface OrchestratorV2Shape {
268268
>;
269269
readonly getShellSnapshot: (options?: {
270270
readonly location?: "active" | "archive";
271+
/** Background sweeps only: skips settled threads. */
272+
readonly unsettledOnly?: boolean;
271273
}) => Effect.Effect<OrchestrationV2ThreadShellSnapshot, OrchestratorV2Error>;
272274
readonly getThreadShell: (
273275
threadId: ThreadId,

‎apps/server/src/orchestration-v2/ProjectionSettlement.test.ts‎

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -344,6 +344,31 @@ for (const [name, testLayer] of [
344344
);
345345
}
346346

347+
for (const [name, testLayer] of [
348+
["sql", SqlLayer],
349+
["memory", layerMemory],
350+
] as const) {
351+
it.effect(`${name}: an unsettled-only shell read skips settled threads`, () =>
352+
Effect.gen(function* () {
353+
const store = yield* ProjectionStoreV2;
354+
const open = yield* createThread("unsettled-open");
355+
const reopened = yield* createThread("unsettled-reopened", { settledOverride: "active" });
356+
yield* createThread("unsettled-manual", { settledOverride: "settled", settledAt: old });
357+
yield* createThread("unsettled-auto", { settledAt: old });
358+
yield* createThread("unsettled-archived", { archivedAt: old });
359+
360+
const shell = yield* store.getShellSnapshot({ location: "active", unsettledOnly: true });
361+
assert.deepEqual(
362+
new Set(shell.threads.map((thread) => thread.id)),
363+
new Set([open, reopened]),
364+
);
365+
assert.equal(shell.archivedThreads.length, 0);
366+
const all = yield* store.getShellSnapshot({ location: "active" });
367+
assert.equal(all.threads.length, 4);
368+
}).pipe(Effect.provide(testLayer)),
369+
);
370+
}
371+
347372
it.effect(
348373
"reads settlement candidates and thread metadata without loading historical or archived payloads",
349374
() =>

‎apps/server/src/orchestration-v2/ProjectionStore.ts‎

Lines changed: 26 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -319,8 +319,13 @@ export interface ProjectionStoreV2Shape {
319319
readonly apply: (
320320
event: OrchestrationV2DomainEvent,
321321
) => Effect.Effect<void, ProjectionStoreV2Error>;
322+
/**
323+
* `unsettledOnly` is for background sweeps, not clients: it skips settled
324+
* threads before any of their run, item or session rows are read.
325+
*/
322326
readonly getShellSnapshot: (options?: {
323327
readonly location?: "active" | "archive";
328+
readonly unsettledOnly?: boolean;
324329
}) => Effect.Effect<OrchestrationV2ThreadShellSnapshot, ProjectionStoreV2Error>;
325330
readonly getThreadShell: (
326331
threadId: ThreadId,
@@ -4677,7 +4682,11 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
46774682
),
46784683
);
46794684

4680-
const selectShellThreadRows = (threadId?: ThreadId, location?: "active" | "archive") =>
4685+
const selectShellThreadRows = (
4686+
threadId?: ThreadId,
4687+
location?: "active" | "archive",
4688+
unsettledOnly = false,
4689+
) =>
46814690
sql<ShellThreadRow>`
46824691
SELECT
46834692
t.thread_id,
@@ -4815,6 +4824,10 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
48154824
: location === "archive"
48164825
? sql` AND json_extract(t.payload_json, '$.archivedAt') IS NOT NULL`
48174826
: sql``
4827+
}${
4828+
unsettledOnly
4829+
? sql` AND json_extract(t.payload_json, '$.settledAt') IS NULL AND json_extract(t.payload_json, '$.settledOverride') IS NOT 'settled'`
4830+
: sql``
48184831
}
48194832
ORDER BY t.updated_at ASC, t.thread_id ASC
48204833
`;
@@ -5169,7 +5182,11 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
51695182
sql
51705183
.withTransaction(
51715184
Effect.gen(function* () {
5172-
const targetThreadRows = yield* selectShellThreadRows(undefined, options?.location);
5185+
const targetThreadRows = yield* selectShellThreadRows(
5186+
undefined,
5187+
options?.location,
5188+
options?.unsettledOnly ?? false,
5189+
);
51735190
const targetThreadIds = new Set(
51745191
targetThreadRows.map((row) => ThreadId.make(row.thread_id)),
51755192
);
@@ -5393,6 +5410,13 @@ export const layerMemory: Layer.Layer<ProjectionStoreV2> = Layer.effect(
53935410
const existing = (yield* Ref.get(replayState)).projections;
53945411
const selectedThreadIds = [...existing.entries()]
53955412
.filter(([, projection]) => {
5413+
if (
5414+
options?.unsettledOnly &&
5415+
(projection.thread.settledAt !== null ||
5416+
projection.thread.settledOverride === "settled")
5417+
) {
5418+
return false;
5419+
}
53965420
if (options?.location === "active") return projection.thread.archivedAt === null;
53975421
if (options?.location === "archive") return projection.thread.archivedAt !== null;
53985422
return true;

‎apps/server/src/orchestration-v2/ThreadPullRequestService.test.ts‎

Lines changed: 70 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ import * as Option from "effect/Option";
1818
import * as PubSub from "effect/PubSub";
1919
import * as Queue from "effect/Queue";
2020
import * as Stream from "effect/Stream";
21+
import * as TestClock from "effect/testing/TestClock";
2122

2223
import { GitManager } from "../git/GitManager.ts";
2324
import { ProjectionSnapshotQuery } from "../orchestration/Services/ProjectionSnapshotQuery.ts";
@@ -164,13 +165,15 @@ describe("ThreadPullRequestServiceV2 reads", () => {
164165
const other = threadShell("other-thread");
165166
const activation = yield* Deferred.make<void>();
166167
const events = yield* PubSub.unbounded<OrchestrationV2DomainEvent>();
167-
// Each read: the thread id for a one-thread read, null for a full read.
168-
const reads = yield* Queue.unbounded<ThreadId | null>();
168+
// Each read: the thread id for a one-thread read, or the full read's options.
169+
const reads = yield* Queue.unbounded<
170+
ThreadId | { readonly location?: string; readonly unsettledOnly?: boolean }
171+
>();
169172
const dependencies = Layer.mergeAll(
170173
Layer.mock(OrchestratorV2)({
171174
streamDomainEvents: Stream.fromPubSub(events),
172-
getShellSnapshot: () =>
173-
Queue.offer(reads, null).pipe(
175+
getShellSnapshot: (options) =>
176+
Queue.offer(reads, options ?? {}).pipe(
174177
Effect.as({
175178
schemaVersion: 2,
176179
snapshotSequence: 1,
@@ -205,8 +208,8 @@ describe("ThreadPullRequestServiceV2 reads", () => {
205208
const service = yield* make;
206209
yield* service.start();
207210
yield* Deferred.succeed(activation, undefined);
208-
// Startup backfill is a full read.
209-
expect(yield* Queue.take(reads)).toBeNull();
211+
// Startup backfill reads every active thread.
212+
expect(yield* Queue.take(reads)).toEqual({ location: "active", unsettledOnly: false });
210213
yield* service.drain;
211214
yield* PubSub.publish(events, {
212215
type: "thread.metadata-updated",
@@ -244,4 +247,65 @@ describe("ThreadPullRequestServiceV2 reads", () => {
244247
}),
245248
),
246249
);
250+
251+
it.effect("periodic sweeps read only unsettled threads once backfill is done", () =>
252+
Effect.scoped(
253+
Effect.gen(function* () {
254+
const settled = {
255+
...threadShell("settled-thread"),
256+
branch: "feature/settled",
257+
settledOverride: "settled" as const,
258+
settledAt: NOW,
259+
};
260+
const activation = yield* Deferred.make<void>();
261+
const reads = yield* Queue.unbounded<{
262+
readonly location?: string;
263+
readonly unsettledOnly?: boolean;
264+
}>();
265+
const dependencies = Layer.mergeAll(
266+
Layer.mock(OrchestratorV2)({
267+
streamDomainEvents: Stream.never,
268+
getShellSnapshot: (options) =>
269+
Queue.offer(reads, options ?? {}).pipe(
270+
Effect.as({
271+
schemaVersion: 2,
272+
snapshotSequence: 1,
273+
// The fake honors unsettledOnly like the store does.
274+
threads: options?.unsettledOnly ? [] : [settled],
275+
archivedThreads: [],
276+
}),
277+
),
278+
}),
279+
Layer.mock(ProjectionSnapshotQuery)({
280+
getProjectShellsWithoutEnrichment: () => Effect.succeed([]),
281+
}),
282+
Layer.mock(GitManager)({}),
283+
Layer.mock(PullRequestService)({}),
284+
Layer.mock(RepositoryIdentityResolver)({}),
285+
Layer.succeed(ServerActivation, Deferred.await(activation)),
286+
Layer.succeed(
287+
Crypto.Crypto,
288+
Crypto.make({
289+
randomBytes: (size) => new Uint8Array(size).fill(1),
290+
digest: (_algorithm, data) => Effect.succeed(data),
291+
}),
292+
),
293+
FileSystem.layerNoop({}),
294+
);
295+
296+
yield* Effect.gen(function* () {
297+
const service = yield* make;
298+
yield* service.start();
299+
yield* Deferred.succeed(activation, undefined);
300+
// Backfill finds the settled branch thread and must read it again.
301+
expect(yield* Queue.take(reads)).toEqual({ location: "active", unsettledOnly: false });
302+
yield* service.drain;
303+
// Its project is gone, so backfill finishes it on the first pass.
304+
yield* TestClock.adjust("1 minute");
305+
expect(yield* Queue.take(reads)).toEqual({ location: "active", unsettledOnly: true });
306+
yield* service.drain;
307+
}).pipe(Effect.provide(dependencies));
308+
}),
309+
),
310+
);
247311
});

‎apps/server/src/orchestration-v2/ThreadPullRequestService.ts‎

Lines changed: 15 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -121,11 +121,16 @@ export const make = Effect.gen(function* () {
121121

122122
/**
123123
* A sweep for one thread reads only that thread's shell, not every thread's.
124-
* Finished runs and checkpoints queue one of these each.
124+
* Finished runs and checkpoints queue one of these each. A sweep over all
125+
* threads reads only active, unsettled ones, since discovery skips the rest;
126+
* backfill looks up settled threads, so its passes read every active thread.
125127
*/
126-
const readThreadSnapshot = (threadId: ThreadId | null) =>
128+
const readThreadSnapshot = ({ threadId, backfill }: RefreshRequest) =>
127129
threadId === null
128-
? orchestrator.getShellSnapshot()
130+
? orchestrator.getShellSnapshot({
131+
location: "active",
132+
unsettledOnly: !(backfill || pendingBackfill.size > 0),
133+
})
129134
: Effect.gen(function* () {
130135
// Read the sequence first. The thread is then at least this new, so a
131136
// sync guarded by the sequence is rejected rather than missing a change.
@@ -138,7 +143,7 @@ export const make = Effect.gen(function* () {
138143
request: RefreshRequest,
139144
) {
140145
const [threadSnapshot, projectShells] = yield* Effect.all([
141-
readThreadSnapshot(request.threadId),
146+
readThreadSnapshot(request),
142147
snapshots.getProjectShellsWithoutEnrichment(),
143148
]);
144149
const projects = new Map(projectShells.map((project) => [project.id, project]));
@@ -152,8 +157,12 @@ export const make = Effect.gen(function* () {
152157
}
153158
}
154159
}
155-
// A single-thread read only shows whether its own thread is gone.
156-
const visibleThreadIds = new Set(threadSnapshot.threads.map((thread) => thread.id));
160+
// A single-thread read only shows whether its own thread is gone. A thread
161+
// with no branch has nothing to look up, and its entry would keep every
162+
// periodic pass on the full read.
163+
const visibleThreadIds = new Set(
164+
threadSnapshot.threads.filter((thread) => thread.branch !== null).map((thread) => thread.id),
165+
);
157166
const checkedIds = request.threadId === null ? [...pendingBackfill.keys()] : [request.threadId];
158167
for (const threadId of checkedIds) {
159168
if (!visibleThreadIds.has(threadId)) pendingBackfill.delete(threadId);

0 commit comments

Comments
 (0)