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 @@ -453,6 +453,7 @@ it.effect(
"cp_agent_list",
"cp_agent_read",
"cp_agent_stop",
"cp_agent_settle",
]),
);

Expand Down
153 changes: 153 additions & 0 deletions apps/server/src/mcp/toolkits/agents/handlers.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ import { AgentLineage } from "../../../orchestration/agentLineage.ts";
import {
OrchestrationCommandInvariantError,
OrchestrationCommandPreviouslyRejectedError,
OrchestrationThreadSettleBlockedError,
} from "../../../orchestration/Errors.ts";
import { AGENT_RUNNING_CAP } from "../../../orchestration/agentProtocol.ts";
import {
Expand All @@ -54,6 +55,7 @@ import {
AgentCreateResult,
AgentListResult,
AgentReadResult,
AgentSettleResult,
AgentStopResult,
AgentsToolkit,
} from "./tools.ts";
Expand Down Expand Up @@ -216,6 +218,8 @@ interface HarnessInput {
readonly failOnce?: OrchestrationCommand["type"];
/** Rejects this command type once; like the engine, its id then stays rejected. */
readonly rejectOnce?: OrchestrationCommand["type"];
/** Rejects every `thread.settle` as the engine does for a thread that became busy. */
readonly blockSettle?: boolean;
/** Holds `thread.create` until the gate opens; `reached` fires when it gets there. */
readonly createGate?: {
readonly reached: Deferred.Deferred<void>;
Expand Down Expand Up @@ -268,6 +272,9 @@ const makeHarness = Effect.fn("makeAgentsToolkitHarness")(function* (input: Harn
yield* Deferred.succeed(input.createGate.reached, undefined);
yield* Deferred.await(input.createGate.open);
}
if (command.type === "thread.settle" && input.blockSettle === true) {
return yield* new OrchestrationThreadSettleBlockedError({ threadId: command.threadId });
}
if (command.type === failOnce) {
failOnce = undefined;
return yield* Effect.die(new Error(`${command.type} failed`));
Expand Down Expand Up @@ -398,6 +405,8 @@ const makeHarness = Effect.fn("makeAgentsToolkitHarness")(function* (input: Harn
run(toolkit.handle("cp_agent_read", params), AgentReadResult),
stop: (params: Parameters<typeof toolkit.handle<"cp_agent_stop">>[1]) =>
run(toolkit.handle("cp_agent_stop", params), AgentStopResult),
settle: (params: Parameters<typeof toolkit.handle<"cp_agent_settle">>[1]) =>
run(toolkit.handle("cp_agent_settle", params), AgentSettleResult),
};
});

Expand Down Expand Up @@ -1179,3 +1188,147 @@ describe("cp_agent_stop", () => {
}),
);
});

describe("cp_agent_settle", () => {
const helper = makeThread("helper", { title: "Helper" });
const foreign = makeThread("foreign", { projectId: OTHER_PROJECT_ID, title: "Foreign" });

it.effect("settles an idle agent with the same command the Settle button sends", () =>
Effect.gen(function* () {
const harness = yield* makeHarness({
caller: COORDINATOR_ID,
threads: [coordinator, helper],
});

const result = yield* harness.settle({ agent: "helper" });

expect(result).toEqual({ threadId: "helper", settled: true, archived: false });
expect(harness.commands).toEqual([
{
type: "thread.settle",
commandId: expect.stringMatching(/^mcp-agent-settle:helper:/),
threadId: "helper",
},
]);
}),
);

it.effect("settles an agent whose turn completed, and archives on request", () =>
Effect.gen(function* () {
const harness = yield* makeHarness({
caller: COORDINATOR_ID,
threads: [coordinator, helper],
messages: {
helper: [agentRequest("request-1"), makeMessage("result-1", "assistant", "Done.")],
},
turns: {
helper: [
{
turnId: "turn-1",
request: "request-1",
result: "result-1",
state: "completed",
requestedAt: "2026-09-01T00:00:00.000Z",
},
],
},
});

const result = yield* harness.settle({ agent: helper.id, archive: true });

expect(result).toEqual({ threadId: "helper", settled: true, archived: true });
expect(commandTypes(harness.commands)).toEqual(["thread.settle", "thread.archive"]);
}),
);

it.effect("refuses a busy agent with a clear error and dispatches nothing", () =>
Effect.gen(function* () {
const harness = yield* makeHarness({
caller: COORDINATOR_ID,
threads: [coordinator, running("helper", { title: "Helper" })],
});

const error = yield* harness.settle({ agent: "helper" }).pipe(Effect.flip);

expect(error._tag).toBe("AgentBusyError");
expect(error.message).toContain("cannot be settled while it has work in flight");
expect(error.message).toContain("cp_agent_stop");
expect(harness.commands).toEqual([]);
}),
);

it.effect("refuses an agent with a first message the provider has not picked up yet", () =>
Effect.gen(function* () {
const harness = yield* makeHarness({ caller: COORDINATOR_ID, threads: [coordinator] });
const created = yield* harness.create({ title: "Pricing", message: "Compare plans." });
const before = harness.commands.length;

const error = yield* harness.settle({ agent: created.threadId }).pipe(Effect.flip);

expect(error._tag).toBe("AgentBusyError");
expect(harness.commands).toHaveLength(before);
}),
);

it.effect("reports a busy agent when the engine blocks the settle after the read", () =>
Effect.gen(function* () {
const harness = yield* makeHarness({
caller: COORDINATOR_ID,
threads: [coordinator, helper],
blockSettle: true,
});

const error = yield* harness.settle({ agent: "helper" }).pipe(Effect.flip);

expect(error._tag).toBe("AgentBusyError");
}),
);

it.effect("refuses an agent the caller does not manage", () =>
Effect.gen(function* () {
const harness = yield* makeHarness({
caller: COORDINATOR_ID,
threads: [coordinator, foreign],
});

const other = yield* harness.settle({ agent: foreign.id }).pipe(Effect.flip);
const self = yield* harness.settle({ agent: COORDINATOR_ID }).pipe(Effect.flip);
const missing = yield* harness.settle({ agent: "Nobody" }).pipe(Effect.flip);

expect(other._tag).toBe("NotYourAgentError");
expect(self._tag).toBe("NotYourAgentError");
expect(missing._tag).toBe("AgentNotFoundError");
expect(harness.commands).toEqual([]);
}),
);

it.effect("is coordinator only: a standing agent cannot settle its own one-off agents", () =>
Effect.gen(function* () {
const harness = yield* makeHarness({
caller: RESEARCH_ID,
threads: [coordinator, research, helper],
creators: { helper: RESEARCH_ID },
});

const error = yield* harness.settle({ agent: "helper" }).pipe(Effect.flip);

expect(error._tag).toBe("SettleCoordinatorOnlyError");
expect(harness.commands).toEqual([]);
}),
);

it.effect("refuses a standing agent, which never settles", () =>
Effect.gen(function* () {
const harness = yield* makeHarness({
caller: COORDINATOR_ID,
threads: [coordinator, research],
});

const error = yield* harness.settle({ agent: "Research" }).pipe(Effect.flip);

expect(error._tag).toBe("AgentToolFailedError");
expect(error.message).toContain("standing");
expect(harness.commands).toEqual([]);
}),
);
});
66 changes: 65 additions & 1 deletion apps/server/src/mcp/toolkits/agents/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,10 @@ import {
agentDeliveryStartId,
agentStopId,
} from "../../../orchestration/agentProtocol.ts";
import { OrchestrationCommandPreviouslyRejectedError } from "../../../orchestration/Errors.ts";
import {
OrchestrationCommandPreviouslyRejectedError,
OrchestrationThreadSettleBlockedError,
} from "../../../orchestration/Errors.ts";
import * as OrchestrationEngine from "../../../orchestration/Services/OrchestrationEngine.ts";
import * as ProjectionSnapshotQuery from "../../../orchestration/Services/ProjectionSnapshotQuery.ts";
import { threadHasQueuedTurnStart } from "../../../orchestration/ThreadSettlementPolicy.ts";
Expand All @@ -58,13 +61,15 @@ import {
import { makeScheduleHandlers } from "./scheduleHandlers.ts";
import {
AgentAmbiguousError,
AgentBusyError,
AgentNotFoundError,
AgentsToolkit,
AgentsUnavailableError,
AgentToolFailedError,
ClientRequestIdConflictError,
ConcurrencyLimitError,
NotYourAgentError,
SettleCoordinatorOnlyError,
StandingAgentNotAllowedError,
UnknownModelError,
type AgentPhase,
Expand All @@ -82,6 +87,7 @@ const readFailed = (cause: unknown) =>

const START_FAILED = "Could not start the agent.";
const isPreviouslyRejected = Schema.is(OrchestrationCommandPreviouslyRejectedError);
const isSettleBlocked = Schema.is(OrchestrationThreadSettleBlockedError);

/** Keeps interrupts as interrupts; any other dispatch failure becomes a tool error. */
const dispatchFailed =
Expand Down Expand Up @@ -586,6 +592,64 @@ const make = Effect.gen(function* () {
archived: input.archive === true,
};
}),

cp_agent_settle: (input) =>
Effect.gen(function* () {
const manager = yield* requireManager();
// Settling is the coordinator's call: a standing agent manages one-offs
// it started, but does not decide when their work is done.
if (manager.role !== "coordinator") return yield* new SettleCoordinatorOnlyError();
const agent = yield* requireAgent(manager, input.agent);
if (isStandingAgent(manager.project, agent)) {
return yield* new AgentToolFailedError({
detail: `'${input.agent}' is a standing agent, and standing agents do not settle.`,
});
}
const turns = yield* turnRows
.listByThreadId({ threadId: agent.id })
.pipe(Effect.mapError(readFailed));
const phase = phaseOf(manager, agent, yield* nowIso);
const busy =
phase === "starting" ||
phase === "running" ||
phase === "waiting_for_approval" ||
phase === "waiting_for_input";
if (busy || hasPendingTurnStart(turns)) {
return yield* new AgentBusyError({
agent: input.agent,
phase: busy ? phase : "starting",
});
}
// The same command the Settle button sends. The decider re-checks that
// the agent is idle, so a turn that began after the read above blocks it.
yield* engine
.dispatch({
type: "thread.settle",
commandId: CommandId.make(`mcp-agent-settle:${agent.id}:${yield* randomUuid}`),
threadId: agent.id,
})
.pipe(
Effect.catchCause(
(cause): Effect.Effect<never, AgentBusyError | AgentToolFailedError> => {
const error = Cause.findErrorOption(cause);
return Option.isSome(error) && isSettleBlocked(error.value)
? Effect.fail(new AgentBusyError({ agent: input.agent, phase: "busy" }))
: dispatchFailed("Could not settle the agent.")(cause);
},
),
);
if (input.archive === true) {
yield* dispatch(
{
type: "thread.archive",
commandId: CommandId.make(`mcp-agent-archive:${agent.id}:${yield* randomUuid}`),
threadId: agent.id,
},
"Settled the agent, but could not archive it.",
);
}
return { threadId: agent.id, settled: true as const, archived: input.archive === true };
}),
});
});

