diff --git a/apps/server/src/mcp/McpHttpServer.test.ts b/apps/server/src/mcp/McpHttpServer.test.ts index 0bb18d91a5..085eccbe49 100644 --- a/apps/server/src/mcp/McpHttpServer.test.ts +++ b/apps/server/src/mcp/McpHttpServer.test.ts @@ -94,6 +94,7 @@ const PairTestLayer = McpHttpServer.PairToolkitRegistrationLive.pipe( Layer.mock(ProjectionSnapshotQuery)({}), Layer.mock(OrchestrationEngineService)({}), Layer.mock(ProviderRegistry)({}), + Layer.mock(ThreadDeletionReactor)({}), ServerSettings.layerTest(), NodeServices.layer, ), diff --git a/apps/server/src/mcp/toolkits/pair/handlers.test.ts b/apps/server/src/mcp/toolkits/pair/handlers.test.ts index 3f46523c98..540b4e57a1 100644 --- a/apps/server/src/mcp/toolkits/pair/handlers.test.ts +++ b/apps/server/src/mcp/toolkits/pair/handlers.test.ts @@ -44,6 +44,7 @@ import { import { ProjectionSnapshotQuery } from "../../../orchestration/Services/ProjectionSnapshotQuery.ts"; import { ProviderRegistry } from "../../../provider/Services/ProviderRegistry.ts"; import * as ServerSettings from "../../../serverSettings.ts"; +import { ThreadDeletionReactor } from "../../../orchestration/Services/ThreadDeletionReactor.ts"; import { PAIR_LEAD_PROTOCOL } from "../../../provider/RuntimeInstructions.ts"; import * as McpInvocationContext from "../../McpInvocationContext.ts"; import { PairToolkitHandlersLive } from "./handlers.ts"; @@ -283,6 +284,8 @@ const makeHarness = Effect.fn("makePairHarness")(function* (options: HarnessOpti } }; + // Dispatches and deletion-reactor drains, in the order they happened. + const order: string[] = []; // Mirrors the receipt store: an accepted command id replays without effect, // a rejected one fails as previously rejected forever. const dispatch: OrchestrationEngineShape["dispatch"] = (command) => @@ -302,6 +305,7 @@ const makeHarness = Effect.fn("makePairHarness")(function* (options: HarnessOpti } accepted.add(command.commandId); yield* Ref.update(commands, (recorded) => [...recorded, command]); + order.push(command.type); apply(command); return { sequence: 1 }; }); @@ -345,6 +349,9 @@ const makeHarness = Effect.fn("makePairHarness")(function* (options: HarnessOpti delegationChildRuntimeMode: options.childRuntimeMode ?? "inherit", }), Layer.succeed(Crypto.Crypto, testCrypto), + Layer.mock(ThreadDeletionReactor)({ + drainThrough: () => Effect.sync(() => void order.push("deletion-reactor drained")), + }), FileSystem.layerNoop({ readFile: (path) => { const content = files.get(path); @@ -382,7 +389,7 @@ const makeHarness = Effect.fn("makePairHarness")(function* (options: HarnessOpti const commandTypes = Ref.get(commands).pipe( Effect.map((recorded) => recorded.map((command) => command.type)), ); - return { commands, commandTypes, shells, files, call }; + return { commands, commandTypes, shells, files, call, order }; }); const tagOf = (effect: Effect.Effect) => @@ -1212,3 +1219,92 @@ describe("pair_stop", () => { }), ); }); + +describe("pair_reset", () => { + it.effect("replaces a finished executor with a fresh one on the same model", () => + Effect.gen(function* () { + const used = makeExecutor({ + latestTurn: completedTurn(), + modelSelection: { instanceId: ANTIGRAVITY, model: "gemini-3-flash" }, + runtimeMode: "approval-required", + // The lead has moved since the executor was created. + branch: "old-branch", + worktreePath: "/wt/old", + }); + const harness = yield* makeHarness({ + shells: [makeShell(LEAD_ID, { branch: "feat/now", worktreePath: "/wt/now" }), used], + }); + expect(yield* harness.call("pair_reset", {})).toEqual({ + threadId: EXECUTOR_ID, + reset: true, + state: "idle", + providerInstanceId: "antigravity", + model: "gemini-3-flash", + }); + const recorded = yield* Ref.get(harness.commands); + expect(recorded.map((command) => command.type)).toEqual(["thread.delete", "thread.create"]); + const [removed, created] = recorded; + expect(removed).toMatchObject({ threadId: EXECUTOR_ID }); + expect(created).toMatchObject({ + type: "thread.create", + threadId: EXECUTOR_ID, + modelSelection: { instanceId: ANTIGRAVITY, model: "gemini-3-flash" }, + runtimeMode: "approval-required", + interactionMode: "default", + branch: "feat/now", + worktreePath: "/wt/now", + }); + // The deletion reactor stops the old provider session after the delete. The + // new thread has the same id, so it must not exist until that has happened, + // or a brief sent straight after the reset could have its session stopped. + expect(harness.order).toEqual(["thread.delete", "deletion-reactor drained", "thread.create"]); + // The first create of a pair uses a fixed id. Reusing it here would replay + // that receipt as a success and create nothing. + expect(created?.commandId).not.toBe(`server:mcp-pair-create:${EXECUTOR_ID}`); + expect(created?.commandId.startsWith(`server:mcp-pair-reset-create:${EXECUTOR_ID}:`)).toBe( + true, + ); + expect(removed?.commandId.startsWith(`server:mcp-pair-reset-delete:${EXECUTOR_ID}:`)).toBe( + true, + ); + }), + ); + + it.effect("leaves an executor that never ran as it is", () => + Effect.gen(function* () { + const harness = yield* makeHarness({ shells: [makeShell(LEAD_ID), makeExecutor()] }); + expect(yield* harness.call("pair_reset", {})).toMatchObject({ + threadId: EXECUTOR_ID, + reset: false, + state: "idle", + }); + expect(yield* harness.commandTypes).toEqual([]); + }), + ); + + it.effect("refuses while the executor is running, and without a pair", () => + Effect.gen(function* () { + const busy = yield* makeHarness({ shells: [makeShell(LEAD_ID), runningExecutor()] }); + expect(yield* tagOf(busy.call("pair_reset", {}))).toBe("PairExecutorBusyError"); + expect(yield* busy.commandTypes).toEqual([]); + + const none = yield* makeHarness(); + expect(yield* tagOf(none.call("pair_reset", {}))).toBe("PairNotActiveError"); + + const off = yield* makeHarness({ + shells: [makeShell(LEAD_ID), makeExecutor({ archivedAt: NOW })], + }); + expect(yield* tagOf(off.call("pair_reset", {}))).toBe("PairNotActiveError"); + }), + ); + + it.effect("works for a paired lead with Pylon delegation off", () => + Effect.gen(function* () { + const harness = yield* makeHarness({ + shells: [makeShell(LEAD_ID), makeExecutor({ latestTurn: completedTurn() })], + }); + const paired = invocation({ capabilities: ["pair", "pull-requests"] }); + expect(yield* harness.call("pair_reset", {}, paired)).toMatchObject({ reset: true }); + }), + ); +}); diff --git a/apps/server/src/mcp/toolkits/pair/handlers.ts b/apps/server/src/mcp/toolkits/pair/handlers.ts index 34766cf849..78ed468392 100644 --- a/apps/server/src/mcp/toolkits/pair/handlers.ts +++ b/apps/server/src/mcp/toolkits/pair/handlers.ts @@ -31,6 +31,7 @@ import * as Semaphore from "effect/Semaphore"; import type { OrchestrationDispatchError } from "../../../orchestration/Errors.ts"; import * as OrchestrationEngine from "../../../orchestration/Services/OrchestrationEngine.ts"; import { markDelegationObservationConsumed } from "../../../orchestration/delegationObservationConsumed.ts"; +import * as ThreadDeletionReactor from "../../../orchestration/Services/ThreadDeletionReactor.ts"; import * as ProjectionSnapshotQuery from "../../../orchestration/Services/ProjectionSnapshotQuery.ts"; import * as ProviderRegistry from "../../../provider/Services/ProviderRegistry.ts"; import { PAIR_LEAD_PROTOCOL } from "../../../provider/RuntimeInstructions.ts"; @@ -82,6 +83,7 @@ import { type PairAwaitResult, type PairHandoffResult, type PairStartResult, + type PairResetResult, type PairStopResult, } from "./tools.ts"; @@ -128,6 +130,7 @@ const requirePairCapability = McpInvocationContext.requireMcpCapability("pair"). const make = Effect.gen(function* () { const engine = yield* OrchestrationEngine.OrchestrationEngineService; const snapshots = yield* ProjectionSnapshotQuery.ProjectionSnapshotQuery; + const deletionReactor = yield* ThreadDeletionReactor.ThreadDeletionReactor; const providers = yield* ProviderRegistry.ProviderRegistry; const serverSettings = yield* ServerSettings.ServerSettingsService; const crypto = yield* Crypto.Crypto; @@ -656,11 +659,89 @@ const make = Effect.gen(function* () { ); }); + const pair_reset = () => + Effect.gen(function* () { + const scope = yield* requirePairCapability; + const executorId = yield* executorIdFor(scope.threadId); + + return yield* withLeadGate(scope.threadId)( + Effect.gen(function* () { + const existing = yield* findExecutor(executorId); + if (Option.isNone(existing) || existing.value.archivedAt !== null) { + return yield* new PairNotActiveError(); + } + const executor = existing.value; + const state = derivePairExecutorState(executor); + if (state === "running") { + return yield* new PairExecutorBusyError(); + } + + if (executor.latestTurn === null && executor.session === null) { + const result: PairResetResult = { + threadId: executorId, + reset: false, + state, + providerInstanceId: executor.modelSelection.instanceId, + model: executor.modelSelection.model, + }; + return result; + } + + const leadOpt = yield* orFail(snapshots.getThreadShellById(scope.threadId)); + if (Option.isNone(leadOpt)) { + return yield* new PairLeadNotFoundError({ threadId: scope.threadId }); + } + const lead = leadOpt.value; + + const uuid = yield* orFail(crypto.randomUUIDv4); + const deleted = yield* engine + .dispatch({ + type: "thread.delete", + commandId: CommandId.make(`server:mcp-pair-reset-delete:${executorId}:${uuid}`), + threadId: executorId, + }) + .pipe(mapDispatch(() => undefined)); + // The deletion reactor stops the old provider session. The new thread + // reuses the id, so wait, or that stop could land on the new session. + yield* orFail(deletionReactor.drainThrough(deleted.sequence)); + + yield* engine + .dispatch({ + type: "thread.create", + commandId: CommandId.make(`server:mcp-pair-reset-create:${executorId}:${uuid}`), + threadId: executorId, + projectId: executor.projectId, + title: executor.title, + modelSelection: executor.modelSelection, + runtimeMode: executor.runtimeMode, + interactionMode: "default", + branch: lead.branch, + worktreePath: lead.worktreePath, + createdAt: yield* nowIso, + }) + .pipe(mapDispatch(() => undefined)); + + protectedRecords.delete(executorId); + instantReadCounts.delete(executorId); + + const result: PairResetResult = { + threadId: executorId, + reset: true, + state: "idle", + providerInstanceId: executor.modelSelection.instanceId, + model: executor.modelSelection.model, + }; + return result; + }), + ); + }); + return PairToolkit.of({ pair_start, pair_handoff, pair_await, pair_stop, + pair_reset, }); }); diff --git a/apps/server/src/mcp/toolkits/pair/tools.ts b/apps/server/src/mcp/toolkits/pair/tools.ts index 4a0f61b56f..9d6c65a18b 100644 --- a/apps/server/src/mcp/toolkits/pair/tools.ts +++ b/apps/server/src/mcp/toolkits/pair/tools.ts @@ -19,6 +19,7 @@ import * as Tool from "effect/unstable/ai/Tool"; import * as Toolkit from "effect/unstable/ai/Toolkit"; import * as OrchestrationEngine from "../../../orchestration/Services/OrchestrationEngine.ts"; +import * as ThreadDeletionReactor from "../../../orchestration/Services/ThreadDeletionReactor.ts"; import * as ProjectionSnapshotQuery from "../../../orchestration/Services/ProjectionSnapshotQuery.ts"; import * as ProviderRegistry from "../../../provider/Services/ProviderRegistry.ts"; import * as ServerSettings from "../../../serverSettings.ts"; @@ -34,6 +35,8 @@ const dependencies = [ ServerSettings.ServerSettingsService, // Reads the lead's protected paths in the worktree the pair shares. FileSystem.FileSystem, + // A reset waits for the old executor's session to be stopped before reusing its id. + ThreadDeletionReactor.ThreadDeletionReactor, ]; const MAX_BRIEF_CHARS = 32_000; @@ -189,6 +192,18 @@ export const PairStopResult = Schema.Struct({ }); export type PairStopResult = typeof PairStopResult.Type; +// ---- pair_reset ---- + +export const PairResetResult = Schema.Struct({ + threadId: Schema.String, + /** False when there was nothing to clear: the executor had never run. */ + reset: Schema.Boolean, + state: PairExecutorState, + providerInstanceId: Schema.String, + model: Schema.String, +}); +export type PairResetResult = typeof PairResetResult.Type; + // ---- Errors ---- export class PairLeadNotFoundError extends Schema.TaggedError()( @@ -401,6 +416,19 @@ const PairAwaitTool = Tool.make("pair_await", { .annotate(Tool.Idempotent, true) .annotate(Tool.OpenWorld, false); +const PairResetTool = Tool.make("pair_reset", { + description: + "Start the executor over with an empty context, on the same model, in your current worktree. Use it when the executor's context is nearly full or it has lost the thread of the work. Its transcript is deleted; the files it changed are untouched. Refused while it is running: stop it first. The next brief must stand on its own, because the executor remembers nothing.", + success: PairResetResult, + failure: PairToolError, + dependencies, +}) + .annotate(Tool.Title, "Reset the pair executor") + .annotate(Tool.Readonly, false) + .annotate(Tool.Destructive, true) + .annotate(Tool.Idempotent, false) + .annotate(Tool.OpenWorld, false); + const PairStopTool = Tool.make("pair_stop", { description: "Request an interrupt of the executor's running turn. interrupted=true means the request was accepted; confirm with pair_await.", @@ -419,4 +447,5 @@ export const PairToolkit = Toolkit.make( PairHandoffTool, PairAwaitTool, PairStopTool, + PairResetTool, ); diff --git a/apps/server/src/provider/RuntimeInstructions.test.ts b/apps/server/src/provider/RuntimeInstructions.test.ts index 71fb4247c0..b625c4647f 100644 --- a/apps/server/src/provider/RuntimeInstructions.test.ts +++ b/apps/server/src/provider/RuntimeInstructions.test.ts @@ -95,6 +95,12 @@ describe("buildRuntimeInstructions", () => { } }); + it("tells a paired lead how to give the executor a fresh start", () => { + for (const rule of ["pair_reset", "remembers nothing"]) { + expect(PAIR_LEAD_PROTOCOL).toContain(rule); + } + }); + it("tells a paired lead that delegating means briefing its executor", () => { // A user who says "use delegation" on a paired thread means the executor. for (const rule of ["delegate_thread", "delegate", "means brief your executor"]) { diff --git a/apps/server/src/provider/RuntimeInstructions.ts b/apps/server/src/provider/RuntimeInstructions.ts index 019224f477..5cbc0079a6 100644 --- a/apps/server/src/provider/RuntimeInstructions.ts +++ b/apps/server/src/provider/RuntimeInstructions.ts @@ -35,7 +35,7 @@ Keep small or tightly coupled work local. Before choosing a delegation method fo export const PAIR_LEAD_PROTOCOL = `You are the lead of a Pylon pair. One executor thread on a faster, cheaper model is linked to this thread and works in your worktree. Hand implementation to it instead of using your own subagents. The executor is reached only through the pair tools on the t3-code MCP server: pair_handoff, pair_await, pair_stop. If they are not in your tool list, search your tools for pair_handoff. An agent started with your harness's own tools (Agent, Task, spawn_agent, or anything in a collaboration namespace) is not your executor: it runs on your model, at your cost, and Pylon cannot see it. Do not use those tools while paired, even when asked to "brief the executor". The same goes for Pylon's fan-out tools: delegate_thread is refused on a paired thread. If the user asks you to delegate or to "use delegation" here, that means brief your executor. Work test-first. Write the contract (types, signatures, stubs) and the failing tests yourself, run them, and confirm they fail for the right reason. Commit them. Then brief the executor with pair_handoff, listing your tests and contract as protectedPaths: make these tests pass without editing them, with the exact files, the behavior wanted, the commands to run, the exclusions, and the report format. Hand off a whole plan step, not small nudges. Where tests cannot express the work, give exact acceptance checks instead. Ask the executor for code only, and write documentation and pull request text yourself. -Wait with pair_await, or end your turn: Pylon wakes you when the executor finishes or needs the user. Never poll in a loop. +Wait with pair_await, or end your turn: Pylon wakes you when the executor finishes or needs the user. Never poll in a loop. If the executor's context is nearly full or it has lost the thread of the work, call pair_reset between briefs: it starts over on the same model and remembers nothing, so the next brief must stand on its own. When it reports, confirm your tests are unchanged (protectedPaths.changed in the pair_await result is empty), read the diff in your worktree, and re-run the checks yourself. The executor's report is not verification, and its prose is less reliable than its code. Send one consolidated correction, at most three rounds, then ask the user. Only you push and open pull requests. Make small or tightly coupled changes yourself. If the executor needs an approval or an answer, tell the user; never approve on their behalf.`; /** diff --git a/docs/internals/delegation.md b/docs/internals/delegation.md index f6d421bb2d..f7a8f8b94d 100644 --- a/docs/internals/delegation.md +++ b/docs/internals/delegation.md @@ -171,6 +171,17 @@ executor that never ran, whose lead is not among the active or archived threads, than a day old. A missing lead alone proves nothing, because that is what every freshly paired draft looks like. The delete uses a deterministic command id, and a failed sweep only logs. +`pair_reset` gives the executor an empty context without ending the pair. One executor serves +every brief, so its context only grows; the first one used in earnest ended a day at 126k of 128k +tokens. No command starts a thread's provider session afresh, and adding one would have meant a +contract and decider change, so the reset deletes the executor thread and creates it again under the +same id, with the same model and runtime mode and the lead's current branch and worktree. The +transcript goes; the files stay. It is refused while the executor runs, does nothing for one that +never ran, and uses unique command ids: the first create of a pair has a fixed id, and reusing it +would replay that receipt and create nothing. Between the delete and the create it waits for the +thread deletion reactor, which stops the old provider session after `thread.deleted`; the new thread +has the same id, so a brief sent straight after the reset could otherwise have its session stopped. + A brief can name the files its lead owns, normally its tests and contract, as `protectedPaths`. The handlers hash each file when the brief is accepted and `pair_await` reports the ones whose content changed or that disappeared, once the executor is no longer running. This is what makes diff --git a/docs/user/agent-delegation.md b/docs/user/agent-delegation.md index caf7af65ae..3bd1780491 100644 --- a/docs/user/agent-delegation.md +++ b/docs/user/agent-delegation.md @@ -152,6 +152,10 @@ executor but cannot lead a pair, because Pylon has no way to pause its own subag Antigravity thread **Pair** is unavailable and says why. Codex can lead, though its own subagents cannot be switched off, so it is told to leave them alone rather than prevented from using them. +A long pair can fill the executor's context. The lead can start it over with an empty one, on the +same model, without ending the pair; the executor's transcript is removed and the files it changed +stay as they are. You can ask for it: “reset your executor.” + The executor appears under its parent in the sidebar like any delegated thread. Archive it to turn the pair off for that thread. Archiving, settling, or deleting the lead does the same to its executor. Rewinding the lead stops the executor first, because both work on the same files. The Pair diff --git a/packages/client-runtime/src/work-log/presentation.test.ts b/packages/client-runtime/src/work-log/presentation.test.ts index 5e2d76d7ef..debe0343f6 100644 --- a/packages/client-runtime/src/work-log/presentation.test.ts +++ b/packages/client-runtime/src/work-log/presentation.test.ts @@ -665,6 +665,7 @@ describe("pair tool presentation", () => { ["pair_handoff", "Handed off to the pair executor", "Handing off to the pair executor"], ["pair_await", "Awaited the pair executor", "Awaiting the pair executor"], ["pair_stop", "Stopped the pair executor", "Stopping the pair executor"], + ["pair_reset", "Reset the pair executor", "Resetting the pair executor"], ])("labels %s with the Pylon tool icon", (tool, completed, running) => { for (const label of [`mcp__t3-code__${tool}`, `t3-code · ${tool}`, tool]) { expect(resolveWorkEntryToolPresentation({ label, toolLifecycleStatus: "completed" })).toEqual( diff --git a/packages/client-runtime/src/work-log/presentation.ts b/packages/client-runtime/src/work-log/presentation.ts index 585a5092fe..00110daead 100644 --- a/packages/client-runtime/src/work-log/presentation.ts +++ b/packages/client-runtime/src/work-log/presentation.ts @@ -87,6 +87,7 @@ const T3_MCP_TOOL_LABELS: Record< pair_handoff: ["Hand off", "Handing off", "Handed off", "to the pair executor"], pair_await: ["Await", "Awaiting", "Awaited", "the pair executor"], pair_stop: ["Stop", "Stopping", "Stopped", "the pair executor"], + pair_reset: ["Reset", "Resetting", "Reset", "the pair executor"], orchestrator_capabilities: ["Get", "Getting", "Got", "orchestration capabilities"], delegate_task: ["Delegate", "Delegating", "Delegated", "a child task"], task_status: ["Get", "Getting", "Got", "delegated task status"],