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: 7 additions & 0 deletions packages/megazord/bin/paperclip-mcp.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,12 @@ const TOOLS = {
description: "Force a harness (optional).",
},
model: { type: "string", description: "Pin the thread model (optional)." },
instance: {
type: "string",
description:
"Pin one account by instanceId (optional; see paperclip_accounts), e.g. " +
"claudeAgent_claude_pessoal. Never the SSB account for general work.",
},
project_title: { type: "string", description: "Target project by exact title (optional)." },
title: { type: "string", description: "Thread title shown in the cockpit (optional)." },
runtime: { type: "string", description: "Runtime mode (default full-access)." },
Expand All @@ -103,6 +109,7 @@ const TOOLS = {
runtimeMode: a.runtime ?? "full-access",
...(a.driver ? { driver: a.driver } : {}),
...(a.model ? { model: a.model } : {}),
...(a.instance ? { instance: a.instance } : {}),
...(a.project_title ? { projectTitle: a.project_title } : {}),
...(a.title ? { title: a.title } : {}),
});
Expand Down
67 changes: 67 additions & 0 deletions packages/megazord/src/dispatch.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -141,6 +141,35 @@ describe("selectInstance — general pool (capability + quota)", () => {
});
});

describe("selectInstance — account pin (instance)", () => {
const WITH_PESSOAL: ReadonlyArray<ProviderInstanceAccount> = [
...ACCOUNTS,
{
instanceId: "claudeAgent_claude_pessoal",
driver: "claudeAgent",
displayName: "Claude Pessoal",
enabled: true,
},
];

it("routes to the pinned second login instead of the first in config order", () => {
const d = selectInstance({
accounts: WITH_PESSOAL,
scope: "general",
driver: "claudeAgent",
instance: "claudeAgent_claude_pessoal",
});
expect(d.instanceId).toBe("claudeAgent_claude_pessoal");
expect(d.reason).toContain("instance pinned=claudeAgent_claude_pessoal");
});

it("refuses to pin the SSB account for general work", () => {
expect(() =>
selectInstance({ accounts: WITH_PESSOAL, scope: "general", instance: "claudeAgent" }),
).toThrow(/not an enabled general account/);
});
});

