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
157 changes: 157 additions & 0 deletions packages/megazord/bin/t3-thread-turn.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,157 @@
#!/usr/bin/env node
/**
* t3-thread-turn — append ONE turn to an existing T3 thread and print its answer,
* in a single call.
*
* The atomic round-trip a thin CLIENT makes against the single OWNER-hosted
* session: where `t3-dispatch.mjs --thread <id>` starts a turn and returns as
* soon as it is STARTED (leaving the caller to run `t3-thread-wait.mjs` and
* compute its own `since`), this does both — send the turn and wait for the
* assistant reply — so a chat channel gets the answer from one process and can
* never latch onto the previous turn's text (the `since` instant is captured
* before the turn is sent). Every surface (WhatsApp bridge, T3 GUI, another
* machine) appends to the SAME thread, so memory carries across turns.
*
* A turn blocked on an approval/question never terminates, so the wait ends at its
* own ceiling with `state: "running"|"awaiting-input", timedOut: true` rather than
* hanging — the caller relays the link (or the question) and the human takes over.
*
* Usage:
* node bin/t3-thread-turn.mjs --thread <threadId> --task "<prompt>" \
* --scope ssb|general [--driver codex|claudeAgent] [--model <id>] \
* [--runtime <mode>] [--interaction default|plan] \
* [--timeout-seconds 900] [--poll-seconds 2] [--base-dir <dir>] [--json]
*
* The token is minted at runtime and NEVER written to disk.
*
* @module megazord/bin/t3-thread-turn
*/
import { MegazordDispatchError, MegazordT3DispatchClient } from "../src/dispatch.ts";

function parseArgs(argv) {
const out = { scope: "general", json: false };
for (let i = 0; i < argv.length; i++) {
const a = argv[i];
const next = () => argv[++i];
switch (a) {
case "--thread":
out.threadId = next();
break;
case "--task":
out.task = next();
break;
case "--scope":
out.scope = next();
break;
case "--driver":
out.driver = next();
break;
case "--model":
out.model = next();
break;
case "--runtime":
out.runtimeMode = next();
break;
case "--interaction":
out.interactionMode = next();
break;
case "--timeout-seconds":
out.timeoutSeconds = Number(next());
break;
case "--poll-seconds":
out.pollSeconds = Number(next());
break;
case "--base-dir":
out.baseDir = next();
break;
case "--json":
out.json = true;
break;
case "-h":
case "--help":
out.help = true;
break;
default:
if (a.startsWith("--")) {
console.error(`unknown flag: ${a}`);
process.exit(2);
}
}
}
return out;
}

const HELP = `t3-thread-turn — append a turn to an existing thread and print its answer

--thread <threadId> the existing thread to append to (required)
--task "<prompt>" the turn text (required)
--scope ssb|general NDA/quota scope, SAME the thread was created under
--driver codex|claudeAgent require a specific harness (match the thread's)
--model <id> pin the model (validated per harness)
--runtime <mode> runtime mode (default: the client's approval-required)
--interaction default|plan interaction mode
--timeout-seconds <n> give up waiting after this long (default 900)
--poll-seconds <n> poll interval (default 2)
--base-dir <dir> T3 base dir (default: ~/.t3)
--json machine-readable output
`;

