Skip to content
Merged
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
1 change: 1 addition & 0 deletions apps/server/src/mcp/McpHttpServer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
),
Expand Down
98 changes: 97 additions & 1 deletion apps/server/src/mcp/toolkits/pair/handlers.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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) =>
Expand All @@ -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 };
});
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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 = <A, E extends { readonly _tag: string }>(effect: Effect.Effect<A, E>) =>
Expand Down Expand Up @@ -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 });
}),
);
});
81 changes: 81 additions & 0 deletions apps/server/src/mcp/toolkits/pair/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -82,6 +83,7 @@ import {
type PairAwaitResult,
type PairHandoffResult,
type PairStartResult,
type PairResetResult,
type PairStopResult,
} from "./tools.ts";

Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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,
});
});

Expand Down
29 changes: 29 additions & 0 deletions apps/server/src/mcp/toolkits/pair/tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand All @@ -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;
Expand Down Expand Up @@ -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<PairLeadNotFoundError>()(
Expand Down Expand Up @@ -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.",
Expand All @@ -419,4 +447,5 @@ export const PairToolkit = Toolkit.make(
PairHandoffTool,
PairAwaitTool,
PairStopTool,
PairResetTool,
);
6 changes: 6 additions & 0 deletions apps/server/src/provider/RuntimeInstructions.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"]) {
Expand Down
2 changes: 1 addition & 1 deletion apps/server/src/provider/RuntimeInstructions.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.`;

/**
Expand Down
11 changes: 11 additions & 0 deletions docs/internals/delegation.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 4 additions & 0 deletions docs/user/agent-delegation.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions packages/client-runtime/src/work-log/presentation.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
Loading
Loading