describe("selectInstance — connector capability (needs)", () => {
it("routes a needs request to the instance exposing the connector", () => {
const d = selectInstance({
Expand Down Expand Up @@ -214,6 +243,7 @@ function makeFetchStub(): {
id: threadId,
deletedAt: threadId === "deleted" ? "2026-09-13T00:00:00.000Z" : null,
archivedAt: null,
modelSelection: { instanceId: "claudeAgent_claude_capiva", model: "claude-opus-5-5" },
latestTurn: {
turnId: "turn-1",
state: "completed",
Expand Down Expand Up @@ -536,6 +566,43 @@ describe("continuing a thread", () => {
});
});

describe("continuing a thread keeps its account", () => {
it("sends the turn on the thread's own account even when the policy would pick another", async () => {
const { fetchImpl, calls } = makeFetchStub();
const c = new MegazordT3DispatchClient({
origin: "http://127.0.0.1:3773",
token: "TESTTOKEN",
environmentId: "env-1",
accounts: ACCOUNTS, // config order puts codex first
fetchImpl,
});
// The stub never completes a new turn: a short wait keeps the test fast; only the account of the sent turn matters here.
const out = await c.sendTurnAndAwait({ threadId: "th-existing", task: "status?", timeoutMs: 200, pollMs: 20 });
expect(out.threadId).toBe("th-existing");
const turn = calls.find((call) => call.url.endsWith("/dispatch"))!.body as {
modelSelection: { instanceId: string; model: string };
};
expect(turn.modelSelection).toEqual({
instanceId: "claudeAgent_claude_capiva",
model: "claude-opus-5-5",
});
});

it("refuses to move a thread to another account", async () => {
const { fetchImpl } = makeFetchStub();
const c = new MegazordT3DispatchClient({
origin: "http://127.0.0.1:3773",
token: "TESTTOKEN",
environmentId: "env-1",
accounts: ACCOUNTS,
fetchImpl,
});
await expect(
c.dispatch({ task: "x", scope: "general", threadId: "th-existing", instance: "codex" }),
).rejects.toThrow(/cannot move it to another account/);
});
});

describe("a turn blocked on a question", () => {
const ASK = {
kind: "user-input.requested",
Expand Down
110 changes: 104 additions & 6 deletions packages/megazord/src/dispatch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -165,6 +165,12 @@ export interface SelectInstanceInput {
readonly driver?: ProviderDriver;
/** Required exclusive connector/MCP the task needs. */
readonly needs?: string;
/**
* Pin ONE account (instanceId), e.g. to spread load over a second login of the
* same harness. Applied after the NDA gate, so a `general` request can never
* pin the SSB account; an id outside the allowed set is a refusal.
*/
readonly instance?: string;
/** Live quota, best-effort. Absent/empty = unknown → deterministic tie-break. */
readonly usage?: ReadonlyArray<InstanceUsage>;
/** Connector inventory per instance, for `needs` routing. */
Expand Down Expand Up @@ -396,6 +402,19 @@ export function selectInstance(input: SelectInstanceInput): DispatchDecision {
}
}

// ── Gate 2c: explicit account pin ──────────────────────────────────────────
const pinnedInstance = (input.instance ?? "").trim();
if (pinnedInstance !== "") {
candidates = candidates.filter((a) => a.instanceId === pinnedInstance);
if (candidates.length === 0) {
throw new MegazordDispatchError(
`refused: account '${pinnedInstance}' is not an enabled ${input.scope} account` +
(input.driver !== undefined ? ` with harness '${input.driver}'` : ""),
{ refusal: true },
);
}
}

// ── Gate 3: quota (least-loaded, deterministic tie-break) ──────────────────
const { account, load } = pickLeastLoaded(candidates, input.usage, saturationPercent);
const model = resolveModel(account, modelByDriver, input.model, input.modelsByDriver);
Expand All @@ -405,6 +424,7 @@ export function selectInstance(input: SelectInstanceInput): DispatchDecision {
if ((input.model ?? "").trim() !== "") reasonParts.push(`model pinned=${input.model}`);
if (input.driver !== undefined) reasonParts.push(`driver=${input.driver}`);
if (input.needs) reasonParts.push(`needs=${input.needs}`);
if (pinnedInstance !== "") reasonParts.push(`instance pinned=${pinnedInstance}`);
reasonParts.push(
load !== undefined
? `least-loaded (${load}% used)`
Expand Down Expand Up @@ -673,6 +693,8 @@ export interface DispatchRequest {
readonly driver?: ProviderDriver;
/** Require an exclusive connector. */
readonly needs?: string;
/** Pin one account (instanceId); see {@link SelectInstanceInput.instance}. */
readonly instance?: string;
/** Pin the model for this thread (validated against the harness's manifest). */
readonly model?: string;
/**
Expand Down Expand Up @@ -865,7 +887,7 @@ export class MegazordT3DispatchClient {
* dispatch. This is the cheap "prove the selection" path — safe for SSB.
*/
async select(
request: Pick<DispatchRequest, "scope" | "driver" | "needs" | "model">,
request: Pick<DispatchRequest, "scope" | "driver" | "needs" | "model" | "instance">,
): Promise<DispatchDecision> {
const pinned = (request.model ?? "").trim();
const [accounts, usage, modelsByDriver] = await Promise.all([
Expand All @@ -882,6 +904,7 @@ export class MegazordT3DispatchClient {
...(pinned !== "" ? { model: pinned, modelsByDriver } : {}),
...(request.driver !== undefined ? { driver: request.driver } : {}),
...(request.needs !== undefined ? { needs: request.needs } : {}),
...(request.instance !== undefined ? { instance: request.instance } : {}),
...(this.capabilities !== undefined ? { capabilities: this.capabilities } : {}),
});
}
Expand All @@ -897,7 +920,14 @@ export class MegazordT3DispatchClient {
});
}
const mode = request.mode ?? "full";
const decision = await this.select(request);
const continuingId = (request.threadId ?? "").trim();
// A follow-up turn stays on the account/model the thread already runs on:
// re-routing it would send the turn to whichever account the policy picks
// now (another login, even another harness), which the thread cannot use.
const decision =
continuingId !== "" && mode === "full"
? await this.threadDecision(continuingId, request)
: await this.select(request);

if (mode === "select") {
return {
Expand Down Expand Up @@ -1039,6 +1069,62 @@ export class MegazordT3DispatchClient {
});
}

/**
* The routing decision of an existing thread: its own `modelSelection`. An
* explicit `instance`/`model` in the request that disagrees is a refusal.
*/
private async threadDecision(
threadId: string,
request: Pick<DispatchRequest, "scope" | "instance" | "model">,
): Promise<DispatchDecision> {
const [origin, token, accounts] = await Promise.all([
this.origin(),
this.token(),
this.accounts(),
]);
let thread: Awaited<ReturnType<MegazordT3DispatchClient["readThread"]>>;
try {
thread = await this.readThread(origin, token, threadId);
} catch {
throw new MegazordDispatchError(
`refused: thread '${threadId}' is not open on this machine (gone or deleted)`,
{ refusal: true },
);
}
const selection = thread.modelSelection;
if (selection === undefined) {
throw new MegazordDispatchError(
`refused: thread '${threadId}' carries no modelSelection to continue on`,
{ refusal: true },
);
}
const account = accounts.find((a) => a.instanceId === selection.instanceId);
const pinned = (request.instance ?? "").trim();
const model = (request.model ?? "").trim();
if ((pinned !== "" && pinned !== selection.instanceId) || (model !== "" && model !== selection.model)) {
throw new MegazordDispatchError(
`refused: thread '${threadId}' runs on ${selection.instanceId}/${selection.model}; ` +
"a follow-up turn cannot move it to another account or model",
{ refusal: true },
);
}
const scope = account !== undefined && this.ssbMatcher(account) ? "ssb" : "general";
if (request.scope !== undefined && request.scope !== scope) {
throw new MegazordDispatchError(
`refused: thread '${threadId}' is a ${scope} thread, not ${request.scope}`,
{ refusal: true },
);
}
return {
instanceId: selection.instanceId,
driver: account?.driver ?? "",
model: selection.model,
scope,
reason: `continuing thread on its own account (${selection.instanceId})`,
consideredInstanceIds: [selection.instanceId],
};
}

/**
* Is this thread still a place a turn can land? A caller holding a remembered
* thread id needs to tell "keep the conversation" from "it is gone, start over".
Expand Down Expand Up @@ -1186,8 +1272,8 @@ export class MegazordT3DispatchClient {
readonly threadId: string;
/** The user turn text to append. */
readonly task: string;
/** NDA/quota scope — the SAME the thread was created under. */
readonly scope: DispatchScope;
/** NDA/quota scope — the SAME the thread was created under. Absent = the thread's own. */
readonly scope?: DispatchScope;
/** Require a specific harness (match the thread's). */
readonly driver?: ProviderDriver;
/** Pin the model for this turn (validated against the harness's manifest). */
Expand Down Expand Up @@ -1218,7 +1304,8 @@ export class MegazordT3DispatchClient {
const since = this.now();
const out = await this.dispatch({
task: input.task,
scope: input.scope,
// Undefined is fine on a continued turn: threadDecision takes the thread's own scope.
scope: input.scope as DispatchScope,
threadId,
...(input.driver !== undefined ? { driver: input.driver } : {}),
...(input.model !== undefined ? { model: input.model } : {}),
Expand Down Expand Up @@ -1247,6 +1334,7 @@ export class MegazordT3DispatchClient {
threadId: string,
): Promise<{
readonly latestTurn?: { turnId: string; state: string; requestedAt: string };
readonly modelSelection?: { instanceId: string; model: string };
readonly messages: ReadonlyArray<Record<string, unknown>>;
readonly activities: ReadonlyArray<Record<string, unknown>>;
readonly deletedAt: string | null;
Expand All @@ -1272,8 +1360,17 @@ export class MegazordT3DispatchClient {
: [];
const deletedAt = (thread["deletedAt"] as string | null | undefined) ?? null;
const archivedAt = (thread["archivedAt"] as string | null | undefined) ?? null;
const rawSelection = thread["modelSelection"] as Record<string, unknown> | null | undefined;
const modelSelection =
rawSelection !== null &&
typeof rawSelection === "object" &&
typeof rawSelection["instanceId"] === "string" &&
typeof rawSelection["model"] === "string"
? { instanceId: rawSelection["instanceId"], model: rawSelection["model"] }
: undefined;
const selection = modelSelection !== undefined ? { modelSelection } : {};
if (rawTurn === null || typeof rawTurn !== "object") {
return { messages, activities, deletedAt, archivedAt };
return { messages, activities, deletedAt, archivedAt, ...selection };
}
const t = rawTurn as Record<string, unknown>;
return {
Expand All @@ -1286,6 +1383,7 @@ export class MegazordT3DispatchClient {
activities,
deletedAt,
archivedAt,
...selection,
};
}

Expand Down