async function main() {
const args = parseArgs(process.argv.slice(2));
if (args.help) {
process.stdout.write(HELP);
return;
}
if (!args.threadId || !args.task) {
console.error("--thread and --task are required");
process.exit(2);
}
if (args.scope !== "ssb" && args.scope !== "general") {
console.error("--scope must be 'ssb' or 'general'");
process.exit(2);
}

const client = new MegazordT3DispatchClient({
...(args.baseDir ? { baseDir: args.baseDir } : {}),
});

try {
const out = await client.sendTurnAndAwait({
threadId: args.threadId,
task: args.task,
scope: args.scope,
...(args.driver ? { driver: args.driver } : {}),
...(args.model ? { model: args.model } : {}),
...(args.runtimeMode ? { runtimeMode: args.runtimeMode } : {}),
...(args.interactionMode ? { interactionMode: args.interactionMode } : {}),
...(Number.isFinite(args.timeoutSeconds)
? { timeoutMs: Math.max(1, args.timeoutSeconds) * 1000 }
: {}),
...(Number.isFinite(args.pollSeconds)
? { pollMs: Math.max(1, args.pollSeconds) * 1000 }
: {}),
});
if (args.json) {
process.stdout.write(JSON.stringify(out, null, 2) + "\n");
return;
}
process.stdout.write(
`route: ${out.instanceId} (${out.driver}, model ${out.model})\n` +
`state: ${out.state}${out.timedOut ? " (timed out)" : ""}\n` +
`thread: ${out.url}\n`,
);
if (out.text !== "") process.stdout.write(`\n${out.text}\n`);
} catch (e) {
if (e instanceof MegazordDispatchError) {
console.error(`t3-thread-turn: ${e.message}`);
if (e.detail?.responseText) console.error(` server: ${e.detail.responseText}`);
process.exit(e.refusal ? 3 : 1);
}
throw e;
}
}

main().catch((e) => {
console.error(e?.stack ?? String(e));
process.exit(1);
});
49 changes: 49 additions & 0 deletions packages/megazord/src/dispatch.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -485,6 +485,55 @@ describe("continuing a thread", () => {
expect(out.timedOut).toBe(false);
expect(out.url).toBe("http://127.0.0.1:3773/env-1/th-existing");
});

it("sendTurnAndAwait appends one turn to the existing thread and returns its answer", async () => {
const { fetchImpl, calls } = makeFetchStub();
let seq = 0;
const c = new MegazordT3DispatchClient({
origin: "http://127.0.0.1:3773",
token: "TESTTOKEN",
environmentId: "env-1",
accounts: ACCOUNTS,
fetchImpl,
uuid: () => `uuid-${++seq}`,
now: () => "2026-09-13T00:00:00.000Z",
});
const out = await c.sendTurnAndAwait({
threadId: "th-existing",
task: "Qual numero eu pedi pra guardar?",
scope: "general",
driver: "claudeAgent",
timeoutMs: 5000,
});
// The round-trip continues the thread (turn.start only — never thread.create)…
const dispatches = calls.filter((cc) => cc.url.endsWith("/dispatch"));
expect(dispatches.map((cc) => (cc.body as { type: string }).type)).toEqual([
"thread.turn.start",
]);
expect((dispatches[0]!.body as { threadId: string }).threadId).toBe("th-existing");
// …and reads the assistant reply back, tagged with the account that ran it.
expect(out.state).toBe("completed");
expect(out.text).toBe("resposta");
expect(out.driver).toBe("claudeAgent");
expect(out.instanceId).toBe("claudeAgent_claude_capiva");
expect(out.sequence).toBe(4242);
expect(out.url).toBe("http://127.0.0.1:3773/env-1/th-existing");
});

it("sendTurnAndAwait refuses a missing threadId before any I/O", async () => {
const { fetchImpl, calls } = makeFetchStub();
const c = new MegazordT3DispatchClient({
origin: "http://127.0.0.1:3773",
token: "TESTTOKEN",
environmentId: "env-1",
accounts: ACCOUNTS,
fetchImpl,
});
await expect(
c.sendTurnAndAwait({ threadId: " ", task: "x", scope: "general" }),
).rejects.toBeInstanceOf(MegazordDispatchError);
expect(calls.length).toBe(0);
});
});