Expand Down
50 changes: 50 additions & 0 deletions apps/server/src/mcp/toolkits/agents/tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,24 @@ export class NotYourAgentError extends Schema.TaggedError<NotYourAgentError>()(
}
}

export class SettleCoordinatorOnlyError extends Schema.TaggedError<SettleCoordinatorOnlyError>()(
"SettleCoordinatorOnlyError",
{},
) {
override get message(): string {
return "Only the Project's coordinator can settle agents.";
}
}

export class AgentBusyError extends Schema.TaggedError<AgentBusyError>()("AgentBusyError", {
agent: Schema.String,
phase: Schema.String,
}) {
override get message(): string {
return `'${this.agent}' is ${this.phase} and cannot be settled while it has work in flight. Wait for its report, or stop it with cp_agent_stop.`;
}
}

export class AgentToolFailedError extends Schema.TaggedError<AgentToolFailedError>()(
"AgentToolFailedError",
{ detail: Schema.String, cause: Schema.optional(Schema.Defect()) },
Expand All @@ -120,6 +138,8 @@ export const AgentToolError = Schema.Union([
AgentNotFoundError,
AgentAmbiguousError,
NotYourAgentError,
SettleCoordinatorOnlyError,
AgentBusyError,
AgentToolFailedError,
]);
export type AgentToolError = typeof AgentToolError.Type;
Expand Down Expand Up @@ -268,6 +288,21 @@ export const AgentStopResult = Schema.Struct({
});
export type AgentStopResult = typeof AgentStopResult.Type;

export const AgentSettleInput = Schema.Struct({
agent: AgentRef,
archive: Schema.optional(
Schema.Boolean.annotate({ description: "Also archive the agent's thread." }),
),
});
export type AgentSettleInput = typeof AgentSettleInput.Type;

export const AgentSettleResult = Schema.Struct({
threadId: ThreadId,
settled: Schema.Literal(true),
archived: Schema.Boolean,
});
export type AgentSettleResult = typeof AgentSettleResult.Type;

const ScheduleRef = TrimmedNonEmptyString.annotate({
description: "The schedule's id, from cp_schedule_list.",
});
Expand Down Expand Up @@ -395,6 +430,20 @@ const AgentStopTool = Tool.make("cp_agent_stop", {
.annotate(Tool.Idempotent, true)
.annotate(Tool.OpenWorld, false);

const AgentSettleTool = Tool.make("cp_agent_settle", {
description:
"Settle an idle agent you manage once its work is done, as the Settle button does. Fails while the agent is working. A settled agent stays readable, and a new message wakes it.",
parameters: AgentSettleInput,
success: AgentSettleResult,
failure: AgentToolError,
dependencies,
})
.annotate(Tool.Title, "Settle an agent")
.annotate(Tool.Readonly, false)
.annotate(Tool.Destructive, false)
.annotate(Tool.Idempotent, true)
.annotate(Tool.OpenWorld, false);

const ScheduleListTool = Tool.make("cp_schedule_list", {
description: "List this Project's schedules with their prompts, cadence and last run.",
success: ScheduleListResult,
Expand Down Expand Up @@ -453,6 +502,7 @@ export const AgentsToolkit = Toolkit.make(
AgentListTool,
AgentReadTool,
AgentStopTool,
AgentSettleTool,
ScheduleListTool,
ScheduleCreateTool,
ScheduleUpdateTool,
Expand Down
Loading
Loading