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
7 changes: 4 additions & 3 deletions apps/server/src/mcp/toolkits/delegation/handlers.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -705,11 +705,12 @@ describe("delegated_thread_status", () => {
const fiber = yield* Effect.forkChild(
harness.call("delegated_thread_status", { delegationKey: "k1", waitSeconds: 5 }),
);
yield* TestClock.adjust("5 seconds");
// A short request is not honored: it waits the whole cap.
yield* TestClock.adjust("45 seconds");
expect(yield* Fiber.join(fiber)).toMatchObject({
state: "queued",
changed: false,
waitedSeconds: 5,
waitedSeconds: 45,
});
}),
);
Expand Down Expand Up @@ -746,7 +747,7 @@ describe("delegated_thread_status", () => {
const runningFiber = yield* Effect.forkChild(
runningHarness.call("delegated_thread_status", { delegationKey: "k1", waitSeconds: 5 }),
);
yield* TestClock.adjust("5 seconds");
yield* TestClock.adjust("45 seconds");
yield* Fiber.join(runningFiber);
const runningCommands = yield* Ref.get(runningHarness.commands);
expect(runningCommands).toHaveLength(0);
Expand Down
3 changes: 2 additions & 1 deletion apps/server/src/mcp/toolkits/delegation/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@ import * as McpInvocationContext from "../../McpInvocationContext.ts";
import {
aggregateFilesChanged,
defaultTitleFor,
delegatedStatusWaitSeconds,
delegatedThreadId,
deriveDelegatedThreadState,
isChildOfParent,
Expand Down Expand Up @@ -347,7 +348,7 @@ const make = Effect.gen(function* () {
const startedAt = yield* Clock.currentTimeMillis;
// Wall-clock bound: each poll's query time counts against the budget so
// the call always returns inside the caller's tool timeout.
const deadline = startedAt + (input.waitSeconds ?? 0) * 1_000;
const deadline = startedAt + delegatedStatusWaitSeconds(input.waitSeconds) * 1_000;
const elapsedSeconds = (now: number) => Math.round((now - startedAt) / 1_000);
let current = shell;
let now = startedAt;
Expand Down
12 changes: 12 additions & 0 deletions apps/server/src/mcp/toolkits/delegation/logic.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@ import { describe, expect, it } from "vite-plus/test";

import {
aggregateFilesChanged,
delegatedStatusWaitSeconds,
MAX_STATUS_WAIT_SECONDS,
defaultTitleFor,
delegatedChildPrefix,
delegatedThreadId,
Expand Down Expand Up @@ -439,3 +441,13 @@ describe("resolveDelegationTarget", () => {
);
});
});

describe("delegated status wait", () => {
it("reads at once for zero or nothing, and waits the whole cap for anything else", () => {
expect(MAX_STATUS_WAIT_SECONDS).toBe(45);
expect(delegatedStatusWaitSeconds(undefined)).toBe(0);
expect(delegatedStatusWaitSeconds(0)).toBe(0);
expect(delegatedStatusWaitSeconds(5)).toBe(45);
expect(delegatedStatusWaitSeconds(45)).toBe(45);
});
});
13 changes: 13 additions & 0 deletions apps/server/src/mcp/toolkits/delegation/logic.ts
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,19 @@ export interface DelegatedFileChange {
readonly deletions: number;
}

/** The longest one `delegated_thread_status` call may wait. Prime Agent cancels MCP calls at 60 s. */
export const MAX_STATUS_WAIT_SECONDS = 45;

/**
* How long one `delegated_thread_status` call waits. Omitted or 0 reads the
* state and returns. Any positive request waits the whole cap: a short wait
* repeated in a loop is polling with a model turn per call, and the call
* already returns the moment the child changes.
*/
export function delegatedStatusWaitSeconds(requested: number | undefined): number {
return requested !== undefined && requested > 0 ? MAX_STATUS_WAIT_SECONDS : 0;
}

export function aggregateFilesChanged(
checkpoints: ReadonlyArray<Pick<OrchestrationCheckpointSummary, "files">>,
): ReadonlyArray<DelegatedFileChange> {
Expand Down
10 changes: 4 additions & 6 deletions apps/server/src/mcp/toolkits/delegation/tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import * as ProviderRegistry from "../../../provider/Services/ProviderRegistry.t
import * as ServerSettings from "../../../serverSettings.ts";
import * as VcsStatusBroadcaster from "../../../vcs/VcsStatusBroadcaster.ts";
import * as McpInvocationContext from "../../McpInvocationContext.ts";
import { MAX_STATUS_WAIT_SECONDS } from "./logic.ts";

// Crypto is not listed: it reaches the handlers as a layer requirement of
// `toLayer(make)`, the same way the pull request toolkit receives it.
Expand All @@ -37,9 +38,6 @@ const dependencies = [
];

const MAX_TASK_CHARS = 32_000;
// Prime Agent's MCP client cancels any call after 60 s, and each poll adds
// query time, so the budget stays well inside that.
const MAX_WAIT_SECONDS = 45;
const MAX_TITLE_CHARS = 200;
const MIN_RESULT_CHARS = 1_000;
const MAX_RESULT_CHARS = 60_000;
Expand Down Expand Up @@ -120,9 +118,9 @@ export type DelegateThreadResult = typeof DelegateThreadResult.Type;
export const DelegatedThreadStatusInput = Schema.Struct({
delegationKey: DelegationKey,
waitSeconds: Schema.optional(
Schema.Int.check(Schema.isBetween({ minimum: 0, maximum: MAX_WAIT_SECONDS })).annotate({
Schema.Int.check(Schema.isBetween({ minimum: 0, maximum: MAX_STATUS_WAIT_SECONDS })).annotate({
description:
"Wait up to this many seconds (at most 45) for the child's state to change before answering. Use 45 when waiting is necessary; do independent work before checking again. Completed or blocked children return immediately.",
"Leave it out or pass 0 to read the state now; any positive value waits the full 45 seconds and returns the moment the child changes, so a shorter wait gains nothing; do not call it in a loop. Completed or blocked children return immediately.",
}),
),
});
Expand Down Expand Up @@ -448,7 +446,7 @@ const DelegateThreadTool = Tool.make("delegate_thread", {

const DelegatedThreadStatusTool = Tool.make("delegated_thread_status", {
description:
"Report a child thread's state (queued, running, completed, interrupted, error, archived) and whether it is waiting on an approval or a question. Pass waitSeconds to wait up to 45 seconds until something changes; use 45 when waiting is necessary. Completed children and children needing approval or input return immediately. Do independent work before checking again; avoid short polling and unchanged progress updates. While Pylon delegation is enabled, child lifecycle changes can queue automatic parent follow-through at an eligible idle boundary; see read_delegation_skill for limits.",
"Report a child thread's state (queued, running, completed, interrupted, error, archived) and whether it is waiting on an approval or a question. Leave it out or pass 0 to read the state now; any positive value waits the full 45 seconds and returns the moment the child changes, so a shorter wait gains nothing; do not call it in a loop. Completed children and children needing approval or input return immediately. Do independent work before checking again; avoid short polling and unchanged progress updates. While Pylon delegation is enabled, child lifecycle changes can queue automatic parent follow-through at an eligible idle boundary; see read_delegation_skill for limits.",
parameters: DelegatedThreadStatusInput,
success: DelegatedThreadStatusResult,
failure: DelegationToolError,
Expand Down
15 changes: 15 additions & 0 deletions apps/server/src/mcp/toolkits/pair/handlers.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -923,6 +923,21 @@ describe("pair_await", () => {
}),
);

it.effect("turns the third instant read of a running executor into a full wait", () =>
Effect.gen(function* () {
const harness = yield* makeHarness({ shells: [makeShell(LEAD_ID), runningExecutor()] });
const instant = { state: "running", waitedSeconds: 0 };
expect(yield* harness.call("pair_await", { maxSeconds: 0 })).toMatchObject(instant);
expect(yield* harness.call("pair_await", { maxSeconds: 0 })).toMatchObject(instant);
const third = yield* Effect.forkChild(harness.call("pair_await", { maxSeconds: 0 }));
yield* TestClock.adjust(`${MAX_PAIR_AWAIT_SECONDS} seconds`);
// The lead in this harness is not Codex, so its cap is 45 seconds.
expect(yield* Fiber.join(third)).toMatchObject({ state: "running", waitedSeconds: 45 });
// The blocked call cleared the count, so a glance works again.
expect(yield* harness.call("pair_await", { maxSeconds: 0 })).toMatchObject(instant);
}),
);

it.effect("reads the state without waiting only when asked for zero seconds", () =>
Effect.gen(function* () {
const harness = yield* makeHarness({ shells: [makeShell(LEAD_ID), runningExecutor()] });
Expand Down
12 changes: 10 additions & 2 deletions apps/server/src/mcp/toolkits/pair/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ import {
isPairLeadSupported,
normalizeProtectedPath,
latestTurnCheckpoints,
pairAwaitBudgetSeconds,
pairAwaitPlan,
pairAwaitCapSeconds,
pairExecutorThreadId,
pairExecutorTitle,
Expand Down Expand Up @@ -136,6 +136,7 @@ const make = Effect.gen(function* () {
// In memory on purpose, lost on restart, and pair_await then reports null
// rather than guessing.
const protectedRecords = new Map<ThreadId, ReadonlyArray<ProtectedPathRecord>>();
const instantReadCounts = new Map<ThreadId, number>();

// One permit per lead: tool calls may run concurrently, and executor lookups
// and busy checks are check-then-act.
Expand Down Expand Up @@ -514,7 +515,14 @@ const make = Effect.gen(function* () {
const leadProvider = providerList.find((p) => p.instanceId === scope.providerInstanceId);
const driver = leadProvider?.driver;
const cap = pairAwaitCapSeconds(driver);
const waitBudget = pairAwaitBudgetSeconds(input.maxSeconds, cap);
const plan = pairAwaitPlan({
requestedSeconds: input.maxSeconds,
capSeconds: cap,
running: state === "running",
instantReads: instantReadCounts.get(executorId) ?? 0,
});
instantReadCounts.set(executorId, plan.instantReads);
const waitBudget = plan.budgetSeconds;

const startedAt = yield* Clock.currentTimeMillis;
const deadline = startedAt + waitBudget * 1_000;
Expand Down
33 changes: 33 additions & 0 deletions apps/server/src/mcp/toolkits/pair/logic.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,8 @@ import {
derivePairExecutorState,
isPairLeadSupported,
latestTurnCheckpoints,
MAX_INSTANT_READS,
pairAwaitPlan,
pairAwaitBudgetSeconds,
normalizeProtectedPath,
isPairExecutorThreadId,
Expand Down Expand Up @@ -255,3 +257,34 @@ describe("the files a lead is shown", () => {
expect(latestTurnCheckpoints([first], null)).toEqual([]);
});
});

describe("pair await plan", () => {
const cap = { capSeconds: 45 };
it("lets a lead glance at a running executor twice, then makes it wait", () => {
expect(MAX_INSTANT_READS).toBe(2);
expect(pairAwaitPlan({ ...cap, requestedSeconds: 0, running: true, instantReads: 0 })).toEqual({
budgetSeconds: 0,
instantReads: 1,
});
expect(pairAwaitPlan({ ...cap, requestedSeconds: 0, running: true, instantReads: 1 })).toEqual({
budgetSeconds: 0,
instantReads: 2,
});
expect(pairAwaitPlan({ ...cap, requestedSeconds: 0, running: true, instantReads: 2 })).toEqual({
budgetSeconds: 45,
instantReads: 0,
});
});

it("clears the count on any wait and whenever the executor is not running", () => {
expect(
pairAwaitPlan({ ...cap, requestedSeconds: undefined, running: true, instantReads: 2 }),
).toEqual({ budgetSeconds: 45, instantReads: 0 });
expect(pairAwaitPlan({ ...cap, requestedSeconds: 10, running: true, instantReads: 1 })).toEqual(
{ budgetSeconds: 45, instantReads: 0 },
);
expect(pairAwaitPlan({ ...cap, requestedSeconds: 0, running: false, instantReads: 2 })).toEqual(
{ budgetSeconds: 0, instantReads: 0 },
);
});
});
30 changes: 30 additions & 0 deletions apps/server/src/mcp/toolkits/pair/logic.ts
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,36 @@ export function derivePairExecutorState(shell: OrchestrationThreadShell): PairEx
return state;
}

/** Instant reads of a running executor a lead may make in a row before one blocks. */
export const MAX_INSTANT_READS = 2;

/**
* What one `pair_await` call does, given how many instant reads of a running
* executor came straight before it. An instant read is `maxSeconds: 0`. The
* third in a row waits the whole cap instead, so a lead cannot poll; the call
* still returns the moment the executor changes. Any wait, and any call that
* finds the executor not running, clears the count.
*/
export function pairAwaitPlan(input: {
readonly requestedSeconds: number | undefined;
readonly capSeconds: number;
readonly running: boolean;
readonly instantReads: number;
}): { readonly budgetSeconds: number; readonly instantReads: number } {
if (!input.running) {
return {
budgetSeconds: pairAwaitBudgetSeconds(input.requestedSeconds, input.capSeconds),
instantReads: 0,
};
}
if (input.requestedSeconds === 0) {
return input.instantReads < MAX_INSTANT_READS
? { budgetSeconds: 0, instantReads: input.instantReads + 1 }
: { budgetSeconds: input.capSeconds, instantReads: 0 };
}
return { budgetSeconds: input.capSeconds, instantReads: 0 };
}

/**
* The checkpoints of the executor's latest turn. One executor serves every
* brief of a pair, so its checkpoints pile up; the lead is reviewing the brief
Expand Down
2 changes: 1 addition & 1 deletion apps/server/src/mcp/toolkits/pair/tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -125,7 +125,7 @@ export const PairAwaitInput = Schema.Struct({
maxSeconds: Schema.optional(
Schema.Int.check(Schema.isBetween({ minimum: 0, maximum: MAX_PAIR_AWAIT_SECONDS })).annotate({
description:
"Leave this out. The call waits as long as your provider's tool timeout allows and returns the moment the executor stops or needs the user, so a blocked call costs no tokens and a shorter wait gains nothing: any value above 0 is raised to that full wait. Pass 0 only to read the state without waiting. If the executor is still running afterwards, end your turn: Pylon wakes you when it finishes or needs the user. Do not call this in a loop.",
"Leave this out. The call waits as long as your provider's tool timeout allows and returns the moment the executor stops or needs the user, so a blocked call costs no tokens and a shorter wait gains nothing: any value above 0 is raised to that full wait. Pass 0 only to read the state without waiting. A third instant read in a row of a running executor waits the full time instead. If the executor is still running afterwards, end your turn: Pylon wakes you when it finishes or needs the user. Do not call this in a loop.",
}),
),
maxChars: Schema.optional(
Expand Down
8 changes: 6 additions & 2 deletions docs/internals/delegation.md
Original file line number Diff line number Diff line change
Expand Up @@ -108,7 +108,10 @@ tokens, up to a cap chosen by the lead's provider: long only where Pylon sets th
tool timeout itself. The lead does not choose the length. A live Codex lead asked for 10 to 20
seconds at a time and looped, which is polling with a model turn per call, so any request other than
an explicit 0 waits the whole cap; the call still returns the moment the executor changes state, and
0 reads the state without waiting. Its `filesChanged` covers the executor's latest turn only: one
0 reads the state without waiting. Instant reads are rationed too: a lead may glance at a running
executor twice in a row, and a third `maxSeconds: 0` waits the whole cap instead (`pairAwaitPlan`), so
no argument turns the tool into a poll. Any wait, or finding the executor not running, clears the
count, which lives in memory beside the protected-path record. Its `filesChanged` covers the executor's latest turn only: one
executor serves every brief of a pair, so listing all of its checkpoints would hand the lead files it
reviewed several briefs ago. `turnCount` still counts them all.

Expand Down Expand Up @@ -185,7 +188,8 @@ brief before anything moves, and a steer never moves a running executor.
- The per-parent semaphore that serializes delegation is in memory, which is enough because one
server process owns the orchestration engine.
- Waiting is polling of the projection inside the tool call, bounded by wall-clock time to 45
seconds; already-settled children and pending approvals/input return immediately. Prime Agent's MCP client cancels any call after 60 seconds, measured in a live run.
seconds. `delegated_thread_status` does not let the caller pick a shorter wait either: omitted or 0
reads the state, any positive value waits the whole 45 seconds (`delegatedStatusWaitSeconds`); already-settled children and pending approvals/input return immediately. Prime Agent's MCP client cancels any call after 60 seconds, measured in a live run.
Pylon disables Prime's autonomous continuation, so a parent cannot be woken when a child finishes.
- The per-parent semaphores are never evicted; one small entry per thread that has delegated.
- Sends and interrupts take the same per-parent gate as delegation, so an interrupt waits behind a
Expand Down
Loading