OR-129 Let an agent delegate a task to a second agent session (#221)

* Let an agent delegate a task to a second agent session

An agent could only do work itself or queue it behind the current turn.
`orx agent spawn "<task>"` now starts a helper in its own top-level
session — visible in Recents, its own transcript, its own worktree — and
resumes this chat with the helper's closing reply when it finishes.

The CLI only writes rows: the child's `chat_sessions` row (carrying a new
`parent_session_id`) and a `chat_spawns` record. The resident `orx up`
watcher starts the child's first turn and later reports back, the same
store-and-watcher split as `orx exp wake` and for the same reason — the
spawning `orx` is a short-lived subprocess with no harness to run a turn
on. It is a CLI verb rather than an MCP tool because the mcp-gate bridge
is Claude-only, while all three harnesses already shell out to `orx` with
`ORX_CHAT_SESSION_ID` exported.

Spawn rows walk pending → starting → running → notifying → done. The two
claimed states carry a token so a watcher that dies mid-step leaves work
reclaimable rather than lost, and a helper whose parent was deleted
retires quietly instead of stranding the row.

Spawned agents cannot spawn their own. A helper that can delegate turns
one request into an unbounded tree of paid sessions, and nothing
downstream bounds it.

Settings carry over only when parent and child share a harness: a model
or permission-mode id from one CLI means nothing to another. Plan never
carries over — a helper is spawned to do the task, not to plan it.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* Harden spawned agents against losing the delegated task

Review of the spawn feature found several ways a delegated task could be
accepted and then quietly not happen.

A helper's completion was read from `ChatHost::is_busy`, a map this
process owns. A second `orx up` — or the same one after a restart — sees
an empty map, so a helper working under it was reported "finished" within
3s and the delegation was abandoned. The gate now consults the durable
turn lease, and `finish_turn` stamps the spawn row, so an idle helper
with no stamp is reported as interrupted rather than as a false
completion carrying half-written text.

A failed start returned the row to Pending and retried every 3s forever.
The brief is persisted before `send_message_showing` can report the turn
didn't start, so each retry appended it again. Retries now skip the
already-recorded brief, stop after three attempts, and wake the parent
with "could not be started" — the CLI had promised it a wake-up.

Claude activates Plan through its permission mode, not the plan axis, so
clearing `plan_mode` still handed a planning parent's helper a mode that
only ever produces a plan. The plan permission id is now dropped when
inheriting.

The closing reply was read from the last assistant message. An answered
prompt card rides its own text-less message and becomes the branch tip,
so an opencode helper that hit one card reported "no written reply" and
its actual answer was lost. A turn that died on a harness error reported
the same thing, reading as "nothing to say" rather than "it failed". The
report now walks back to the newest message with text and distinguishes
a reply, a failure, and genuine silence.

Depth was capped at one level; breadth was not, so a looping parent could
create as many paid sessions and worktrees as it liked. Five in flight.

Also: the transcript's squash key ignored the spawned ids, collapsing two
adjacent spawns into one row and stranding the second helper; following a
spawn card set the one session filter that hides archived rows; the
delegation playbook told agents to hand helpers a node this session owns,
which cardinal rule 1 and one-branch-one-owner both forbid, and never
told them to bound a helper's compute.

The playbook placeholder test checked a hardcoded token list, so it could
not catch a newly added token — it now scans the rendered prompt, and
fails on the `{max_spawns}` this change introduced if its substitution is
removed.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* Never retire a delegation before its report is delivered

Round-two review found four more ways a spawned task could be lost.

The give-up notification was fire-and-forget: it returned silently when
the parent's turn slot was taken, then retired the row anyway. The parent
is busy in exactly the expected case — it spawned from inside a turn, and
the third failed attempt lands ~9s later — so the promised "could not be
started" usually never arrived. Both wake-ups now go through one
`deliver_wake_up`, which retires the row only once the message actually
started a turn and otherwise leaves it for the next tick.

`--no-wake` settled to Done the instant the turn started, and the
in-flight cap counts only rows that are not Done. A looping parent could
spawn five, wait one tick, and spawn five more forever. Fire-and-forget
rows now stay Running until the helper is genuinely idle, so "in flight"
means the same thing for both kinds.

`interrupted` was read from the listing snapshot while the liveness check
was live. A helper finishing cleanly while an earlier iteration awaited a
whole turn was reported as interrupted and its reply thrown away — the
same loss the stamp was added to prevent, inverted. It is re-read after
the liveness check.

A user pressing Stop reaches the same `finish_turn`, so a manually
stopped helper was stamped finished and its half-written text handed back
as a closing reply. `spawn_outcome` now recognizes the interrupt marker.

The wake-up also claimed the helper's commits were "on its own branch".
Session worktrees start detached, so a helper only has a branch if it ran
`orx create-experiment`, and for the delegations the playbook recommends
there are usually no commits at all. Restored the accurate wording, in
the report and the playbook, and the clause is omitted entirely when the
worktree isn't there.

