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
19 changes: 19 additions & 0 deletions apps/server/src/mcp/toolkits/pair/handlers.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -392,6 +392,25 @@ const tagOf = <A, E extends { readonly _tag: string }>(effect: Effect.Effect<A,
);

describe("pair toolkit gate", () => {
it.effect("lets a paired lead brief, await and stop with Pylon delegation off", () =>
Effect.gen(function* () {
// A session is paired because the user switched the pair on for this thread.
// Only an agent starting a pair by itself still needs the delegation setting.
const harness = yield* makeHarness({ shells: [makeShell(LEAD_ID), makeExecutor()] });
const paired = invocation({ capabilities: ["pair", "pull-requests"] });
expect(
yield* harness.call("pair_handoff", { messageKey: "m-1", text: "Step one" }, paired),
).toMatchObject({ accepted: true });
expect(yield* harness.call("pair_await", { maxSeconds: 0 }, paired)).toMatchObject({
threadId: EXECUTOR_ID,
});
expect(yield* harness.call("pair_start", {}, paired).pipe(Effect.flip)).toMatchObject({
_tag: "McpCapabilityUnavailableError",
capability: "delegation",
});
}),
);

it.effect("requires the delegation capability for every tool", () =>
Effect.gen(function* () {
const harness = yield* makeHarness({ shells: [makeShell(LEAD_ID), makeExecutor()] });
Expand Down
10 changes: 7 additions & 3 deletions apps/server/src/mcp/toolkits/pair/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,10 @@ const mapDispatch =
}),
);

const requirePairCapability = McpInvocationContext.requireMcpCapability("pair").pipe(
Effect.catch(() => McpInvocationContext.requireMcpCapability("delegation")),
);

const make = Effect.gen(function* () {
const engine = yield* OrchestrationEngine.OrchestrationEngineService;
const snapshots = yield* ProjectionSnapshotQuery.ProjectionSnapshotQuery;
Expand Down Expand Up @@ -343,7 +347,7 @@ const make = Effect.gen(function* () {
readonly steer?: boolean | undefined;
}) =>
Effect.gen(function* () {
const scope = yield* McpInvocationContext.requireMcpCapability("delegation");
const scope = yield* requirePairCapability;
if (!isValidDelegationKey(input.messageKey)) {
return yield* new PairKeyInvalidError();
}
Expand Down Expand Up @@ -495,7 +499,7 @@ const make = Effect.gen(function* () {
readonly maxChars?: number | undefined;
}) =>
Effect.gen(function* () {
const scope = yield* McpInvocationContext.requireMcpCapability("delegation");
const scope = yield* requirePairCapability;
const executorId = yield* executorIdFor(scope.threadId);
const existing = yield* findExecutor(executorId);
if (Option.isNone(existing)) {
Expand Down Expand Up @@ -601,7 +605,7 @@ const make = Effect.gen(function* () {

const pair_stop = () =>
Effect.gen(function* () {
const scope = yield* McpInvocationContext.requireMcpCapability("delegation");
const scope = yield* requirePairCapability;
const executorId = yield* executorIdFor(scope.threadId);

return yield* withLeadGate(scope.threadId)(
Expand Down
4 changes: 2 additions & 2 deletions apps/server/src/mcp/toolkits/pair/tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,8 @@
* Pair toolkit declaration: one lead thread and one persistent executor thread
* that works in the lead's worktree. The executor's id is the delegated child
* id for the reserved key `pair`, so no contract field or table records the
* link. Every tool requires the `delegation` capability, which keeps
* `enableAgentDelegation` as the single kill switch.
* link. Starting a pair requires the `delegation` capability; driving an
* existing pair accepts either the `pair` or `delegation` capability.
*
* @module mcp/toolkits/pair/tools
*/
Expand Down
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { pairExecutorThreadId } from "@t3tools/shared/delegatedThreads";
import {
DEFAULT_SERVER_SETTINGS,
EventId,
Expand Down Expand Up @@ -447,6 +448,28 @@ describe("DelegationFollowThroughReactor", () => {
}),
),
);
it.effect("still wakes a lead for its pair executor when delegation is disabled", () =>
Effect.scoped(
Effect.gen(function* () {
const executor = pairExecutorThreadId(PARENT);
const h = yield* makeHarness(
[shell(PARENT), shell(executor, true), shell(CHILD, true)],
false,
);
// A fan-out child finishing stays silent: that is what the setting is for.
h.replace(shell(CHILD));
yield* h.emit(CHILD);
assert.strictEqual(h.wakes().length, 0);
h.replace(shell(executor));
yield* h.emit(executor);
assert.strictEqual(h.wakes().length, 1);
assert.deepStrictEqual(
h.wakes()[0]!.children.map((child) => child.threadId),
[executor],
);
}),
),
);
it.effect("enforces the persisted three-turn budget without a wake loop", () =>
Effect.scoped(
Effect.gen(function* () {
Expand Down
10 changes: 8 additions & 2 deletions apps/server/src/orchestration/DelegationFollowThroughReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import {
type OrchestrationEvent,
} from "@t3tools/contracts";
import { makeDrainableWorker } from "@t3tools/shared/DrainableWorker";
import { pairExecutorThreadId } from "@t3tools/shared/delegatedThreads";
import { resolveProjectSettings } from "@t3tools/shared/projectSettings";
import * as Cause from "effect/Cause";
import * as Context from "effect/Context";
Expand All @@ -26,6 +27,7 @@ import * as ProjectionSnapshotQuery from "./Services/ProjectionSnapshotQuery.ts"
import {
DELEGATION_OBSERVED_ACTIVITY_KIND,
delegationObservationReceipt,
followThroughChildren,
observeDelegatedChild,
isActionableDelegationObservation,
isDelegationParentEligible,
Expand Down Expand Up @@ -119,9 +121,13 @@ export const make = Effect.gen(function* () {
yield* settingsService.getSettings,
parent.projectId,
).settings;
if (!settings.enableAgentDelegation) return;
const snapshot = yield* snapshots.getShellSnapshot();
const children = snapshot.threads.filter((child) => isChildOfParent(child.id, parent.id));
const pairExecutorId = pairExecutorThreadId(parent.id);
const children = followThroughChildren({
delegationEnabled: settings.enableAgentDelegation,
pairExecutorId,
children: snapshot.threads.filter((child) => isChildOfParent(child.id, parent.id)),
});
if (children.length === 0) return;
const detailOption = yield* snapshots.getThreadDetailById(parent.id, {
activityKinds: [DELEGATION_OBSERVED_ACTIVITY_KIND, DELIVERED, PAUSED],
Expand Down
10 changes: 8 additions & 2 deletions apps/server/src/orchestration/Layers/OrchestrationEngine.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,9 @@ import type {
ThreadId,
} from "@t3tools/contracts";
import { CommandId, OrchestrationCommand } from "@t3tools/contracts";
import { pairExecutorThreadId } from "@t3tools/shared/delegatedThreads";
import { resolveProjectSettings } from "@t3tools/shared/projectSettings";
import { isFollowThroughAdmitted } from "../delegationFollowThrough.logic.ts";
import { ServerSettingsService } from "../../serverSettings.ts";
import * as Cause from "effect/Cause";
import * as Clock from "effect/Clock";
Expand Down Expand Up @@ -413,8 +415,12 @@ const makeOrchestrationEngine = Effect.gen(function* () {
return {
enabled:
parent !== undefined &&
resolveProjectSettings(settings, parent.projectId).settings
.enableAgentDelegation,
isFollowThroughAdmitted({
delegationEnabled: resolveProjectSettings(settings, parent.projectId).settings
.enableAgentDelegation,
pairExecutorId: pairExecutorThreadId(command.threadId),
childThreadIds: command.children.map((c) => c.threadId),
}),
messageExists: Option.isSome(message),
deliveredNotificationIds,
};
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,8 @@ import { deriveDelegatedThreadState } from "../mcp/toolkits/delegation/logic.ts"
import {
DELEGATION_OBSERVED_ACTIVITY_KIND,
consumedDelegationObservation,
followThroughChildren,
isFollowThroughAdmitted,
delegationObservationReceipt,
observeDelegatedChild,
isActionableDelegationObservation,
Expand Down Expand Up @@ -256,3 +258,45 @@ describe("an observation a parent read in-turn", () => {
});
});
});

describe("who may wake a parent when Pylon delegation is off", () => {
const executor = { id: "delegated:parent:pairpairpairpair" };
const fanOut = { id: "delegated:parent:0123456789abcdef" };

it("keeps every child while delegation is on, and only the pair executor while it is off", () => {
const children = [executor, fanOut];
expect(
followThroughChildren({ delegationEnabled: true, pairExecutorId: executor.id, children }),
).toEqual(children);
expect(
followThroughChildren({ delegationEnabled: false, pairExecutorId: executor.id, children }),
).toEqual([executor]);
expect(
followThroughChildren({
delegationEnabled: false,
pairExecutorId: executor.id,
children: [fanOut],
}),
).toEqual([]);
});

it("admits a wake with delegation off only when every child named is the pair executor", () => {
const base = { pairExecutorId: executor.id };
expect(
isFollowThroughAdmitted({ ...base, delegationEnabled: true, childThreadIds: [fanOut.id] }),
).toBe(true);
expect(
isFollowThroughAdmitted({ ...base, delegationEnabled: false, childThreadIds: [executor.id] }),
).toBe(true);
expect(
isFollowThroughAdmitted({
...base,
delegationEnabled: false,
childThreadIds: [executor.id, fanOut.id],
}),
).toBe(false);
expect(isFollowThroughAdmitted({ ...base, delegationEnabled: false, childThreadIds: [] })).toBe(
false,
);
});
});
31 changes: 31 additions & 0 deletions apps/server/src/orchestration/delegationFollowThrough.logic.ts
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,37 @@ export function isActionableDelegationObservation(observation: DelegationObserva
return observation.phase !== "running";
}

/**
* The children whose updates may wake this parent. The Pylon delegation setting
* covers what an agent starts on its own. A pair is switched on per thread by
* the user, so its executor wakes the lead whether or not that setting is on.
*/
export function followThroughChildren<T extends { readonly id: string }>(input: {
readonly delegationEnabled: boolean;
readonly pairExecutorId: string;
readonly children: ReadonlyArray<T>;
}): ReadonlyArray<T> {
if (input.delegationEnabled) {
return input.children;
}
return input.children.filter((child) => child.id === input.pairExecutorId);
}

/** Whether a follow-through wake naming these children may be admitted. */
export function isFollowThroughAdmitted(input: {
readonly delegationEnabled: boolean;
readonly pairExecutorId: string;
readonly childThreadIds: ReadonlyArray<string>;
}): boolean {
if (input.delegationEnabled) {
return true;
}
return (
input.childThreadIds.length > 0 &&
input.childThreadIds.every((id) => id === input.pairExecutorId)
);
}

/** The activity kind that records what the follow-through reactor last saw of a child. */
export const DELEGATION_OBSERVED_ACTIVITY_KIND = "delegation.child-state";

Expand Down
15 changes: 13 additions & 2 deletions apps/server/src/provider/Layers/ProviderService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7112,7 +7112,7 @@ describe("agent browser access", () => {
}).pipe(Effect.provide(NodeServices.layer)),
);

it.effect("marks a session as paired only when delegation is on and its executor exists", () =>
it.effect("marks a session as paired whenever its executor exists and is not archived", () =>
Effect.gen(function* () {
const lead = asThreadId("thread-pair-lead");
const executorOf = (threadId: ThreadId) =>
Expand All @@ -7127,7 +7127,18 @@ describe("agent browser access", () => {
["paired", lead, true, [executorOf(lead)], ["delegation", "pair", "pull-requests"], false],
["no executor", lead, true, [], ["delegation", "pull-requests"], false],
["only a fan-out child", lead, true, [otherChild], ["delegation", "pull-requests"], false],
["delegation off", lead, false, [executorOf(lead)], ["pull-requests"], false],
// The user switched this pair on for this thread, so it does not need the
// setting that lets agents start threads on their own.
["delegation off", lead, false, [executorOf(lead)], ["pair", "pull-requests"], false],
["delegation off, no executor", lead, false, [], ["pull-requests"], false],
[
"delegation off, pair turned off",
lead,
false,
[executorOf(lead)],
["pull-requests"],
true,
],
["the executor itself", executorOf(lead), true, [], ["pull-requests"], false],
// Turning a pair off archives an executor that has history. The lead
// must get its own subagents back, not stay in paired mode.
Expand Down
24 changes: 12 additions & 12 deletions apps/server/src/provider/Layers/ProviderService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1217,18 +1217,18 @@ const makeProviderService = Effect.fn("makeProviderService")(function* (
// A delegated child never receives delegation; its id carries the prefix.
if (access.delegation && !threadId.startsWith("delegated:")) {
capabilities.add("delegation");
if (Option.isSome(projectionQuery)) {
const executorId = pairExecutorThreadId(threadId, (input) =>
NodeCrypto.createHash("sha256").update(input).digest("hex"),
);
const executor = yield* projectionQuery.value
.getThreadShellById(executorId)
.pipe(Effect.orElseSucceed(() => Option.none()));
// An archived executor is a pair that was turned off; the lead gets its
// own subagents and the delegation instructions back.
if (Option.isSome(executor) && executor.value.archivedAt === null) {
capabilities.add("pair");
}
}
if (!threadId.startsWith("delegated:") && Option.isSome(projectionQuery)) {
const executorId = pairExecutorThreadId(threadId, (input) =>
NodeCrypto.createHash("sha256").update(input).digest("hex"),
);
const executor = yield* projectionQuery.value
.getThreadShellById(executorId)
.pipe(Effect.orElseSucceed(() => Option.none()));
// An archived executor is a pair that was turned off; the lead gets its
// own subagents and the delegation instructions back.
if (Option.isSome(executor) && executor.value.archivedAt === null) {
capabilities.add("pair");
}
}
return capabilities;
Expand Down
32 changes: 1 addition & 31 deletions apps/web/src/components/chat/pairControl.logic.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@ import { describe, expect, it } from "vite-plus/test";

import {
pairLockedReason,
pairSetupReason,
pairToggleStep,
resolveExecutorSelection,
shouldRestartLeadSession,
Expand Down Expand Up @@ -38,11 +37,7 @@ const on = (phase: Extract<PairState, { kind: "on" }>["phase"]): PairState => ({
modelSelection: SELECTION,
activity: null,
});
const base = {
executorSelection: SELECTION,
childRuntimeMode: "inherit" as const,
delegationEnabled: true,
};
const base = { executorSelection: SELECTION, childRuntimeMode: "inherit" as const };

describe("pairLockedReason", () => {
it("makes a change wait while the lead is mid-turn, and only then", () => {
Expand All @@ -56,32 +51,7 @@ describe("pairLockedReason", () => {
});
});

describe("pairSetupReason", () => {
it("asks for Pylon delegation before a pair can start, and never blocks turning one off", () => {
expect(pairSetupReason({ delegationEnabled: false, state: off })).toBe(
"Turn on Pylon delegation in Settings → Integrations to pair.",
);
expect(pairSetupReason({ delegationEnabled: true, state: off })).toBeNull();
expect(pairSetupReason({ delegationEnabled: false, state: on("idle") })).toBeNull();
});
});

describe("pairToggleStep", () => {
it("does not start a pair the server would not honor", () => {
expect(
pairToggleStep({ ...base, delegationEnabled: false, on: true, state: off, lead: lead() }),
).toBeNull();
expect(
pairToggleStep({
...base,
delegationEnabled: false,
on: false,
state: on("idle"),
lead: lead(),
}),
).toEqual({ kind: "delete", threadId: EXECUTOR });
});

it("creates the executor in the lead's location when turned on", () => {
expect(pairToggleStep({ ...base, on: true, state: off, lead: lead() })).toEqual({
kind: "create",
Expand Down
Loading
Loading