describe("a turn blocked on a question", () => {
Expand Down
91 changes: 91 additions & 0 deletions packages/megazord/src/dispatch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -739,6 +739,20 @@ export interface TurnWaitOutcome {
readonly pendingInput?: PendingUserInput;
}

/**
* The result of a full round-trip: append one turn to an existing thread AND read
* that turn's answer back. It is the {@link TurnWaitOutcome} plus which
* (instance, harness, model) the owner ran the turn on and the dispatch sequence
* of the `thread.turn.start` that opened it.
*/
export interface TurnRoundTripOutcome extends TurnWaitOutcome {
readonly instanceId: string;
readonly driver: ProviderDriver | string;
readonly model: string;
/** Dispatch sequence of the `thread.turn.start` that opened the round-trip. */
readonly sequence: number;
}

const TERMINAL_TURN_STATES: ReadonlyArray<string> = ["completed", "error", "interrupted"];
const DEFAULT_WAIT_TIMEOUT_MS = 900_000;
const DEFAULT_WAIT_POLL_MS = 2_000;
Expand Down Expand Up @@ -1149,6 +1163,83 @@ export class MegazordT3DispatchClient {
}
}

/**
* Append one turn to an EXISTING thread and return that turn's answer — the
* atomic round-trip a thin CLIENT (the WhatsApp bridge, the T3 GUI, another
* machine) makes against the single OWNER-hosted session. Exactly one thread is
* the orchestrator's persistent conversation; every surface's turn appends to
* it, so memory carries across turns and the reply flows back to whoever asked.
*
* It composes {@link dispatch} (continue the thread — `thread.turn.start`, no
* `thread.create`) with {@link awaitTurn} (poll the snapshot for the assistant
* reply). The `since` instant is captured BEFORE the turn is sent, so the wait
* can never latch onto the PREVIOUS turn's answer — the race that would make a
* continuity proof lie.
*
* Refused (never dispatched) if the thread is gone/deleted: a remembered id that
* no longer accepts turns is an error, not a silent new thread. Pass the SAME
* `scope`/`driver` the thread was created under — routing still runs to resolve
* the model selection, and a mismatched scope would resolve a wrong account.
*/
async sendTurnAndAwait(input: {
/** The existing owner-hosted thread every surface appends to. */
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;
/** Require a specific harness (match the thread's). */
readonly driver?: ProviderDriver;
/** Pin the model for this turn (validated against the harness's manifest). */
readonly model?: string;
/** Runtime mode override for this turn. */
readonly runtimeMode?: RuntimeMode;
/** Interaction mode. `plan` keeps the agent in planning (no edits). */
readonly interactionMode?: InteractionMode;
/** Wait ceiling. Default {@link DEFAULT_WAIT_TIMEOUT_MS}. */
readonly timeoutMs?: number;
/** Poll interval. Default {@link DEFAULT_WAIT_POLL_MS}. */
readonly pollMs?: number;
}): Promise<TurnRoundTripOutcome> {
const threadId = input.threadId.trim();
if (threadId === "") {
throw new MegazordDispatchError(
"refused: sendTurnAndAwait needs an existing threadId to append the turn to",
{ refusal: true },
);
}
if (input.task.trim() === "") {
throw new MegazordDispatchError("refused: task text is required for a round-trip turn", {
refusal: true,
});
}
// Captured BEFORE the turn is sent: awaitTurn ignores any turn requested
// before this instant, so it cannot return the prior turn's answer.
const since = this.now();
const out = await this.dispatch({
task: input.task,
scope: input.scope,
threadId,
...(input.driver !== undefined ? { driver: input.driver } : {}),
...(input.model !== undefined ? { model: input.model } : {}),
...(input.runtimeMode !== undefined ? { runtimeMode: input.runtimeMode } : {}),
...(input.interactionMode !== undefined ? { interactionMode: input.interactionMode } : {}),
});
const wait = await this.awaitTurn({
threadId,
since,
...(input.timeoutMs !== undefined ? { timeoutMs: input.timeoutMs } : {}),
...(input.pollMs !== undefined ? { pollMs: input.pollMs } : {}),
});
return {
...wait,
instanceId: out.instanceId,
driver: out.driver,
model: out.model,
sequence: out.sequence ?? -1,
};
}

/** GET one thread's snapshot (the same read the UI does). */
private async readThread(
origin: string,
Expand Down