Skip to content

Commit ace70ca

Browse files
perf(server): port #13673 to V2 — settling a thread closes its idle shells
Main's fix closed a settled thread's idle terminals from the V1 ProviderCommandReactor, which V2 deleted. V2's settle only detached the provider session, so idle shells kept holding the worktree. The merge already brought TerminalManager.closeIdle and the setup-script cleanup (ProjectSetupScriptRunner is shared). ThreadSettlementServiceV2 already consumes V2 domain events, so it now handles thread.settled: it re-reads the thread and, if it is still settled, calls closeIdle. A terminal running a command stays open. Covers manual and automatic settlement, which both emit thread.settled. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
1 parent 433ee8b commit ace70ca

2 files changed

Lines changed: 129 additions & 1 deletion

File tree

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

Lines changed: 102 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,10 +3,13 @@ import * as DateTime from "effect/DateTime";
33
import {
44
DEFAULT_SERVER_SETTINGS,
55
ProjectId,
6+
EventId,
67
ProviderInstanceId,
78
ThreadId,
89
type OrchestrationProjectShell,
10+
type OrchestrationV2AppThread,
911
type OrchestrationV2Command,
12+
type OrchestrationV2DomainEvent,
1013
type OrchestrationV2ShellSnapshot,
1114
type OrchestrationV2ThreadShell,
1215
type PullRequestSummary,
@@ -32,6 +35,7 @@ import {
3235
} from "../pullRequest/PullRequestService.ts";
3336
import { ServerActivation } from "../serverActivation.ts";
3437
import { ServerSettingsService } from "../serverSettings.ts";
38+
import { TerminalManager } from "../terminal/Manager.ts";
3539
import { ProjectionSnapshotQuery } from "../orchestration/Services/ProjectionSnapshotQuery.ts";
3640
import { OrchestratorV2, type OrchestratorV2Shape } from "./Orchestrator.ts";
3741
import { ProjectionStoreV2 } from "./ProjectionStore.ts";
@@ -410,6 +414,8 @@ interface HarnessOptions {
410414
readonly pullRequestSummary?: PullRequestService["Service"]["summary"];
411415
readonly existingWorktreePaths?: ReadonlyArray<string>;
412416
readonly onDispatch?: (command: AutoSettleCommand) => Effect.Effect<void>;
417+
/** Threads `getThread` returns when a `thread.settled` event is handled. */
418+
readonly currentThreads?: ReadonlyArray<OrchestrationV2AppThread>;
413419
}
414420

415421
const makeHarness = Effect.fn("makeThreadSettlementHarness")(function* (options: HarnessOptions) {
@@ -433,6 +439,8 @@ const makeHarness = Effect.fn("makeThreadSettlementHarness")(function* (options:
433439
>([]);
434440
const summaryRecovery = yield* Ref.make<ReadonlyArray<boolean | undefined>>([]);
435441
const invalidatedCwds = yield* Ref.make<ReadonlyArray<string>>([]);
442+
const domainEvents = yield* PubSub.unbounded<OrchestrationV2DomainEvent>();
443+
const closedIdle = yield* Queue.unbounded<{ readonly threadId: string }>();
436444

437445
const updateSettings = (patch: ServerSettingsPatch) =>
438446
Effect.gen(function* () {
@@ -499,11 +507,20 @@ const makeHarness = Effect.fn("makeThreadSettlementHarness")(function* (options:
499507
Effect.andThen(Ref.get(snapshots)),
500508
Effect.map((snapshot) => snapshot.threads),
501509
),
510+
getThread: (threadId) => {
511+
const thread = options.currentThreads?.find((candidate) => candidate.id === threadId);
512+
return thread
513+
? Effect.succeed(thread)
514+
: Effect.die(new Error(`Unexpected thread read: ${threadId}`));
515+
},
502516
}),
503517
Layer.mock(OrchestratorV2)({
504-
streamDomainEvents: Stream.empty,
518+
streamDomainEvents: Stream.fromPubSub(domainEvents),
505519
dispatch,
506520
}),
521+
Layer.mock(TerminalManager)({
522+
closeIdle: (input) => Queue.offer(closedIdle, input).pipe(Effect.asVoid),
523+
}),
507524
Layer.mock(GitManager)({
508525
branchPullRequest,
509526
invalidateStatus: (cwd) => Ref.update(invalidatedCwds, (cwds) => [...cwds, cwd]),
@@ -532,6 +549,8 @@ const makeHarness = Effect.fn("makeThreadSettlementHarness")(function* (options:
532549
summaryCalls,
533550
summaryRecovery,
534551
invalidatedCwds,
552+
closedIdle,
553+
publishEvent: (event: OrchestrationV2DomainEvent) => PubSub.publish(domainEvents, event),
535554
updateSettings,
536555
publishMerge: PubSub.publish(mergedPullRequests, {
537556
projectId: PROJECT_ID,
@@ -905,3 +924,85 @@ describe("ThreadSettlementServiceV2 worker", () => {
905924
),
906925
);
907926
});
927+
928+
describe("ThreadSettlementServiceV2 terminals", () => {
929+
const appThread = (
930+
id: string,
931+
settledOverride: OrchestrationV2AppThread["settledOverride"],
932+
): OrchestrationV2AppThread => {
933+
const threadId = ThreadId.make(id);
934+
return {
935+
createdBy: "user",
936+
creationSource: "web",
937+
id: threadId,
938+
projectId: PROJECT_ID,
939+
title: id,
940+
providerInstanceId: ProviderInstanceId.make("codex"),
941+
modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5" },
942+
runtimeMode: "full-access",
943+
interactionMode: "default",
944+
branch: null,
945+
worktreePath: null,
946+
activeProviderThreadId: null,
947+
lineage: { parentThreadId: null, relationshipToParent: null, rootThreadId: threadId },
948+
forkedFrom: null,
949+
createdAt: DateTime.makeUnsafe(NOW),
950+
updatedAt: DateTime.makeUnsafe(NOW),
951+
archivedAt: null,
952+
settledOverride,
953+
settledAt: settledOverride === "settled" ? DateTime.makeUnsafe(NOW) : null,
954+
lastVisitedAt: null,
955+
deletedAt: null,
956+
};
957+
};
958+
const settledEvent = (thread: OrchestrationV2AppThread): OrchestrationV2DomainEvent => ({
959+
type: "thread.settled",
960+
id: EventId.make(`event:settled:${thread.id}`),
961+
threadId: thread.id,
962+
occurredAt: DateTime.makeUnsafe(NOW),
963+
payload: thread,
964+
});
965+
966+
it.effect("closes the thread's idle terminals when it settles", () =>
967+
Effect.scoped(
968+
Effect.gen(function* () {
969+
const thread = appThread("settled-thread", "settled");
970+
const fixture = yield* makeHarness({
971+
snapshot: makeSnapshot([]),
972+
currentThreads: [thread],
973+
});
974+
975+
yield* Effect.gen(function* () {
976+
const service = yield* ThreadSettlementService.ThreadSettlementServiceV2;
977+
yield* startHarness(service, fixture.activation, fixture.snapshotReads);
978+
yield* fixture.publishEvent(settledEvent(thread));
979+
assert.deepStrictEqual(yield* Queue.take(fixture.closedIdle), { threadId: thread.id });
980+
}).pipe(Effect.provide(fixture.layer));
981+
}),
982+
),
983+
);
984+
985+
it.effect("keeps the terminals of a thread re-engaged before the event ran", () =>
986+
Effect.scoped(
987+
Effect.gen(function* () {
988+
const settled = appThread("reengaged-thread", "settled");
989+
const reengaged = appThread("reengaged-thread", "active");
990+
const marker = appThread("marker-thread", "settled");
991+
const fixture = yield* makeHarness({
992+
snapshot: makeSnapshot([]),
993+
currentThreads: [reengaged, marker],
994+
});
995+
996+
yield* Effect.gen(function* () {
997+
const service = yield* ThreadSettlementService.ThreadSettlementServiceV2;
998+
yield* startHarness(service, fixture.activation, fixture.snapshotReads);
999+
yield* fixture.publishEvent(settledEvent(settled));
1000+
// Events run in order, so the first close belongs to the later
1001+
// event only if the re-engaged thread was skipped.
1002+
yield* fixture.publishEvent(settledEvent(marker));
1003+
assert.deepStrictEqual(yield* Queue.take(fixture.closedIdle), { threadId: marker.id });
1004+
}).pipe(Effect.provide(fixture.layer));
1005+
}),
1006+
),
1007+
);
1008+
});

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

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import * as GitManager from "../git/GitManager.ts";
2222
import * as PullRequestService from "../pullRequest/PullRequestService.ts";
2323
import * as ServerSettings from "../serverSettings.ts";
2424
import { forkParked } from "../serverActivation.ts";
25+
import * as TerminalManager from "../terminal/Manager.ts";
2526
import * as ProjectionSnapshotQuery from "../orchestration/Services/ProjectionSnapshotQuery.ts";
2627
import { OrchestratorV2 } from "./Orchestrator.ts";
2728
import { ProjectionStoreV2, type ProjectionSettlementCandidate } from "./ProjectionStore.ts";
@@ -257,6 +258,7 @@ export const make = Effect.gen(function* () {
257258
const pullRequests = yield* PullRequestService.PullRequestService;
258259
const crypto = yield* Crypto.Crypto;
259260
const fileSystem = yield* FileSystem.FileSystem;
261+
const terminals = yield* TerminalManager.TerminalManager;
260262

261263
const sweep = Effect.fn("ThreadSettlementServiceV2.sweep")(function* (
262264
mergedPullRequest: PullRequestService.PullRequestMergeEvent | null,
@@ -506,8 +508,33 @@ export const make = Effect.gen(function* () {
506508
runSweep(null, threadId),
507509
);
508510

511+
// Settling closes the thread's shells that sit at an idle prompt, so they stop
512+
// holding the worktree. A terminal running a command (a dev server, an
513+
// editor) stays for the user to close.
514+
const closeIdleTerminals = Effect.fn("ThreadSettlementServiceV2.closeIdleTerminals")(
515+
function* (threadId: ThreadId) {
516+
// A thread re-engaged before this event ran keeps its shells.
517+
const thread = yield* projections.getThread(threadId);
518+
if (thread.settledOverride !== "settled") return;
519+
yield* terminals.closeIdle({ threadId });
520+
},
521+
(effect, threadId) =>
522+
effect.pipe(
523+
Effect.catchCause((cause) =>
524+
Cause.hasInterruptsOnly(cause)
525+
? Effect.failCause(cause)
526+
: Effect.logWarning("closing idle terminals after settlement failed", {
527+
threadId,
528+
cause: Cause.pretty(cause),
529+
}),
530+
),
531+
),
532+
);
533+
509534
const processEvent = (event: OrchestrationV2DomainEvent) => {
510535
switch (event.type) {
536+
case "thread.settled":
537+
return closeIdleTerminals(event.threadId);
511538
case "thread.pull-request-synced":
512539
case "provider-session.detached":
513540
return worker.enqueue(event.threadId);

0 commit comments

Comments
 (0)