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
9 changes: 9 additions & 0 deletions crates/agent/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1382,6 +1382,15 @@ pub async fn start_session(
provider: ProviderKind,
opts: SessionOptions,
) -> Result<SessionHandle, AgentError> {
// Spawning into a missing directory fails with the OS's "not found" (or
// Windows' "invalid directory"), which every provider then reports as its
// binary being missing. Worktrees are routinely removed under old threads.
if !opts.cwd.is_dir() {
return Err(AgentError::Spawn(format!(
"working directory `{}` no longer exists",
opts.cwd.display()
)));
}
match provider {
ProviderKind::Codex => codex::start(opts).await,
ProviderKind::ClaudeCode => claude::start(opts).await,
Expand Down
56 changes: 39 additions & 17 deletions crates/runtime/src/app/orchestrate.rs
Original file line number Diff line number Diff line change
Expand Up @@ -953,6 +953,8 @@ impl AppState {
});
let state = if running {
"running"
} else if trailing_start_error(timeline).is_some() {
"failed"
} else {
match timeline.last_turn_status {
Some(TurnStatus::Completed) => "completed",
Expand Down Expand Up @@ -1020,6 +1022,9 @@ impl AppState {
"last_output_tail": tail_chars(&final_message, 600),
"updated_at": meta.updated_at,
});
if let Some(error) = trailing_start_error(timeline) {
status["start_error"] = serde_json::json!(error);
}
if let Some(usage) = usage.as_ref() {
status["tokens"] = token_usage_json(usage);
}
Expand Down Expand Up @@ -1099,23 +1104,32 @@ impl AppState {
return;
}
let turn = timeline.turns.len();
if state.callback_last_turn.get(&child_id).copied() == Some(turn) {
return;
}
// A report pushed via the child's report_result tool supersedes
// the last-message digest and is delivered in full; consuming it
// here keeps the fallback per turn.
let reported = state.child_reported_results.remove(&child_id);
let text = assemble_callback_text(
&child_id,
&title,
status,
&final_assistant_message(&timeline),
reported.as_deref(),
timeline.usage.as_ref(),
result_max_chars,
auto_archive,
);
// A failed start folds into the turn that was already reported,
// so it is not deduplicated against that turn's callback; and its
// final message is that old turn's, so it is not repeated.
let text = if let Some(error) = trailing_start_error(&timeline) {
format!(
"[orchestrate] thread {child_id} (\"{title}\") failed to start: {error}"
)
} else {
if state.callback_last_turn.get(&child_id).copied() == Some(turn) {
return;
}
// A report pushed via the child's report_result tool supersedes
// the last-message digest and is delivered in full; consuming it
// here keeps the fallback per turn.
let reported = state.child_reported_results.remove(&child_id);
assemble_callback_text(
&child_id,
&title,
status,
&final_assistant_message(&timeline),
reported.as_deref(),
timeline.usage.as_ref(),
result_max_chars,
auto_archive,
)
};
state.callback_last_turn.insert(child_id.clone(), turn);
state.deliver_orchestrate_callback_to_parent(&parent_id, text, cx);
if auto_archive
Expand Down Expand Up @@ -1743,6 +1757,14 @@ pub(super) fn final_assistant_message(timeline: &Timeline) -> String {
parts.concat()
}

/// The error of a provider start that failed after the thread's last entry.
fn trailing_start_error(timeline: &Timeline) -> Option<&str> {
match &timeline.entries.last()?.content {
EntryContent::ProviderStartError { error } => Some(error),
_ => None,
}
}

pub(super) fn tail_chars(text: &str, max: usize) -> String {
let count = text.chars().count();
text.chars().skip(count.saturating_sub(max)).collect()
Expand Down
99 changes: 99 additions & 0 deletions crates/runtime/src/app/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3666,6 +3666,105 @@ fn orchestrate_send_reactivates_the_child_and_its_settled_parent() {
});
}

/// An orchestrator returning to a long-idle child whose worktree has since been
/// removed must hear that the child could not start and why, not nothing (the
/// failure folds into the already-reported turn) or that turn's old output.
#[test]
fn send_to_child_whose_cwd_was_removed_reports_the_start_failure() {
let cx = &mut TestAppContext::default();
let test_store = TestStore::new("tcode-orchestrate-send-missing-cwd-test");
let store = (*test_store).clone();
let state = cx.new_entity(TestClientState::new(store));
let (parent_commands, parent_receiver) = smol::channel::unbounded();
let (child_commands, _child_receiver) = smol::channel::unbounded();
let cwd = std::env::temp_dir().join(format!("tcode-removed-{}", uuid::Uuid::new_v4()));

state.update(cx, |state, cx| {
let mut parent = live_session(ProviderKind::Codex, parent_commands);
parent.meta.id = "parent".into();
parent.turn_in_flight = true;
state.sessions.push(parent.meta.clone());
state
.residents
.parked
.insert(parent.meta.id.clone(), parent);

let mut child = live_session(ProviderKind::Codex, child_commands);
child.meta.id = "child".into();
child.meta.parent_session_id = Some("parent".into());
child.meta.archive_on_complete = false;
child.meta.cwd = cwd.clone();
child.turn_in_flight = true;
state.sessions.push(child.meta.clone());
state.residents.parked.insert(child.meta.id.clone(), child);

state.on_event("child", persisted_assistant_event("old report"), cx);
state.on_event(
"child",
AgentEvent::TurnCompleted {
turn_id: "turn-1".into(),
status: TurnStatus::Completed,
usage: None,
},
cx,
);
});
cx.run_until(|state| state.callback_last_turn.contains_key("child"));
assert!(matches!(
parent_receiver.try_recv(),
Ok(SessionCommand::Steer { text, .. }) if text.ends_with("\nold report")
));

state.update(cx, |state, cx| {
state.resident_mut("child").unwrap().shutdown_to_idle();
let (reply, response) = smol::channel::bounded(1);
state.handle_orchestrate_op(
orchestrate_mcp::OrchestrateOp::Send {
parent_id: "parent".into(),
thread_id: "child".into(),
message: "continue".into(),
fast: None,
},
reply,
cx,
);
assert!(response.try_recv().unwrap().is_ok());
});
cx.run_until(|_| !parent_receiver.is_empty());

let Ok(SessionCommand::Steer { text, .. }) = parent_receiver.try_recv() else {
panic!("the parent must be told the child failed to start");
};
assert_eq!(
text,
format!(
"[orchestrate] thread child (\"{}\") failed to start: failed to spawn provider process: working directory `{}` no longer exists",
state.read(|state| state.find_meta("child").unwrap().title.clone()),
cwd.display()
)
);

state.update(cx, |state, cx| {
let (reply, response) = smol::channel::bounded(1);
state.handle_orchestrate_op(
orchestrate_mcp::OrchestrateOp::Status {
parent_id: "parent".into(),
thread_id: Some("child".into()),
},
reply,
cx,
);
let status = response.try_recv().unwrap().unwrap();
assert_eq!(status[0]["state"], "failed");
assert!(
status[0]["start_error"]
.as_str()
.unwrap()
.contains("no longer exists")
);
});
}

#[test]
fn orchestrate_send_fast_switch_persists_and_schedules_restart() {
let cx = &mut TestAppContext::default();
Expand Down
Loading