diff --git a/README.md b/README.md index 7f709b7..b33a5bf 100644 --- a/README.md +++ b/README.md @@ -154,6 +154,8 @@ printf '%s' "New findings that require more work..." \ The send command does not report success from the HTTP response alone. It waits until the exact message is visible in T3's thread projection. Archived threads are rejected. +By default, `send` refuses a busy thread with `THREAD_BUSY` and dispatches nothing. A thread is busy while a turn runs or while an earlier message still waits for its turn; `error.details` says which. Wait with `threads wait` and send again, or pass `--if-busy inject` to send into the running turn. The provider then folds the message into that turn or queues it, which `--wait` follows either way. The busy check is a snapshot, not a lock, so callers that send to the same thread at once must take turns themselves. + Add `--wait` to wait for the turn that handles the message and print its reply: ```bash diff --git a/skills/t3thread/SKILL.md b/skills/t3thread/SKILL.md index 2aea6be..6446ade 100644 --- a/skills/t3thread/SKILL.md +++ b/skills/t3thread/SKILL.md @@ -90,7 +90,7 @@ Rules for sending: - A settled thread needs `--wake-settled`. The user's explicit instruction to message this thread authorizes it. - Archived threads cannot receive messages. -- If the thread is mid-turn, the provider either folds the message into the running turn or queues a new turn. `--wait` handles both. +- `send` refuses a busy thread with `THREAD_BUSY` (exit code 4): a turn runs or an earlier message waits. Wait with `threads wait`, then send. Pass `--if-busy inject` only when the user wants to steer the running turn; the provider then folds the message in or queues it, and `--wait` follows either way. - `THREAD_WAIT_TIMEOUT` (exit code 6) means the message was sent. Never resend it. Keep waiting with `t3code threads wait --thread --timeout 540`. - `THREAD_TURN_NOT_VERIFIED` (exit code 5) means T3 has not shown the message yet. Do not retry automatically; read the thread first. diff --git a/skills/use-t3code-cli/SKILL.md b/skills/use-t3code-cli/SKILL.md index 14b0ca8..a0ccb73 100644 --- a/skills/use-t3code-cli/SKILL.md +++ b/skills/use-t3code-cli/SKILL.md @@ -89,7 +89,7 @@ printf '%s' "$THREAD_MESSAGE" \ | t3code --json threads send --thread "$TARGET_THREAD_ID" --stdin ``` -Sending is an external state change. Keep the target and message within the caller's authorization. A settled thread requires interactive confirmation or `--wake-settled`; JSON and stdin workflows are non-interactive, so use that override only when waking the inspected target is authorized. Archived threads cannot receive a turn. +Sending is an external state change. Keep the target and message within the caller's authorization. A settled thread requires interactive confirmation or `--wake-settled`; JSON and stdin workflows are non-interactive, so use that override only when waking the inspected target is authorized. Archived threads cannot receive a turn. A busy thread, with a running turn or a message waiting for its turn, gets `THREAD_BUSY` unless you pass `--if-busy inject` to send into the running turn. Add `--wait` to get the reply. Give the shell call a longer timeout than `--timeout`: diff --git a/src/cli.ts b/src/cli.ts index 737c955..5de6a33 100644 --- a/src/cli.ts +++ b/src/cli.ts @@ -202,6 +202,7 @@ interface SettingsCommandOptions { interface ThreadSendCommandOptions extends PromptOptions, ThreadWaitCommandOptions, SettingsCommandOptions { wakeSettled?: boolean; wait?: boolean; + ifBusy: "reject" | "inject"; } interface ThreadRequestCommandOptions extends ThreadWaitCommandOptions { @@ -592,7 +593,12 @@ addReplyOptions( .option("--prompt-file ", "Read the message from a UTF-8 file.") .option("--stdin", "Read the message from stdin.") .option("--wake-settled", "Explicitly allow this message to wake a settled thread.") - .option("--wait", "Wait for the turn that handles the message and print its reply."), + .option("--wait", "Wait for the turn that handles the message and print its reply.") + .addOption( + new Option("--if-busy ", "reject: refuse while a turn runs or a message waits; inject: send into the running turn.") + .choices(["reject", "inject"]) + .default("reject"), + ), ), ).action((options: ThreadSendCommandOptions) => action(async () => { @@ -605,6 +611,7 @@ addReplyOptions( ...(!context.json && !options.stdin ? { confirmSettled: confirmSettledThread } : {}), ...(options.wait ? { wait: waitOptions(options) } : {}), settings: settingsChange(options), + ifBusy: options.ifBusy, }); const changed = result.settings ? describeChanges(result.settings) : ""; const sent = `${changed ? `Changed ${changed}. ` : ""}Sent message ${result.message.messageId} to thread ${result.thread.id}; T3 accepted and projected the turn.`; diff --git a/src/service.ts b/src/service.ts index 473d359..a7ecc46 100644 --- a/src/service.ts +++ b/src/service.ts @@ -11,6 +11,7 @@ import { discoverRuntime } from "./runtime.js"; import { T3ThreadApi, type ThreadSettlementState } from "./threadApi.js"; import { changeSettingsWithApi, + busyState, hasSettingsChange, settingsSummary, type ThreadSettingsChange, @@ -85,6 +86,11 @@ export interface ThreadSendOptions { wait?: ThreadWaitOptions; /** Change the thread's model, effort, speed, or modes before the message starts its turn. */ settings?: ThreadSettingsChange; + /** + * What to do when the thread is busy: `reject` (the default) refuses to send; `inject` sends into the + * running turn, where the provider folds the message in or queues it. + */ + ifBusy?: "reject" | "inject"; } interface EffectiveT3Settings { @@ -534,6 +540,23 @@ export async function sendThreadMessage(config: CliConfig, options: ThreadSendOp } } + const busy = busyState(thread); + if (busy && (options.ifBusy ?? "reject") === "reject") { + throw new CliError( + "THREAD_BUSY", + `Thread ${threadId} ${busy.turnRunning ? "is running a turn" : "has a message waiting for its turn"}. Wait for it with threads wait, or pass --if-busy inject to send into the running turn.`, + { + exitCode: 4, + details: { + threadId, + ...busy, + sessionStatus: thread.session?.status ?? null, + latestTurnState: thread.latestTurn?.state ?? null, + }, + }, + ); + } + const settings = hasSettingsChange(options.settings) ? await changeSettingsWithApi(api, adapter, thread, options.settings) : null; diff --git a/src/threadControls.ts b/src/threadControls.ts index 7b4b66d..b56c673 100644 --- a/src/threadControls.ts +++ b/src/threadControls.ts @@ -22,7 +22,7 @@ import { type ThreadWaitOptions, type ThreadWaitView, } from "./threadSupport.js"; -import { pendingRequests, type PendingQuestion, type PendingRequest } from "./transcript.js"; +import { pendingRequests, queuedMessages, type PendingQuestion, type PendingRequest } from "./transcript.js"; import type { CliConfig, InteractionMode, ModelSelection, RuntimeMode, T3Thread } from "./types.js"; /** Settings a caller asks to change on an existing thread; anything left out stays as it is. */ @@ -65,6 +65,13 @@ function liveSession(thread: T3Thread): boolean { return thread.session != null && thread.session.status !== "stopped"; } +/** What keeps a thread busy: a running turn, or messages that still wait for their turn. */ +export function busyState(thread: T3Thread): { turnRunning: boolean; queuedMessages: number } | null { + const running = turnRunning(thread); + const queued = queuedMessages(thread).length; + return running || queued > 0 ? { turnRunning: running, queuedMessages: queued } : null; +} + /** T3 binds a conversation to its provider once the thread has a session or any history. */ function conversationStarted(thread: T3Thread): boolean { return thread.session != null || thread.latestTurn != null || (thread.messages?.length ?? 0) > 0; diff --git a/src/threads.test.ts b/src/threads.test.ts index c575833..4c7ea18 100644 --- a/src/threads.test.ts +++ b/src/threads.test.ts @@ -742,6 +742,38 @@ describe("thread controls", () => { expect(harness.commands[0]).toMatchObject({ type: "thread.meta.update" }); }); + it("refuses to send into a busy thread unless the caller injects", async () => { + const harness = await testHarness([makeThread("target", running())]); + + await expect(sendThreadMessage(harness.config, { threadId: "target", prompt: "Also check the docs" })).rejects.toMatchObject({ + code: "THREAD_BUSY", + exitCode: 4, + details: { turnRunning: true, queuedMessages: 0 }, + }); + expect(harness.commands).toEqual([]); + + const injected = await sendThreadMessage(harness.config, { threadId: "target", prompt: "Also check the docs", ifBusy: "inject" }); + + expect(harness.commands).toEqual([expect.objectContaining({ type: "thread.turn.start", threadId: "target" })]); + expect(injected.verification).toMatchObject({ accepted: true }); + }); + + it("counts a message waiting for its turn as busy", async () => { + const harness = await testHarness([makeThread("target", { + latestTurn: { turnId: "turn-1", state: "completed", requestedAt: at(0), startedAt: at(0), completedAt: at(2), assistantMessageId: null }, + messages: [ + ...running().messages, + { id: "answer", role: "assistant", text: "Done", turnId: "turn-1", streaming: false, createdAt: at(1), updatedAt: at(1) }, + { id: "queued", role: "user", text: "Next", turnId: null, streaming: false, createdAt: at(5), updatedAt: at(5) }, + ], + })]); + + await expect(sendThreadMessage(harness.config, { threadId: "target", prompt: "And then this" })).rejects.toMatchObject({ + code: "THREAD_BUSY", + details: { turnRunning: false, queuedMessages: 1 }, + }); + }); + it("refuses to send with new settings while a turn runs", async () => { const harness = await testHarness([makeThread("target", running())], { catalog: CATALOG }); diff --git a/tests/cli.test.mjs b/tests/cli.test.mjs index 23c4ade..d5077b5 100644 --- a/tests/cli.test.mjs +++ b/tests/cli.test.mjs @@ -84,6 +84,16 @@ describe("CLI parsing", () => { }); }); + it("rejects an unknown busy-thread mode", async () => { + const result = await run(["--json", "threads", "send", "--thread", "thread-1", "--prompt", "x", "--if-busy", "queue"]); + + expect(result.code).toBe(2); + expect(JSON.parse(result.stderr).error).toEqual({ + code: "INVALID_USAGE", + message: "option '--if-busy ' argument 'queue' is invalid. Allowed choices are reject, inject.", + }); + }); + it("rejects an unknown read detail", async () => { const result = await run(["--json", "threads", "read", "--thread", "thread-1", "--detail", "verbose"]);