Also: a reclaimed `starting` claim now counts as an attempt, so a watcher
that dies inside the send cannot retry forever; `ChatSpawn` no longer
carries a `state` no reader used, and the INSERT hardcodes `pending`;
branch-local dev slots get ALTERs for the renamed columns.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
Myles Anderson
2026-08-20 10:48:41 -07:00
committed by GitHub
co-authored by Claude Opus 5
parent 95ea6b8e5a
commit a0c30df38a
17 changed files with 1767 additions and 2308 deletions
+34
View File
@@ -113,6 +113,39 @@ of the same clone. Git state is shared between you:
- Your worktree starts **detached on the baseline tip**; check out your
experiment's branch before editing.
### Delegating with `orx agent spawn`
You can start one of those sessions yourself. `orx agent spawn "<task>"`
creates a **new top-level session** in this project — visible to the user in
the sidebar, on its own worktree, with its own transcript — and hands it the
task. This chat is resumed with the helper's closing reply when it finishes
(pass `--no-wake` if you don't need to hear back). Use `--stdin` for a brief
too long for a command line. You may have {max_spawns} helpers in flight at
once, and a session that was itself spawned cannot spawn helpers of its own —
if `orx agent spawn` tells you that, do the task here instead.
Delegate work that is genuinely **separate from the node you are on**: a
literature sweep, a survey of an unfamiliar codebase, a write-up of results you
already have. Do it yourself when it is a step in the loop you are running —
the round trip costs more than the step.
Two rules make a helper safe to start:
- **Never hand it a branch this session holds.** One branch, one owner (above)
applies across spawned sessions too, and a frozen node may not be edited by
anyone (cardinal rule 1). If the helper needs to change code, tell it to
create its own node with `orx create-experiment {id} --parent <expId>` and
work there.
- **Say which runs it may launch, if any.** The helper reads this same
playbook, so it will otherwise assume the full research loop is open to it
and may launch paid runs you never see. Name them ("one `orx exp run` on
`<expId>`, nothing else") or forbid them ("do not run `orx exp run`").
The helper starts with an **empty transcript and cannot see this conversation**,
so the brief must stand alone: name the project, the experiment id, the branch,
the metric, and what "done" looks like. Its edits stay in its own worktree; the
wake-up names where. Nothing merges into yours by itself.
## Cardinal rules
Breaking any of these silently invalidates results — they are not style
@@ -162,6 +195,7 @@ preferences.
| `orx exp cancel <expId>` | Cancel the in-flight run. |
| `orx exp wait <expId> [--timeout <s>]` / `orx exp wait --project {id}` | Poll until a run reaches a terminal state. Exits **non-zero** after `--timeout` seconds (default 1800) with nothing changed — that means "still running", not an error. |
| `orx exp wake <expId>` | Resume this local chat after the experiment's latest run succeeds or fails. Cancelled runs do not wake the chat. |
| `orx agent spawn "<task>" [--title "<t>"] [--stdin] [--no-wake]` | Hand a self-contained task to a helper agent in its own session and worktree. This chat resumes with the helper's reply when it finishes. |
| `orx runs {id} [--experiment <expId>]` | Run table, newest first. Run ids come from here. |
| `orx logs <runId> [--head] [--bytes <n>] [--range <s>:<e>]` | Read a run's log (tail by default). |
| `orx lit "<query>" [--source alphaxiv\|openalex\|biorxiv]` / `orx paper <id\|url>` | Literature search across alphaXiv, OpenAlex, and bioRxiv (public, no login): **`orx-lit`** skill. Preferred over web search for academic/research queries — start here. |
+231
View File
@@ -0,0 +1,231 @@
//! The `agent` command group: delegate work to a second agent session.
//!
//! orx agent spawn "<task>" start a helper agent on its own top-level session
//!
//! Only meaningful inside a local `orx up` agent session. `ORX_LOCAL_SESSION`
//! marks the process as one; `ORX_CHAT_SESSION_ID` names the session doing the
//! spawning (see `local::chat::set_chat_session_env`). Both are needed — the
//! cloud opencode plugin exports the session id too, for run attribution.
//!
//! This command only writes the child's session row and a `chat_spawns` record;
//! it never runs the child itself. The resident `orx up` picks the record up,
//! starts the helper's first turn, and (unless `--no-wake`) wakes the parent
//! when the helper is done. Same store-and-watcher split as `orx exp wake`, and for
//! the same reason: the CLI is a short-lived subprocess with no harness of its
//! own to run a turn on.
use std::io::Read;
use crate::error::{anyhow, Result};
use crate::local::harness::PermissionMode;
use crate::store::{now_ms, ChatSpawn, Store, StoredChatSession};
use crate::AgentCommand;
/// Helpers one session may have in flight at once.
pub(crate) const MAX_LIVE_SPAWNS: i64 = 5;
pub async fn run(args: crate::AgentArgs) -> Result<()> {
let store = Store::open()?;
match args.command {
AgentCommand::Spawn {
task,
stdin,
title,
harness,
model,
no_wake,
} => spawn(&store, task, stdin, title, harness, model, !no_wake),
}
}
/// Read the task from the positional argument or, with `--stdin`, from the
/// whole of stdin (agents write multi-paragraph briefs as heredocs).
fn task_text(task: Option<String>, stdin: bool) -> Result<String> {
if stdin {
if task.is_some() {
return Err(anyhow!(
"Pass the task as an argument or --stdin, not both."
));
}
let mut buf = String::new();
std::io::stdin()
.read_to_string(&mut buf)
.map_err(|e| anyhow!("Could not read the task from stdin: {e}"))?;
return non_empty(buf);
}
non_empty(task.unwrap_or_default())
}
fn non_empty(text: String) -> Result<String> {
let trimmed = text.trim();
if trimmed.is_empty() {
return Err(anyhow!(
"Describe the task for the spawned agent: `orx agent spawn \"<task>\"`."
));
}
Ok(trimmed.to_string())
}
/// Why this session may not spawn right now, if it may not. Depth and breadth
/// are the two ways one request becomes an unbounded tree of paid sessions, and
/// nothing downstream of here bounds either.
fn spawn_refusal(parent: &StoredChatSession, live: i64) -> Option<String> {
if parent.parent_session_id.is_some() {
return Some(
"This session was itself spawned by another agent, and spawned agents cannot spawn \
their own. Do the task here, or report back so the session that spawned you can \
delegate it."
.to_string(),
);
}
(live >= MAX_LIVE_SPAWNS).then(|| {
format!(
"You already have {live} agents in flight, the most one session may run at once. \
Wait for one to report back before spawning another."
)
})
}
fn spawn(
store: &Store,
task: Option<String>,
stdin: bool,
title: Option<String>,
harness: Option<String>,
model: Option<String>,
wake_parent: bool,
) -> Result<()> {
if !crate::local::chat::in_local_session() {
return Err(anyhow!(
"`orx agent spawn` is only available inside a local `orx up` agent session."
));
}
let parent_id = crate::local::chat::launching_chat_session()
.ok_or_else(|| anyhow!("This agent session has no chat id to spawn from."))?;
let parent = store
.get_chat_session(&parent_id)?
.ok_or_else(|| anyhow!("The current chat session no longer exists."))?;
let prompt = task_text(task, stdin)?;
if let Some(refusal) = spawn_refusal(&parent, store.count_live_chat_spawns(&parent_id)?) {
return Err(anyhow!(refusal));
}
let harness = harness.unwrap_or_else(|| parent.harness.clone());
if !crate::local::harness::is_chat_harness(&harness) {
return Err(anyhow!("unknown harness: {harness}"));
}
// Settings only carry over when the child runs the same harness; a model or
// permission-mode id from one CLI is meaningless to another.
let inherits = harness == parent.harness;
// Claude activates Plan through its permission mode, not the plan axis, so
// clearing `plan_mode` alone would still hand a planning parent's helper a
// mode that only ever produces a plan.
let plan_permission =
crate::local::harness::permission_id_for_mode(&harness, PermissionMode::Plan);
let title = title
.map(|t| t.trim().to_string())
.filter(|t| !t.is_empty());
let session = StoredChatSession {
id: format!("chat_{}", uuid::Uuid::new_v4()),
project_id: parent.project_id.clone(),
harness,
native_session_id: None,
// "user" = explicitly chosen, so auto-titling leaves it alone.
title_source: title.is_some().then(|| "user".to_string()),
title,
model: model.or_else(|| inherits.then(|| parent.model.clone()).flatten()),
permission_mode: inherits
.then(|| parent.permission_mode.clone())
.flatten()
.filter(|mode| Some(mode) != plan_permission.as_ref()),
plan_mode: false,
plan_reset_pending: false,
reasoning_level: inherits.then(|| parent.reasoning_level.clone()).flatten(),
archived: false,
context_usage_json: None,
bootstrap_context: None,
active_leaf_id: None,
parent_session_id: Some(parent_id.clone()),
created_at: now_ms(),
updated_at: now_ms(),
};
// One transaction: a session row without its spawn row is an empty session
// in the user's sidebar that nothing will ever start.
let tx = store.begin()?;
store.create_chat_session(&session)?;
store.create_chat_spawn(&ChatSpawn {
session_id: session.id.clone(),
parent_session_id: parent_id,
prompt,
wake_parent,
attempts: 0,
finished_at: None,
})?;
tx.commit()?;
println!("Spawned agent session {}.", session.id);
println!("It starts within a few seconds and works in its own git worktree.");
if wake_parent {
println!("This chat will be resumed with its result when it finishes.");
} else {
println!("You will NOT be told when it finishes; check its session yourself.");
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::{spawn_refusal, task_text, MAX_LIVE_SPAWNS};
use crate::store::StoredChatSession;
fn parent(parent_session_id: Option<&str>) -> StoredChatSession {
StoredChatSession {
id: "chat_parent".into(),
project_id: "p1".into(),
harness: "codex".into(),
native_session_id: None,
title: None,
title_source: None,
model: None,
permission_mode: None,
plan_mode: false,
plan_reset_pending: false,
reasoning_level: None,
archived: false,
context_usage_json: None,
bootstrap_context: None,
active_leaf_id: None,
parent_session_id: parent_session_id.map(str::to_string),
created_at: 1,
updated_at: 1,
}
}
#[test]
fn a_spawned_session_may_not_spawn_its_own() {
let refusal = spawn_refusal(&parent(Some("chat_grandparent")), 0)
.expect("a spawned session must be refused");
assert!(refusal.contains("cannot spawn"), "{refusal}");
// Depth is refused regardless of how few helpers are in flight.
assert!(spawn_refusal(&parent(None), 0).is_none());
}
#[test]
fn one_session_may_only_run_so_many_helpers_at_once() {
assert!(spawn_refusal(&parent(None), MAX_LIVE_SPAWNS - 1).is_none());
let refusal =
spawn_refusal(&parent(None), MAX_LIVE_SPAWNS).expect("the cap must refuse one more");
assert!(refusal.contains("in flight"), "{refusal}");
}
#[test]
fn a_task_is_required_and_comes_from_one_place() {
assert_eq!(
task_text(Some(" Sweep the literature ".into()), false).unwrap(),
"Sweep the literature"
);
assert!(task_text(None, false).is_err());
assert!(task_text(Some(" ".into()), false).is_err());
// --stdin and a positional together are ambiguous, so neither is used.
assert!(task_text(Some("from the args".into()), true).is_err());
}
}
+1
View File
@@ -15,6 +15,7 @@
//! return `Err(anyhow!(...))` (clap already enforces required positionals, so
//! most of those usage guards are unnecessary in the Rust port).
pub mod agent;
pub mod app;
pub mod compute;
pub mod create_experiment;
+1
View File
@@ -4575,6 +4575,7 @@ async fn create_chat_session(
context_usage_json: None,
bootstrap_context: None,
active_leaf_id: None,
parent_session_id: None,
created_at: now_ms(),
updated_at: now_ms(),
};
+745 -1
View File
@@ -23,7 +23,7 @@ use crate::error::{anyhow, Result};
use crate::local::harness::ResumeAction;
use crate::local::model::LocalProject;
use crate::local::opencode::AgentHost;
use crate::store::{now_ms, Store, StoredChatMessage, StoredChatSession};
use crate::store::{now_ms, ChatSpawnState, Store, StoredChatMessage, StoredChatSession};
/// Min interval between mid-turn persist+broadcast flushes (streaming parts
/// can update many times a second; the final flush is always unconditional).
@@ -1100,6 +1100,7 @@ pub fn session_json(s: &StoredChatSession, busy: bool) -> Value {
"busy": busy,
"contextUsage": context_usage,
"activeLeafId": s.active_leaf_id,
"parentSessionId": s.parent_session_id,
})
}
@@ -2734,6 +2735,36 @@ impl ChatHost {
.await
}
/// Deliver a spawned helper's brief as the opening message of its session.
/// Visible, unlike a wake-up: the brief *is* that transcript's starting
/// point, and hiding it would leave the helper apparently working unbidden.
/// `record_brief` is false on a retry, where it is already in the
/// transcript — this persists the bubble before it can report a failure.
async fn send_spawn_task(
self: &Arc<Self>,
session_id: &str,
text: String,
record_brief: bool,
guard: TurnGuard,
) -> Result<TurnSubmission> {
self.send_message_showing(
session_id,
vec![AnnotatedText {
text: text.clone(),
annotations: Vec::new(),
}],
TranscriptDisplay {
text: Some(text),
annotations: None,
record_user_message: record_brief,
},
TurnOverrides::default(),
TurnAttachments::Uploaded(Vec::new()),
TurnAdmission::Preclaimed(guard),
)
.await
}
/// Persists one displayed user turn, then expands and contextualizes each
/// raw annotated message separately for the harness. `transcript_text`
/// overrides the displayed text; an empty override records no user message.
@@ -3208,6 +3239,9 @@ impl ChatHost {
self.cancel_pending_permissions(session_id);
let session = if let Ok(store) = Store::open() {
let _ = store.touch_chat_session(session_id);
// Without this a crash mid-turn is indistinguishable from a
// completed one once the turn lease lapses.
let _ = store.mark_chat_spawn_finished(session_id);
store.get_chat_session(session_id).ok().flatten()
} else {
None
@@ -4745,6 +4779,330 @@ async fn process_run_wakeups(
Ok(())
}
// --- spawned agents -------------------------------------------------------------
/// How much of a helper's closing reply, and of the brief echoed back with it,
/// rides into the parent's context. Both are agent-authored and unbounded; the
/// user can open the helper's session for the rest.
const SPAWN_REPORT_LIMIT: usize = 4000;
const SPAWN_BRIEF_LIMIT: usize = 500;
/// Give up starting a helper after this many failed attempts, so a permanently
/// failing spawn reports once instead of retrying every tick forever.
const SPAWN_START_ATTEMPTS: i64 = 3;
fn truncated(text: &str, limit: usize) -> String {
if text.chars().count() <= limit {
return text.to_string();
}
text.chars().take(limit).collect::<String>() + "… (truncated)"
}
/// What the helper left behind on the branch on screen.
enum SpawnOutcome {
Reply(String),
Failed(String),
/// The turn was stopped — by the user, or by an `orx up` that died holding
/// it. Whatever text is there is half-written, not an answer.
Interrupted,
Silent,
}
/// Read the helper's closing words. Walks back from the tip rather than reading
/// only the last assistant row: an answered prompt card rides its own
/// text-less assistant message and would otherwise mask the real reply.
fn spawn_outcome(store: &Store, session: &StoredChatSession) -> Result<SpawnOutcome> {
let messages = store.list_chat_messages(&session.id)?;
let path = active_path(&messages, session.active_leaf_id.as_deref());
let mut failure = None;
// Stop at the newest user message: anything above it belongs to an earlier
// turn and is not this task's answer.
for message in path.iter().rev().take_while(|m| m.role != "user") {
let parts: Vec<WirePart> = serde_json::from_str(&message.parts_json)?;
// A user Stop persists this marker and no text; the half-written text
// above it must not be handed back as a closing reply.
if parts
.iter()
.any(|part| part.tool.as_deref() == Some("interrupted"))
{
return Ok(SpawnOutcome::Interrupted);
}
let text = parts
.iter()
.filter(|part| part.kind == "text")
.filter_map(|part| part.text.as_deref())
.map(str::trim)
.filter(|text| !text.is_empty())
.collect::<Vec<_>>()
.join("\n\n");
if !text.is_empty() {
return Ok(SpawnOutcome::Reply(truncated(&text, SPAWN_REPORT_LIMIT)));
}
// Keep the newest error, but keep looking: a harness that errors after
// answering should still hand the answer back.
if failure.is_none() {
failure = parts
.iter()
.find(|part| part.tool.as_deref() == Some("error"))
.and_then(|part| part.state.as_ref())
.and_then(|state| state.error.clone().or_else(|| state.output.clone()));
}
}
Ok(match failure {
Some(error) => SpawnOutcome::Failed(truncated(error.trim(), SPAWN_REPORT_LIMIT)),
None => SpawnOutcome::Silent,
})
}
/// Where the helper's work is. Deliberately claims no branch: session worktrees
/// start detached, so a helper only has one if it ran `orx create-experiment`.
fn spawn_workspace(store: &Store, session: &StoredChatSession) -> String {
let Ok(Some(project)) = store.get_local_project(&session.project_id) else {
return String::new();
};
let path = crate::local::git::existing_session_worktree_path(&project, &session.id);
if !path.exists() {
return String::new();
}
format!(
"\n\nIts edits are in its own worktree at `{}` — read them there. Nothing was merged \
into yours.",
path.display()
)
}
fn spawn_report_text(
store: &Store,
spawn: &crate::store::ChatSpawn,
interrupted: bool,
) -> Result<String> {
let child = store
.get_chat_session(&spawn.session_id)?
.ok_or_else(|| anyhow!("spawned session disappeared"))?;
let named = child
.title
.as_deref()
.map(|title| format!(" (\"{title}\")"))
.unwrap_or_default();
let brief = truncated(spawn.prompt.trim(), SPAWN_BRIEF_LIMIT);
let stopped = "Its session holds however far it got. Re-delegate it if you still need the \
task done.";
let (headline, closing) = if interrupted {
("was interrupted before it finished", stopped.to_string())
} else {
match spawn_outcome(store, &child)? {
SpawnOutcome::Reply(reply) => {
("has finished", format!("Its closing reply:\n\n{reply}"))
}
SpawnOutcome::Failed(error) => (
"failed",
format!("It ended on an error and did NOT do the task:\n\n{error}"),
),
SpawnOutcome::Interrupted => ("was stopped before it finished", stopped.to_string()),
SpawnOutcome::Silent => (
"has finished",
"It wrote no reply, so treat the task as unconfirmed and check its session."
.to_string(),
),
}
};
Ok(format!(
"[orx] The agent you spawned for `{}`{named} {headline}.\n\nIt was asked to: {brief}\n\n\
{closing}{}",
spawn.session_id,
spawn_workspace(store, &child),
))
}
fn spawn_start_failure_text(spawn: &crate::store::ChatSpawn) -> String {
format!(
"[orx] The agent you spawned for `{}` could NOT be started, so the task was not done. \
It was asked to: {}\n\nDo the task here, or delegate it again.",
spawn.session_id,
truncated(spawn.prompt.trim(), SPAWN_BRIEF_LIMIT),
)
}
/// Start the first turn of every freshly spawned helper, then wake the parent
/// that asked for each finished one.
///
/// Both halves are claim-guarded in the store rather than in memory: the CLI
/// that wrote the row is long gone, and a crash between claiming and delivering
/// has to be recoverable by the next `orx up`.
async fn process_chat_spawns(
chat: &Arc<ChatHost>,
mut store: Store,
data_dir_move_in_progress: Option<&std::sync::atomic::AtomicBool>,
) -> Result<()> {
let moving = || {
data_dir_move_in_progress.is_some_and(|flag| flag.load(std::sync::atomic::Ordering::SeqCst))
};
if moving() {
return Ok(());
}
store.prune_chat_spawns()?;
for spawn in store.list_chat_spawns(ChatSpawnState::Pending)? {
if moving() {
return Ok(());
}
// Out of retries: the helper will never run, and all that is left is
// telling the parent so. Retried like any other wake-up rather than
// fired once — the parent is usually still mid-turn from delegating.
if spawn.attempts >= SPAWN_START_ATTEMPTS {
if !spawn.wake_parent {
let Some(token) = store.claim_chat_spawn(
&spawn.session_id,
ChatSpawnState::Pending,
ChatSpawnState::Waking,
)?
else {
continue;
};
store.settle_chat_spawn(&spawn.session_id, &token, ChatSpawnState::Done)?;
continue;
}
store = deliver_wake_up(
chat,
store,
&spawn,
ChatSpawnState::Pending,
spawn_start_failure_text(&spawn),
)
.await?;
continue;
}
let Some(mut guard) = TurnGuard::claim_hidden(chat, &spawn.session_id).await else {
continue;
};
let Some(token) = store.claim_chat_spawn(
&spawn.session_id,
ChatSpawnState::Pending,
ChatSpawnState::Starting,
)?
else {
guard.release().await;
continue;
};
// Broadcast before the turn so the session appears in every open
// dashboard's Recents as it starts working, not after its first flush.
chat.emit_session(store.get_chat_session(&spawn.session_id)?)
.await;
// A retry must not re-record the brief: `send_message_showing` persists
// the user message before it can report the turn didn't start.
let record_brief = store.list_chat_messages(&spawn.session_id)?.is_empty();
let started = chat
.send_spawn_task(&spawn.session_id, spawn.prompt.clone(), record_brief, guard)
.await;
// Running even for a fire-and-forget spawn: the second loop retires it
// once the helper is actually idle, so `--no-wake` cannot slip past the
// in-flight cap by retiring the instant its turn starts.
let next = match started {
Ok(TurnSubmission::Started) => ChatSpawnState::Running,
outcome => {
if let Err(err) = &outcome {
eprintln!(
"orx up: could not start spawned agent {}: {err}",
spawn.session_id
);
}
store.record_chat_spawn_attempt(&spawn.session_id)?;
ChatSpawnState::Pending
}
};
store.settle_chat_spawn(&spawn.session_id, &token, next)?;
}
for spawn in store.list_chat_spawns(ChatSpawnState::Running)? {
if moving() {
return Ok(());
}
// The durable lease, not just this process's turn map: a helper running
// under another `orx up` (or under one that has since restarted) is
// still working, and reporting it finished would abandon the task.
if chat.is_busy(&spawn.session_id).await || store.chat_turn_leased(&spawn.session_id)? {
continue;
}
if !spawn.wake_parent || store.get_chat_session(&spawn.parent_session_id)?.is_none() {
// Nobody to tell — fire-and-forget, or a deleted parent. The
// helper's own session stays either way.
let Some(token) = store.claim_chat_spawn(
&spawn.session_id,
ChatSpawnState::Running,
ChatSpawnState::Waking,
)?
else {
continue;
};
store.settle_chat_spawn(&spawn.session_id, &token, ChatSpawnState::Done)?;
continue;
}
// Re-read rather than trusting the listing: earlier iterations await
// whole turns, and a helper that finished cleanly during that window
// would otherwise be reported as interrupted and its reply discarded.
let interrupted = store
.get_chat_spawn(&spawn.session_id)?
.is_some_and(|row| row.finished_at.is_none());
let text = match spawn_report_text(&store, &spawn, interrupted) {
Ok(text) => text,
Err(err) => {
eprintln!("orx up: could not summarize spawned agent: {err}");
let Some(token) = store.claim_chat_spawn(
&spawn.session_id,
ChatSpawnState::Running,
ChatSpawnState::Waking,
)?
else {
continue;
};
store.settle_chat_spawn(&spawn.session_id, &token, ChatSpawnState::Done)?;
continue;
}
};
store = deliver_wake_up(chat, store, &spawn, ChatSpawnState::Running, text).await?;
}
Ok(())
}
/// Wake the parent and retire the row, but only together: a row that stays in
/// `from` is retried next tick, which is what keeps a busy parent from costing
/// the delegation its report.
/// Takes the `Store` by value and hands it back: it is `!Sync`, so a borrow
/// held across these awaits would make the watcher future non-`Send`.
async fn deliver_wake_up(
chat: &Arc<ChatHost>,
store: Store,
spawn: &crate::store::ChatSpawn,
from: ChatSpawnState,
text: String,
) -> Result<Store> {
let Some(guard) = TurnGuard::claim_hidden(chat, &spawn.parent_session_id).await else {
return Ok(store);
};
let Some(token) = store.claim_chat_spawn(&spawn.session_id, from, ChatSpawnState::Waking)?
else {
return Ok(store);
};
let next = match chat
.send_hidden_message(&spawn.parent_session_id, text, guard)
.await
{
Ok(TurnSubmission::Started) => ChatSpawnState::Done,
outcome => {
if let Err(err) = &outcome {
if !chat.is_busy(&spawn.parent_session_id).await {
eprintln!("orx up: could not wake a spawned agent's parent: {err}");
}
}
from
}
};
if !store.settle_chat_spawn(&spawn.session_id, &token, next)? {
return Err(anyhow!(
"spawn claim expired before the parent wake-up was recorded"
));
}
Ok(store)
}
/// Resume explicitly subscribed agent sessions after a run finishes. Busy and
/// draining sessions retain their durable wake-up until they become idle.
pub async fn watch_runs(
@@ -4770,6 +5128,12 @@ pub async fn watch_runs(
{
eprintln!("orx up: run watcher: {err}");
}
let Ok(store) = Store::open() else { continue };
if let Err(err) =
process_chat_spawns(&chat, store, Some(data_dir_move_in_progress.as_ref())).await
{
eprintln!("orx up: spawn watcher: {err}");
}
}
}
@@ -5747,6 +6111,7 @@ mod bridge_tests {
context_usage_json: None,
bootstrap_context: None,
active_leaf_id: None,
parent_session_id: None,
created_at: 1,
updated_at: 1,
}
@@ -5880,6 +6245,7 @@ mod run_wakeup_tests {
context_usage_json: None,
bootstrap_context: None,
active_leaf_id: None,
parent_session_id: None,
created_at: 1,
updated_at: 1,
})
@@ -5929,6 +6295,384 @@ with other project runs using `orx runs p1` and inspect this run's logs using `o
assert!(run_wakeup_text(&run("cancelled")).is_none());
}
fn assistant_message(id: &str, parent: Option<&str>, text: &str) -> StoredChatMessage {
StoredChatMessage {
id: id.into(),
session_id: "child".into(),
role: "assistant".into(),
parts_json: serde_json::to_string(&[WirePart::text(format!("{id}-part"), text)])
.unwrap(),
created_at: 1,
parent_id: parent.map(str::to_string),
base_native_session_id: None,
result_native_session_id: None,
}
}
/// `create_chat_spawn` always inserts `pending`, so a test that needs a
/// later state advances the row the way the watcher would.
fn spawn_row(store: &Store, child: &str, parent: &str, state: ChatSpawnState) {
store
.create_chat_spawn(&crate::store::ChatSpawn {
session_id: child.into(),
parent_session_id: parent.into(),
prompt: "Sweep the literature".into(),
wake_parent: true,
attempts: 0,
finished_at: None,
})
.unwrap();
if matches!(state, ChatSpawnState::Running) {
let token = store
.claim_chat_spawn(child, ChatSpawnState::Pending, ChatSpawnState::Starting)
.unwrap()
.unwrap();
store
.settle_chat_spawn(child, &token, ChatSpawnState::Running)
.unwrap();
// A Running row stands for a turn that reached `finish_turn`; the
// crash and Stop cases are exercised in the report-text tests.
store.mark_chat_spawn_finished(child).unwrap();
}
}
fn spawn_fixture(child: &str, parent: &str) -> crate::store::ChatSpawn {
crate::store::ChatSpawn {
session_id: child.into(),
parent_session_id: parent.into(),
prompt: "Sweep the literature".into(),
wake_parent: true,
attempts: 0,
finished_at: Some(1),
}
}
fn bare_host() -> Arc<ChatHost> {
Arc::new(ChatHost::new(
Arc::new(crate::local::opencode::AgentHost::new(None)),
Arc::new(crate::local::codex::CodexHost::new()),
Arc::new(crate::local::claude::ClaudeHost::new()),
))
}
#[tokio::test]
async fn a_helper_whose_turn_is_already_held_keeps_its_row_pending() {
let (store, dir) = temp_store("spawn-busy");
session(&store, "child");
spawn_row(&store, "child", "parent", ChatSpawnState::Pending);
let host = bare_host();
host.turns
.lock()
.await
.insert("child".into(), TurnState::Draining);
drop(store);
process_chat_spawns(&host, Store::open_at(dir.clone()).unwrap(), None)
.await
.unwrap();
let store = Store::open_at(dir.clone()).unwrap();
assert_eq!(
store
.list_chat_spawns(ChatSpawnState::Pending)
.unwrap()
.len(),
1
);
drop(store);
std::fs::remove_dir_all(dir).unwrap();
}
#[tokio::test]
async fn a_working_helper_is_not_reported_until_it_goes_idle() {
let (store, dir) = temp_store("spawn-working");
session(&store, "child");
session(&store, "parent");
spawn_row(&store, "child", "parent", ChatSpawnState::Running);
let host = bare_host();
host.turns
.lock()
.await
.insert("child".into(), TurnState::Reserved);
drop(store);
process_chat_spawns(&host, Store::open_at(dir.clone()).unwrap(), None)
.await
.unwrap();
let store = Store::open_at(dir.clone()).unwrap();
assert_eq!(
store
.list_chat_spawns(ChatSpawnState::Running)
.unwrap()
.len(),
1
);
assert!(!host.is_busy("parent").await);
drop(store);
std::fs::remove_dir_all(dir).unwrap();
}
#[tokio::test]
async fn a_helper_still_running_under_another_orx_up_is_not_reported_finished() {
let (store, dir) = temp_store("spawn-leased");
session(&store, "child");
session(&store, "parent");
spawn_row(&store, "child", "parent", ChatSpawnState::Running);
// Another process holds the turn: this one's `turns` map knows nothing
// about it, so only the durable lease can stop a false completion.
store
.claim_chat_turn("child", "other-process-token")
.unwrap();
let host = bare_host();
drop(store);
process_chat_spawns(&host, Store::open_at(dir.clone()).unwrap(), None)
.await
.unwrap();
let store = Store::open_at(dir.clone()).unwrap();
assert_eq!(
store
.list_chat_spawns(ChatSpawnState::Running)
.unwrap()
.len(),
1,
"a leased helper must stay Running"
);
assert!(!host.is_busy("parent").await);
drop(store);
let _ = std::fs::remove_dir_all(dir);
}
#[test]
fn an_interrupted_helper_is_not_reported_as_finished() {
let (store, dir) = temp_store("spawn-interrupted");
session(&store, "child");
store
.upsert_chat_message(&assistant_message("a1", None, "halfway through"))
.unwrap();
let spawn = crate::store::ChatSpawn {
session_id: "child".into(),
parent_session_id: "parent".into(),
prompt: "Sweep the literature".into(),
wake_parent: true,
attempts: 0,
finished_at: None,
};
let text = spawn_report_text(&store, &spawn, true).unwrap();
assert!(
text.contains("was interrupted before it finished"),
"{text}"
);
assert!(!text.contains("has finished"), "{text}");
// The partial text must not be handed back as a closing answer.
assert!(!text.contains("halfway through"), "{text}");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_helper_that_errored_is_reported_as_failed_not_silent() {
let (store, dir) = temp_store("spawn-errored");
session(&store, "child");
// Exactly what `TurnCtx::push_error` writes.
let parts = vec![WirePart::tool(
"err-0",
"error",
"error",
Some("model 'gpt-5' is not supported".into()),
)];
store
.upsert_chat_message(&StoredChatMessage {
id: "a1".into(),
session_id: "child".into(),
role: "assistant".into(),
parts_json: serde_json::to_string(&parts).unwrap(),
created_at: 1,
parent_id: None,
base_native_session_id: None,
result_native_session_id: None,
})
.unwrap();
let child = store.get_chat_session("child").unwrap().unwrap();
let SpawnOutcome::Failed(error) = spawn_outcome(&store, &child).unwrap() else {
panic!("an error part must not read as a silent turn");
};
assert!(error.contains("not supported"), "{error}");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn an_answered_prompt_card_does_not_mask_the_real_reply() {
let (store, dir) = temp_store("spawn-card");
session(&store, "child");
store
.upsert_chat_message(&assistant_message("a1", None, "the answer"))
.unwrap();
// A resolved card rides its own text-less assistant message and becomes
// the branch tip; reading only the tip would lose the answer above it.
store
.upsert_chat_message(&StoredChatMessage {
id: "a2".into(),
session_id: "child".into(),
role: "assistant".into(),
parts_json: "[]".into(),
created_at: 2,
parent_id: Some("a1".into()),
base_native_session_id: None,
result_native_session_id: None,
})
.unwrap();
store
.set_chat_session_active_leaf("child", Some("a2"))
.unwrap();
let child = store.get_chat_session("child").unwrap().unwrap();
let SpawnOutcome::Reply(reply) = spawn_outcome(&store, &child).unwrap() else {
panic!("expected the reply above the card");
};
assert_eq!(reply, "the answer");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_long_reply_and_a_long_brief_are_both_bounded() {
let (store, dir) = temp_store("spawn-truncate");
session(&store, "child");
store
.upsert_chat_message(&assistant_message("a1", None, &"x".repeat(9000)))
.unwrap();
let spawn = crate::store::ChatSpawn {
session_id: "child".into(),
parent_session_id: "parent".into(),
prompt: "y".repeat(9000),
wake_parent: true,
attempts: 0,
finished_at: Some(1),
};
let text = spawn_report_text(&store, &spawn, false).unwrap();
assert_eq!(text.matches("… (truncated)").count(), 2, "{text}");
assert!(text.chars().count() < SPAWN_REPORT_LIMIT + SPAWN_BRIEF_LIMIT + 600);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_finished_helper_whose_parent_is_gone_retires_quietly() {
let (store, dir) = temp_store("spawn-orphan");
session(&store, "child");
spawn_row(&store, "child", "deleted_parent", ChatSpawnState::Running);
let host = bare_host();
drop(store);
process_chat_spawns(&host, Store::open_at(dir.clone()).unwrap(), None)
.await
.unwrap();
let store = Store::open_at(dir.clone()).unwrap();
assert!(store
.list_chat_spawns(ChatSpawnState::Running)
.unwrap()
.is_empty());
assert_eq!(
store.list_chat_spawns(ChatSpawnState::Done).unwrap().len(),
1
);
drop(store);
std::fs::remove_dir_all(dir).unwrap();
}
#[test]
fn a_spawned_agent_reports_its_last_reply_on_the_branch_on_screen() {
let (store, dir) = temp_store("spawn-report");
session(&store, "child");
store
.upsert_chat_message(&assistant_message("a1", None, "first pass"))
.unwrap();
store
.upsert_chat_message(&assistant_message("a2", Some("a1"), "the answer"))
.unwrap();
// A re-sampled sibling that is *not* the branch on screen must not win.
store
.upsert_chat_message(&assistant_message("a3", Some("a1"), "discarded fork"))
.unwrap();
store
.set_chat_session_active_leaf("child", Some("a2"))
.unwrap();
let child = store.get_chat_session("child").unwrap().unwrap();
let SpawnOutcome::Reply(reply) = spawn_outcome(&store, &child).unwrap() else {
panic!("expected a written reply");
};
assert_eq!(reply, "the answer");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn the_parent_is_told_what_it_asked_for_and_what_came_back() {
let (store, dir) = temp_store("spawn-text");
session(&store, "child");
store
.set_chat_session_title("child", "Lit sweep", "user")
.unwrap();
store
.upsert_chat_message(&assistant_message("a1", None, "Rank 8 wins."))
.unwrap();
let spawn = spawn_fixture("child", "parent");
let text = spawn_report_text(&store, &spawn, false).unwrap();
// No workspace clause: the fixture's project has no checkout on disk,
// which is also what a helper that never wrote anything looks like.
assert_eq!(
text,
"[orx] The agent you spawned for `child` (\"Lit sweep\") has finished.\n\n\
It was asked to: Sweep the literature\n\nIts closing reply:\n\nRank 8 wins."
);
// An untitled helper still reads as a sentence.
session(&store, "untitled");
store
.upsert_chat_message(&StoredChatMessage {
id: "b1".into(),
session_id: "untitled".into(),
role: "assistant".into(),
parts_json: serde_json::to_string(&[WirePart::text("p", "Rank 8 wins.")]).unwrap(),
created_at: 1,
parent_id: None,
base_native_session_id: None,
result_native_session_id: None,
})
.unwrap();
assert!(
spawn_report_text(&store, &spawn_fixture("untitled", "parent"), false)
.unwrap()
.starts_with("[orx] The agent you spawned for `untitled` has finished.")
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_helper_that_wrote_nothing_reads_as_silent_not_as_a_reply() {
let (store, dir) = temp_store("spawn-silent");
session(&store, "child");
let child = store.get_chat_session("child").unwrap().unwrap();
assert!(matches!(
spawn_outcome(&store, &child).unwrap(),
SpawnOutcome::Silent
));
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn hidden_transcript_override_creates_no_user_parts() {
assert!(transcript_parts("", &[], &[]).is_empty());
+3
View File
@@ -330,6 +330,7 @@ fn seed_at(
context_usage_json: None,
bootstrap_context: Some(BOOTSTRAP_CONTEXT.into()),
active_leaf_id: Some(ASSISTANT_MESSAGE_ID.into()),
parent_session_id: None,
created_at: 1_785_824_322_614,
updated_at: 1_785_879_263_859,
};
@@ -369,6 +370,7 @@ fn seed_at(
context_usage_json: None,
bootstrap_context: Some(FIGURE_BOOTSTRAP_CONTEXT.into()),
active_leaf_id: Some(FIGURE_ASSISTANT_MESSAGE_ID.into()),
parent_session_id: None,
created_at: 1_785_824_322_630,
updated_at: 1_785_879_263_858,
};
@@ -411,6 +413,7 @@ fn seed_at(
context_usage_json: None,
bootstrap_context: Some(LITERATURE_BOOTSTRAP_CONTEXT.into()),
active_leaf_id: Some(LITERATURE_ASSISTANT_MESSAGE_ID.into()),
parent_session_id: None,
created_at: 1_785_824_322_634,
updated_at: 1_785_879_263_857,
};
+20 -20
View File
@@ -276,6 +276,10 @@ fn playbook_md(project: &LocalProject) -> String {
.replace("{run_guidance}", run_guidance)
.replace("{compute_guidance}", compute_guidance)
.replace("{skills_scope}", skills_scope)
.replace(
"{max_spawns}",
&crate::commands::agent::MAX_LIVE_SPAWNS.to_string(),
)
}
/// Keep the files we drop into the checkout out of `git status` / accidental
@@ -688,26 +692,22 @@ mod tests {
#[test]
fn playbook_has_no_unresolved_placeholders() {
let md = sample_playbook();
// Every current token must be substituted — a typo'd or newly added
// token that playbook_md doesn't know about fails here.
for token in [
"{name}",
"{id}",
"{repo}",
"{baseline}",
"{paper_line}",
"{compute_bullet}",
"{artifacts}",
"{skills_list}",
"{launch_step}",
"{backends_intro}",
"{run_invocation}",
"{run_guidance}",
"{compute_guidance}",
"{skills_scope}",
] {
assert!(!md.contains(token), "unresolved placeholder {token}");
}
// Scanned, not listed: a NEWLY ADDED token playbook_md doesn't know
// about is exactly the case a hardcoded list cannot catch, and the
// agent would read the literal `{token}` as instruction.
let leftover: Vec<&str> = md
.lines()
.flat_map(|line| {
line.match_indices('{').filter_map(move |(i, _)| {
let rest = &line[i + 1..];
let end = rest.find('}')?;
let token = &rest[..end];
(!token.is_empty() && token.chars().all(|c| c.is_ascii_lowercase() || c == '_'))
.then_some(&line[i..=i + end + 1])
})
})
.collect();
assert!(leftover.is_empty(), "unresolved placeholders: {leftover:?}");
for retired in ["{files}", "{memory}"] {
assert!(!md.contains(retired), "retired placeholder {retired}");
}
+36
View File
@@ -73,6 +73,9 @@ enum Command {
/// Operate on one local project.
Project(ProjectArgs),
/// Delegate a task to a second agent session.
Agent(AgentArgs),
/// List a project's runs.
Runs(RunsArgs),
@@ -370,6 +373,37 @@ pub struct InstanceDeleteArgs {
pub sandbox_id: String,
}
#[derive(Args, Debug)]
pub struct AgentArgs {
#[command(subcommand)]
pub command: AgentCommand,
}
#[derive(Subcommand, Debug)]
pub enum AgentCommand {
/// Hand a task to a helper agent running in its own top-level session.
Spawn {
/// What the helper agent should do. Write it as a self-contained brief:
/// the helper starts with an empty transcript and cannot see this chat.
task: Option<String>,
/// Read the task from stdin instead, for long multi-paragraph briefs.
#[arg(long)]
stdin: bool,
/// Name the session in the sidebar. Defaults to an auto-generated title.
#[arg(long)]
title: Option<String>,
/// Harness for the helper (defaults to this session's).
#[arg(long)]
harness: Option<String>,
/// Model for the helper (defaults to this session's).
#[arg(long)]
model: Option<String>,
/// Do not resume this chat when the helper finishes.
#[arg(long)]
no_wake: bool,
},
}
#[derive(Args, Debug)]
pub struct ExpArgs {
#[command(subcommand)]
@@ -785,6 +819,7 @@ fn command_name(command: &Command) -> &'static str {
Command::Projects(_) => "projects",
Command::Orgs(_) => "orgs",
Command::Project(_) => "project",
Command::Agent(_) => "agent",
Command::Runs(_) => "runs",
Command::Logs(_) => "logs",
Command::CreateExperiment(_) => "create-experiment",
@@ -825,6 +860,7 @@ async fn dispatch(command: Command) -> error::Result<()> {
Command::Projects(args) => commands::projects::run(args).await,
Command::Orgs(args) => commands::orgs::run(args).await,
Command::Project(args) => commands::project::run(args).await,
Command::Agent(args) => commands::agent::run(args).await,
Command::Runs(args) => commands::runs::run(args).await,
Command::Logs(args) => commands::logs::run(args).await,
Command::CreateExperiment(args) => commands::create_experiment::run(args).await,
+380 -5
View File
@@ -235,7 +235,55 @@ pub enum RunWakeupRegistration {
AlreadyDelivered,
}
/// A helper session an agent created with `orx agent spawn`, and the task it
/// was spawned to do. The CLI only writes the row; the resident `orx up`
/// watcher is what starts the child's first turn and reports back.
#[derive(Debug, Clone)]
pub struct ChatSpawn {
pub session_id: String,
pub parent_session_id: String,
pub prompt: String,
/// Whether the parent gets a wake-up once the helper finishes.
pub wake_parent: bool,
/// Failed attempts to start the helper's first turn.
pub attempts: i64,
/// Set by `finish_turn` when the helper's turn ends. Absent on a row whose
/// turn is still running — or whose `orx up` died mid-turn, which is how a
/// crash is told from a completion.
pub finished_at: Option<i64>,
}
/// Lifecycle of a [`ChatSpawn`]: `Pending → Starting → Running → Waking →
/// Done`. The two `-ing` states are claims — one watcher owns the transition
/// and holds a token, so a crash mid-step is reclaimable rather than lost.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ChatSpawnState {
/// Row written; the child's first turn has not been delivered yet.
Pending,
/// Claimed for delivery of the brief to the helper.
Starting,
/// The helper is working on its task.
Running,
/// Claimed for delivery of the wake-up to the parent.
Waking,
/// Terminal: the parent has been told, or never asked to be.
Done,
}
impl ChatSpawnState {
pub fn as_str(self) -> &'static str {
match self {
Self::Pending => "pending",
Self::Starting => "starting",
Self::Running => "running",
Self::Waking => "waking",
Self::Done => "done",
}
}
}
const RUN_WAKEUP_CLAIM_TTL_MS: i64 = 60 * 1000;
const CHAT_SPAWN_CLAIM_TTL_MS: i64 = 60 * 1000;
const CHAT_TURN_LEASE_TTL_MS: i64 = 60 * 1000;
pub struct Store {
@@ -323,6 +371,7 @@ impl Store {
archived INTEGER NOT NULL DEFAULT 0,
context_usage_json TEXT,
bootstrap_context TEXT,
parent_session_id TEXT,
created_at INTEGER NOT NULL,
updated_at INTEGER NOT NULL
);
@@ -350,6 +399,20 @@ impl Store {
);
CREATE INDEX IF NOT EXISTS idx_chat_run_wakeups_requested
ON chat_run_wakeups(requested_at);
CREATE TABLE IF NOT EXISTS chat_spawns (
session_id TEXT PRIMARY KEY,
parent_session_id TEXT NOT NULL,
prompt TEXT NOT NULL,
wake_parent INTEGER NOT NULL DEFAULT 1,
state TEXT NOT NULL DEFAULT 'pending',
requested_at INTEGER NOT NULL,
claim_token TEXT,
claimed_at INTEGER,
attempts INTEGER NOT NULL DEFAULT 0,
finished_at INTEGER
);
CREATE INDEX IF NOT EXISTS idx_chat_spawns_state
ON chat_spawns(state, requested_at);
CREATE TABLE IF NOT EXISTS chat_turn_leases (
chat_session_id TEXT PRIMARY KEY,
claim_token TEXT NOT NULL,
@@ -405,6 +468,10 @@ impl Store {
"ALTER TABLE chat_messages ADD COLUMN base_native_session_id TEXT",
"ALTER TABLE chat_messages ADD COLUMN result_native_session_id TEXT",
"ALTER TABLE chat_sessions ADD COLUMN active_leaf_id TEXT",
"ALTER TABLE chat_sessions ADD COLUMN parent_session_id TEXT",
"ALTER TABLE chat_spawns ADD COLUMN wake_parent INTEGER NOT NULL DEFAULT 1",
"ALTER TABLE chat_spawns ADD COLUMN attempts INTEGER NOT NULL DEFAULT 0",
"ALTER TABLE chat_spawns ADD COLUMN finished_at INTEGER",
] {
let _ = conn.execute(ddl, []);
}
@@ -935,6 +1002,150 @@ impl Store {
Ok(delivered == 1)
}
pub fn create_chat_spawn(&self, spawn: &ChatSpawn) -> Result<()> {
self.conn.execute(
"INSERT INTO chat_spawns
(session_id, parent_session_id, prompt, wake_parent, state, requested_at)
VALUES (?1, ?2, ?3, ?4, 'pending', ?5)",
params![
spawn.session_id,
spawn.parent_session_id,
spawn.prompt,
spawn.wake_parent,
now_ms(),
],
)?;
Ok(())
}
/// Return abandoned claims to their prior state and drop rows whose child
/// session is gone. A missing *parent* is not pruned here: the child keeps
/// working, and the wake step retires the row on its own.
pub fn prune_chat_spawns(&self) -> Result<()> {
self.clear_stale_data_dir_move_lease()?;
let stale = now_ms() - CHAT_SPAWN_CLAIM_TTL_MS;
self.conn.execute(
"UPDATE chat_spawns
SET state = CASE state WHEN 'starting' THEN 'pending' ELSE 'running' END,
attempts = attempts + CASE state WHEN 'starting' THEN 1 ELSE 0 END,
claim_token = NULL, claimed_at = NULL
WHERE state IN ('starting', 'waking') AND claimed_at < ?1
AND NOT EXISTS (SELECT 1 FROM data_dir_move_lease WHERE id = 1)",
params![stale],
)?;
self.conn.execute(
"DELETE FROM chat_spawns
WHERE NOT EXISTS (SELECT 1 FROM data_dir_move_lease WHERE id = 1)
AND NOT EXISTS (
SELECT 1 FROM chat_sessions WHERE chat_sessions.id = chat_spawns.session_id
)",
[],
)?;
Ok(())
}
/// Spawns sitting in one state, oldest request first.
pub fn list_chat_spawns(&self, state: ChatSpawnState) -> Result<Vec<ChatSpawn>> {
let mut stmt = self.conn.prepare(
"SELECT session_id, parent_session_id, prompt, wake_parent, attempts, finished_at
FROM chat_spawns WHERE state = ?1 ORDER BY requested_at, session_id",
)?;
let rows = stmt.query_map(params![state.as_str()], row_to_chat_spawn)?;
Ok(rows.collect::<std::result::Result<Vec<_>, _>>()?)
}
pub fn get_chat_spawn(&self, session_id: &str) -> Result<Option<ChatSpawn>> {
let mut stmt = self.conn.prepare(
"SELECT session_id, parent_session_id, prompt, wake_parent, attempts, finished_at
FROM chat_spawns WHERE session_id = ?1",
)?;
let mut rows = stmt.query_map(params![session_id], row_to_chat_spawn)?;
Ok(rows.next().transpose()?)
}
/// Take ownership of a spawn's next step, moving it into the claimed state
/// and returning the token that proves the claim. `None` means another
/// watcher got there first.
pub fn claim_chat_spawn(
&self,
session_id: &str,
from: ChatSpawnState,
to: ChatSpawnState,
) -> Result<Option<String>> {
let token = uuid::Uuid::new_v4().to_string();
let claimed = self.conn.execute(
"UPDATE chat_spawns SET state = ?3, claim_token = ?4, claimed_at = ?5
WHERE session_id = ?1 AND state = ?2",
params![session_id, from.as_str(), to.as_str(), token, now_ms()],
)?;
Ok((claimed == 1).then_some(token))
}
/// Finish (or hand back) a claimed step. `false` means the claim expired
/// and someone else reclaimed the row, so the caller's work is NOT recorded
/// and must not be treated as done.
pub fn settle_chat_spawn(
&self,
session_id: &str,
token: &str,
to: ChatSpawnState,
) -> Result<bool> {
let settled = self.conn.execute(
"UPDATE chat_spawns SET state = ?3, claim_token = NULL, claimed_at = NULL
WHERE session_id = ?1 AND claim_token = ?2",
params![session_id, token, to.as_str()],
)?;
Ok(settled == 1)
}
/// Whether any `orx up` currently holds a turn lease on this session.
///
/// The durable counterpart to `ChatHost::is_busy`, which reads a map this
/// process owns: a helper started by a different `orx up` — or by one that
/// has since restarted — is busy without this process knowing it.
pub fn chat_turn_leased(&self, chat_session_id: &str) -> Result<bool> {
let leased = self.conn.query_row(
"SELECT 1 FROM chat_turn_leases WHERE chat_session_id = ?1 AND heartbeat_at >= ?2",
params![chat_session_id, now_ms() - CHAT_TURN_LEASE_TTL_MS],
|_| Ok(()),
);
Ok(matches!(leased, Ok(())))
}
/// Record that a spawned helper's turn ended. An unstamped row whose lease
/// has lapsed is how the watcher recognizes an `orx up` that died mid-turn,
/// so the parent hears "interrupted" rather than a false completion.
pub fn mark_chat_spawn_finished(&self, chat_session_id: &str) -> Result<()> {
self.conn.execute(
"UPDATE chat_spawns SET finished_at = ?2
WHERE session_id = ?1 AND finished_at IS NULL",
params![chat_session_id, now_ms()],
)?;
Ok(())
}
pub fn record_chat_spawn_attempt(&self, chat_session_id: &str) -> Result<i64> {
self.conn.execute(
"UPDATE chat_spawns SET attempts = attempts + 1 WHERE session_id = ?1",
params![chat_session_id],
)?;
Ok(self.conn.query_row(
"SELECT attempts FROM chat_spawns WHERE session_id = ?1",
params![chat_session_id],
|row| row.get(0),
)?)
}
/// Helpers this parent has in flight, for the fan-out cap.
pub fn count_live_chat_spawns(&self, parent_session_id: &str) -> Result<i64> {
Ok(self.conn.query_row(
"SELECT count(*) FROM chat_spawns
WHERE parent_session_id = ?1 AND state != 'done'",
params![parent_session_id],
|row| row.get(0),
)?)
}
pub fn claim_chat_turn(&self, chat_session_id: &str, token: &str) -> Result<bool> {
self.clear_stale_data_dir_move_lease()?;
self.conn.execute(
@@ -1227,6 +1438,11 @@ impl Store {
/// experiments) in one transaction. GitHub repo and cache clone are kept.
pub fn delete_local_project(&self, id: &str) -> Result<()> {
let tx = self.begin()?;
self.conn.execute(
"DELETE FROM chat_spawns WHERE session_id IN
(SELECT id FROM chat_sessions WHERE project_id = ?1)",
params![id],
)?;
self.conn.execute(
"DELETE FROM chat_messages WHERE session_id IN
(SELECT id FROM chat_sessions WHERE project_id = ?1)",
@@ -1339,8 +1555,8 @@ impl Store {
self.conn.execute(
"INSERT INTO chat_sessions (id, project_id, harness, native_session_id, title, title_source, model,
permission_mode, plan_mode, plan_reset_pending, reasoning_level, archived, bootstrap_context,
active_leaf_id, created_at, updated_at)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16)",
active_leaf_id, parent_session_id, created_at, updated_at)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17)",
params![
s.id,
s.project_id,
@@ -1356,6 +1572,7 @@ impl Store {
s.archived,
s.bootstrap_context,
s.active_leaf_id,
s.parent_session_id,
s.created_at,
s.updated_at,
],
@@ -1388,6 +1605,8 @@ impl Store {
"DELETE FROM chat_messages WHERE session_id = ?1",
params![id],
)?;
self.conn
.execute("DELETE FROM chat_spawns WHERE session_id = ?1", params![id])?;
self.conn
.execute("DELETE FROM chat_sessions WHERE id = ?1", params![id])?;
Ok(())
@@ -1679,8 +1898,10 @@ pub struct StoredChatSession {
pub native_session_id: Option<String>,
pub title: Option<String>,
/// Who wrote `title`: `"fallback"` (first-line placeholder), `"generated"`
/// (harness auto-title), `"user"` (Rename). NULL on legacy rows, which the
/// conditional setter treats as "unknown, don't overwrite".
/// (harness auto-title), `"user"` (explicitly chosen — a rename, or an
/// agent's `orx agent spawn --title`, neither of which auto-titling may
/// overwrite). NULL on legacy rows, which the conditional setter treats as
/// "unknown, don't overwrite".
pub title_source: Option<String>,
pub model: Option<String>,
/// Permission-mode wire id (`"auto"` / `"plan"` / …); None = harness default.
@@ -1703,6 +1924,9 @@ pub struct StoredChatSession {
/// Tip of the branch the UI is currently showing. Forked turns make the
/// transcript a tree; this picks which path through it is live.
pub active_leaf_id: Option<String>,
/// Session that spawned this one with `orx agent spawn`. `None` for
/// sessions the user started from the dashboard.
pub parent_session_id: Option<String>,
pub created_at: i64,
pub updated_at: i64,
}
@@ -1775,7 +1999,18 @@ fn upsert_chat_message_with(conn: &Connection, m: &StoredChatMessage) -> Result<
const CHAT_SESSION_COLS: &str = "id, project_id, harness, native_session_id, title, model, \
permission_mode, plan_mode, plan_reset_pending, reasoning_level, archived, context_usage_json, \
created_at, updated_at, title_source, bootstrap_context, active_leaf_id";
created_at, updated_at, title_source, bootstrap_context, active_leaf_id, parent_session_id";
fn row_to_chat_spawn(row: &rusqlite::Row<'_>) -> std::result::Result<ChatSpawn, rusqlite::Error> {
Ok(ChatSpawn {
session_id: row.get(0)?,
parent_session_id: row.get(1)?,
prompt: row.get(2)?,
wake_parent: row.get(3)?,
attempts: row.get(4)?,
finished_at: row.get(5)?,
})
}
fn row_to_chat_session(
row: &rusqlite::Row<'_>,
@@ -1798,6 +2033,7 @@ fn row_to_chat_session(
title_source: row.get(14)?,
bootstrap_context: row.get(15)?,
active_leaf_id: row.get(16)?,
parent_session_id: row.get(17)?,
})
}
@@ -2195,11 +2431,150 @@ mod tests {
context_usage_json: None,
bootstrap_context: None,
active_leaf_id: None,
parent_session_id: None,
created_at: 1,
updated_at: 1,
}
}
fn chat_spawn_fixture(session_id: &str, parent: &str) -> ChatSpawn {
ChatSpawn {
session_id: session_id.into(),
parent_session_id: parent.into(),
prompt: "Sweep the literature for LoRA rank ablations".into(),
wake_parent: true,
attempts: 0,
finished_at: None,
}
}
#[test]
fn spawned_sessions_record_their_parent() {
let dir = std::env::temp_dir().join(format!("orx-store-spawnrow-{}", uuid::Uuid::new_v4()));
let store = Store::open_at(dir.clone()).unwrap();
store
.create_chat_session(&chat_session_fixture("chat_parent"))
.unwrap();
let mut child = chat_session_fixture("chat_child");
child.parent_session_id = Some("chat_parent".into());
store.create_chat_session(&child).unwrap();
assert!(store
.get_chat_session("chat_parent")
.unwrap()
.unwrap()
.parent_session_id
.is_none());
assert_eq!(
store
.get_chat_session("chat_child")
.unwrap()
.unwrap()
.parent_session_id
.as_deref(),
Some("chat_parent")
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn only_one_claimant_advances_a_spawn() {
let dir =
std::env::temp_dir().join(format!("orx-store-spawnclaim-{}", uuid::Uuid::new_v4()));
let store = Store::open_at(dir.clone()).unwrap();
store
.create_chat_session(&chat_session_fixture("chat_child"))
.unwrap();
store
.create_chat_spawn(&chat_spawn_fixture("chat_child", "chat_parent"))
.unwrap();
let pending = store.list_chat_spawns(ChatSpawnState::Pending).unwrap();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].parent_session_id, "chat_parent");
assert!(pending[0].wake_parent);
let token = store
.claim_chat_spawn(
"chat_child",
ChatSpawnState::Pending,
ChatSpawnState::Starting,
)
.unwrap()
.expect("first claim wins");
assert!(store
.claim_chat_spawn(
"chat_child",
ChatSpawnState::Pending,
ChatSpawnState::Starting
)
.unwrap()
.is_none());
// A settle under the wrong token must not move the row.
assert!(!store
.settle_chat_spawn("chat_child", "not-the-token", ChatSpawnState::Running)
.unwrap());
assert!(store
.settle_chat_spawn("chat_child", &token, ChatSpawnState::Running)
.unwrap());
assert_eq!(
store
.list_chat_spawns(ChatSpawnState::Running)
.unwrap()
.len(),
1
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn pruning_reclaims_stale_spawns_and_drops_deleted_sessions() {
let dir =
std::env::temp_dir().join(format!("orx-store-spawnprune-{}", uuid::Uuid::new_v4()));
let store = Store::open_at(dir.clone()).unwrap();
for id in ["chat_stuck", "chat_gone"] {
store
.create_chat_session(&chat_session_fixture(id))
.unwrap();
store
.create_chat_spawn(&chat_spawn_fixture(id, "chat_parent"))
.unwrap();
}
// A watcher that claimed the start and then died.
store
.claim_chat_spawn(
"chat_stuck",
ChatSpawnState::Pending,
ChatSpawnState::Starting,
)
.unwrap()
.unwrap();
store
.conn
.execute(
"UPDATE chat_spawns SET claimed_at = ?1 WHERE session_id = 'chat_stuck'",
params![now_ms() - CHAT_SPAWN_CLAIM_TTL_MS - 1],
)
.unwrap();
// A session removed WITHOUT its spawn row (the `delete_local_project`
// path) is what prune's orphan sweep is actually for.
store
.conn
.execute("DELETE FROM chat_sessions WHERE id = 'chat_gone'", [])
.unwrap();
store.prune_chat_spawns().unwrap();
let pending = store.list_chat_spawns(ChatSpawnState::Pending).unwrap();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].session_id, "chat_stuck");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn title_source_roundtrips_and_defaults_to_none() {
let dir = std::env::temp_dir().join(format!("orx-store-titlesrc-{}", uuid::Uuid::new_v4()));
-1034
View File
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
-1030
View File
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
+2 -2
View File
@@ -26,8 +26,8 @@
html { background: #ffffff; }
html[data-theme="dark"] { background: #0e0c0c; }
</style>
<script type="module" crossorigin src="/assets/index-E_PWOcfp.js"></script>
<link rel="stylesheet" crossorigin href="/assets/index-Dlmhe_xj.css">
<script type="module" crossorigin src="/assets/index-CSO9Nj35.js"></script>
<link rel="stylesheet" crossorigin href="/assets/index-DrR-RKJ6.css">
</head>
<body>
<div
+5 -1
View File
@@ -1277,7 +1277,8 @@ export interface ChatSession {
harness: HarnessId;
title: string | null;
/** Who wrote `title`: `"fallback"` (first-line placeholder), `"generated"`
* (harness auto-title), `"user"` (rename). Null on legacy sessions. */
* (harness auto-title), `"user"` (explicitly chosen — a rename, or an agent's
* `orx agent spawn --title`). Null on legacy sessions. */
titleSource?: string | null;
model: string | null;
permissionMode: string | null;
@@ -1286,6 +1287,9 @@ export interface ChatSession {
reasoningLevel: string | null;
/** Hidden from the default Recents list, but fully intact and resumable. */
archived: boolean;
/** Session whose agent spawned this one with `orx agent spawn`; null for
* sessions the user started themselves. */
parentSessionId?: string | null;
createdAt: number;
updatedAt: number;
busy: boolean;
+102 -1
View File
@@ -930,6 +930,8 @@ interface ToolActivity {
litCall?: NonNullable<ReturnType<typeof parseOrxLit>>;
runIds?: string[];
experimentIds?: string[];
/** Chat sessions `orx agent spawn` created in this tool call. */
spawnedSessionIds?: string[];
}
type OpenTranscriptFile = (path: string, line?: number, exp?: string, ref?: string) => void;
@@ -1247,6 +1249,9 @@ function commandFilePath(
}
const UUID_PATTERN = "[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}";
/** Session ids as `orx agent spawn` prints them, so the spawn card can link to
* the session it started. The id is only in the output — never the command. */
const SPAWNED_SESSION_PATTERN = new RegExp(`\\bchat_(${UUID_PATTERN})\\b`, "gi");
const RUN_TARGET_PATTERN = `(?:${UUID_PATTERN}|[0-9a-f]{8})`;
interface ShellCommandSegment {
@@ -1374,6 +1379,16 @@ function commandInvokesOrx(command: string, args: string): boolean {
return new RegExp(`(?:^|[\\s($;])orx\\s+${args}\\b`, "i").test(command);
}
function spawnedSessionIds(output: string | undefined): string[] {
if (!output) return [];
const ids = new Set<string>();
for (const match of output.slice(0, TOOL_OUTPUT_SCAN_LIMIT).matchAll(SPAWNED_SESSION_PATTERN)) {
ids.add(match[0].toLowerCase());
if (ids.size >= TOOL_TARGET_LIMIT) break;
}
return [...ids];
}
function idsFromToolOutput(output: string | undefined, resource: "runs" | "experiments"): string[] {
if (!output) return [];
const ids = new Set<string>();
@@ -1595,6 +1610,14 @@ function toolActivity(part: ChatPart): ToolActivity {
return { kind: litCall.kind === "lit" ? "search" : "read", label, litCall };
}
if (commandInvokesOrx(command, "agent\\s+spawn")) {
return {
kind: "agent",
label: "Delegated a task to a new agent",
spawnedSessionIds: spawnedSessionIds(toolOutput),
};
}
const shellSegments = shellCommandSegments(command);
const shellInvocations = shellSegments.map((segment) => shellInvocation(segment.raw));
const readsExperimentStatus = commandInvokesOrx(command, "exp\\s+status");
@@ -1937,6 +1960,7 @@ function ToolActivityLabel({
activity,
onOpenFile,
onOpenRun,
onOpenSpawnedSession,
runExperimentName,
onOpenExperiment,
experimentName,
@@ -1944,6 +1968,7 @@ function ToolActivityLabel({
activity: ToolActivity;
onOpenFile?: OpenTranscriptFile;
onOpenRun?: (runId: string) => void;
onOpenSpawnedSession?: (sessionId: string) => void;
runExperimentName?: (runId: string) => string;
onOpenExperiment?: (experimentId: string) => void;
experimentName?: (experimentId: string) => string;
@@ -1998,6 +2023,42 @@ function ToolActivityLabel({
</>
);
}
if (activity.spawnedSessionIds?.length && onOpenSpawnedSession) {
const sessionIds = activity.spawnedSessionIds;
const single = sessionIds.length === 1;
const visibleSessionIds = sessionIds.slice(0, 3);
const hiddenSessions = sessionIds.slice(visibleSessionIds.length).map((sessionId, index) => ({
id: sessionId,
label: `agent ${visibleSessionIds.length + index + 1}`,
}));
return (
<>
{single ? "Delegated a task to " : "Delegated tasks to "}
{visibleSessionIds.map((sessionId, index) => (
<span key={sessionId}>
{index > 0 && ", "}
<button
className="tool-target"
title="Open the session this agent spawned"
onClick={(event) => {
event.preventDefault();
event.stopPropagation();
onOpenSpawnedSession(sessionId);
}}
>
{single ? "a new agent" : `agent ${index + 1}`}
</button>
</span>
))}
{hiddenSessions.length > 0 && (
<>
{", "}
<ToolTargetOverflow items={hiddenSessions} onOpen={onOpenSpawnedSession} targetType="agent sessions" />
</>
)}
</>
);
}
if (activity.runIds?.length) {
const runIds = runExperimentName
? activity.runIds.filter((runId) => Boolean(runExperimentName(runId)))
@@ -2124,6 +2185,7 @@ function activityInProgress(activity: ToolActivity): ToolActivity {
[/^Checked /, "Checking "],
[/^Built /, "Building "],
[/^Cancelled /, "Cancelling "],
[/^Delegated /, "Delegating "],
];
let label = activity.label;
for (const [pattern, replacement] of replacements) {
@@ -2149,6 +2211,7 @@ function permissionActivityLabel(tool: string | undefined, input: Record<string,
[/^Updated /, "Update "],
[/^Created /, "Create "],
[/^Deleted /, "Delete "],
[/^Delegated /, "Delegate "],
[/^Ran /, "Run "],
[/^Started /, "Start "],
[/^Waited /, "Wait "],
@@ -2277,6 +2340,7 @@ function squashableToolPartKey(part: ChatPart): string | null {
activity.litCall?.kind === "paper" ? activity.litCall.id ?? null : null,
activity.runIds ?? null,
activity.experimentIds ?? null,
activity.spawnedSessionIds ?? null,
]);
}
@@ -2301,6 +2365,7 @@ function ToolRow({
repeatCount = 1,
onOpenFile,
onOpenRun,
onOpenSpawnedSession,
runExperimentName,
onOpenExperiment,
experimentName,
@@ -2309,6 +2374,7 @@ function ToolRow({
repeatCount?: number;
onOpenFile?: OpenTranscriptFile;
onOpenRun?: (runId: string) => void;
onOpenSpawnedSession?: (sessionId: string) => void;
runExperimentName?: (runId: string) => string;
onOpenExperiment?: (experimentId: string) => void;
experimentName?: (experimentId: string) => string;
@@ -2333,6 +2399,7 @@ function ToolRow({
activity={activity}
onOpenFile={onOpenFile}
onOpenRun={onOpenRun}
onOpenSpawnedSession={onOpenSpawnedSession}
runExperimentName={runExperimentName}
onOpenExperiment={onOpenExperiment}
experimentName={experimentName}
@@ -2383,6 +2450,7 @@ function ToolGroup({
pendingTail,
onOpenFile,
onOpenRun,
onOpenSpawnedSession,
runExperimentName,
onOpenExperiment,
experimentName,
@@ -2391,6 +2459,7 @@ function ToolGroup({
pendingTail?: boolean;
onOpenFile?: OpenTranscriptFile;
onOpenRun?: (runId: string) => void;
onOpenSpawnedSession?: (sessionId: string) => void;
runExperimentName?: (runId: string) => string;
onOpenExperiment?: (experimentId: string) => void;
experimentName?: (experimentId: string) => string;
@@ -2428,6 +2497,7 @@ function ToolGroup({
activity={pendingActivity}
onOpenFile={onOpenFile}
onOpenRun={onOpenRun}
onOpenSpawnedSession={onOpenSpawnedSession}
runExperimentName={runExperimentName}
onOpenExperiment={onOpenExperiment}
experimentName={experimentName}
@@ -2443,6 +2513,7 @@ function ToolGroup({
part={parts[0]}
onOpenFile={onOpenFile}
onOpenRun={onOpenRun}
onOpenSpawnedSession={onOpenSpawnedSession}
runExperimentName={runExperimentName}
onOpenExperiment={onOpenExperiment}
experimentName={experimentName}
@@ -2465,6 +2536,7 @@ function ToolGroup({
activity={pendingActivity}
onOpenFile={onOpenFile}
onOpenRun={onOpenRun}
onOpenSpawnedSession={onOpenSpawnedSession}
runExperimentName={runExperimentName}
onOpenExperiment={onOpenExperiment}
experimentName={experimentName}
@@ -2504,6 +2576,7 @@ function ToolGroup({
repeatCount={count}
onOpenFile={onOpenFile}
onOpenRun={onOpenRun}
onOpenSpawnedSession={onOpenSpawnedSession}
runExperimentName={runExperimentName}
onOpenExperiment={onOpenExperiment}
experimentName={experimentName}
@@ -2879,6 +2952,7 @@ const Message = memo(function Message({
pendingTailToolId,
onOpenFile,
onOpenRun,
onOpenSpawnedSession,
runExperimentName,
onOpenExperiment,
experimentName,
@@ -2901,6 +2975,7 @@ const Message = memo(function Message({
pendingTailToolId?: string | null;
onOpenFile?: OpenTranscriptFile;
onOpenRun?: (runId: string) => void;
onOpenSpawnedSession?: (sessionId: string) => void;
runExperimentName?: (runId: string) => string;
onOpenExperiment?: (experimentId: string) => void;
experimentName?: (experimentId: string) => string;
@@ -3048,6 +3123,7 @@ const Message = memo(function Message({
pendingTailToolId,
onOpenFile,
onOpenRun,
onOpenSpawnedSession,
runExperimentName,
onOpenExperiment,
experimentName,
@@ -3073,6 +3149,7 @@ function renderParts(
pendingTailToolId?: string | null;
onOpenFile?: OpenTranscriptFile;
onOpenRun?: (runId: string) => void;
onOpenSpawnedSession?: (sessionId: string) => void;
runExperimentName?: (runId: string) => string;
onOpenExperiment?: (experimentId: string) => void;
experimentName?: (experimentId: string) => string;
@@ -3087,6 +3164,7 @@ function renderParts(
pendingTailToolId,
onOpenFile,
onOpenRun,
onOpenSpawnedSession,
runExperimentName,
onOpenExperiment,
experimentName,
@@ -3111,6 +3189,7 @@ function renderParts(
pendingTail={toolRun.some((part) => part.id === pendingTailToolId)}
onOpenFile={onOpenFile}
onOpenRun={onOpenRun}
onOpenSpawnedSession={onOpenSpawnedSession}
runExperimentName={runExperimentName}
onOpenExperiment={onOpenExperiment}
experimentName={experimentName}
@@ -3501,6 +3580,7 @@ const Transcript = memo(function Transcript({
busy,
onOpenFile,
onOpenRun,
onOpenSpawnedSession,
runExperimentName,
onOpenExperiment,
experimentName,
@@ -3520,6 +3600,7 @@ const Transcript = memo(function Transcript({
busy: boolean;
onOpenFile?: OpenTranscriptFile;
onOpenRun?: (runId: string) => void;
onOpenSpawnedSession?: (sessionId: string) => void;
runExperimentName?: (runId: string) => string;
onOpenExperiment?: (experimentId: string) => void;
experimentName?: (experimentId: string) => string;
@@ -3564,6 +3645,7 @@ const Transcript = memo(function Transcript({
pendingTailToolId={pendingTailTool?.messageId === m.id ? pendingTailTool.toolId : null}
onOpenFile={onOpenFile}
onOpenRun={onOpenRun}
onOpenSpawnedSession={onOpenSpawnedSession}
runExperimentName={runExperimentName}
onOpenExperiment={onOpenExperiment}
experimentName={experimentName}
@@ -3742,7 +3824,9 @@ function SessionRow({
className={`session-row relative flex items-center gap-2 w-full text-left py-[7px] px-2.5 rounded-md text-md text-text cursor-pointer select-none [&:hover]:bg-surface [&.active]:bg-surface [&.active]:font-medium [&_.session-dot]:w-3.5 [&_.session-dot]:inline-flex [&_.session-dot]:items-center [&_.session-dot]:justify-center [&_.session-dot]:shrink-0 [&_.session-title]:flex-1 [&_.session-title]:min-w-0 [&_.session-title]:overflow-hidden [&_.session-title]:text-ellipsis [&_.session-title]:whitespace-nowrap [&.unread_.session-title]:font-semibold [&_.session-time]:text-2xs [&_.session-time]:text-muted [&_.session-time]:shrink-0 [&_.session-menu-btn]:hidden [&_.session-menu-btn]:items-center [&_.session-menu-btn]:justify-center [&_.session-menu-btn]:w-4 [&_.session-menu-btn]:h-4 [&_.session-menu-btn]:-my-0.5 [&_.session-menu-btn]:mx-0 [&_.session-menu-btn]:rounded-sm [&_.session-menu-btn]:text-muted [&_.session-menu-btn]:shrink-0 [&_.session-menu-btn:hover]:text-text [&_.session-menu-btn:hover]:bg-panel [&:hover_.session-menu-btn]:inline-flex [&:focus-within_.session-menu-btn]:inline-flex [&.menu-open_.session-menu-btn]:inline-flex [&:hover_.session-time]:hidden [&:focus-within_.session-time]:hidden [&.menu-open_.session-time]:hidden [&_.busy-dot]:w-[7px] [&_.busy-dot]:h-[7px] [&_.busy-dot]:rounded-full [&_.busy-dot]:bg-primary [&_.busy-dot]:animate-[or-pulse_1.2s_infinite] [&_.busy-dot]:shrink-0 [&_.unread-dot]:w-[7px] [&_.unread-dot]:h-[7px] [&_.unread-dot]:rounded-full [&_.unread-dot]:bg-primary [&_.unread-dot]:shrink-0 [&_.busy-dot.waiting]:animate-none [&_.session-title-input]:flex-1 [&_.session-title-input]:min-w-0 [&_.session-title-input]:py-px [&_.session-title-input]:px-[5px] [&_.session-title-input]:-my-0.5 [&_.session-title-input]:mx-0 [&_.session-title-input]:[font:inherit] [&_.session-title-input]:text-text [&_.session-title-input]:bg-background [&_.session-title-input]:border [&_.session-title-input]:border-primary [&_.session-title-input]:rounded-sm [&_.session-title-input]:outline-none [&.editing]:bg-surface [&.editing]:cursor-default [&.editing_.session-menu-btn]:hidden [&.editing_.session-time]:hidden ${active ? "active" : ""} ${unread ? "unread" : ""} ${open ? "menu-open" : ""} ${
editing ? "editing" : ""
}`}
title={`${HARNESS_LABELS[session.harness]}${session.model ? ` · ${session.model}` : ""}`}
title={`${HARNESS_LABELS[session.harness]}${session.model ? ` · ${session.model}` : ""}${
session.parentSessionId ? " · Spawned by another agent" : ""
}`}
onClick={() => {
// While editing, a body click is a no-op; blur/Enter/Esc drive it.
if (editing) return;
@@ -3772,6 +3856,9 @@ function SessionRow({
unread && <span className="unread-dot" />
)}
</span>
{session.parentSessionId && !editing && (
<Users className="text-muted shrink-0" size={12} aria-hidden />
)}
{editing ? (
<input
ref={inputRef}
@@ -5194,6 +5281,19 @@ export function ChatPanel({
onSelectMainView("chat");
}, [onSelectMainView]);
/** Follow a spawn card into the session it started. Spawned sessions are
* ordinary top-level sessions, so this is just a switch in the rail — via
* "All", because selecting a row the active filter hides would leave the
* thread keyed to a session with no row (see `setArchived`). */
const openSpawnedSession = useCallback(
(sessionId: string) => {
setSessionFilter("all");
setActiveId(sessionId);
onSelectMainView("chat");
},
[onSelectMainView],
);
useEffect(() => {
const onKeyDown = (event: KeyboardEvent) => {
if (
@@ -5474,6 +5574,7 @@ export function ChatPanel({
busy={busy}
onOpenFile={openFileInSession}
onOpenRun={onOpenRun}
onOpenSpawnedSession={openSpawnedSession}
runExperimentName={runExperimentName}
onOpenExperiment={onOpenExperiment}
experimentName={experimentName}