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
89 changes: 65 additions & 24 deletions cmux-tui/crates/cmux-tui-core/src/mux.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24870,12 +24870,23 @@ mod tests {
assert_eq!(records[0].source, AgentSource::Hook);
assert_eq!(records[0].session.as_deref(), Some("racing-hook"));
assert_eq!(mux.resource_agent_projection_count_for_test().unwrap(), 1);
// The socket report commits its own revision only when it wins the
// race; a socket report that lands after the hook is retained by the
// hook-owned record without a new revision. Either way the hook's
// commit is the last batch.
let batches = mux.resource_events_after(revision).unwrap().batches;
assert_eq!(batches.len(), 2);
assert_eq!(batches[0].revision, revision + 1);
assert_eq!(batches[1].revision, revision + 2);
assert_eq!(batches[1].changes[0]["value"]["source"], "hook");
assert_eq!(batches[1].changes[0]["value"]["state"], "blocked");
assert_eq!(
u64::try_from(batches.len()).unwrap(),
hook_commit.revision - revision,
"one batch per committed revision"
);
for (offset, batch) in batches.iter().enumerate() {
assert_eq!(batch.revision, revision + 1 + u64::try_from(offset).unwrap());
}
let last = batches.last().unwrap();
assert_eq!(last.revision, hook_commit.revision);
assert_eq!(last.changes[0]["value"]["source"], "hook");
assert_eq!(last.changes[0]["value"]["state"], "blocked");
}

#[test]
Expand Down Expand Up @@ -24933,6 +24944,16 @@ mod tests {
assert_eq!(records[0].state, AgentState::Idle);
}

/// The hook projector's live record for `terminal_id`. `list_agents` reads
/// the journal-folded roster, which direct `apply_agent_hook_record` calls
/// bypass, so projector tests inspect the projector's own output.
fn hook_projected_agent(
mux: &Mux,
terminal_id: &TerminalPublicId,
) -> Option<TerminalAgentRecord> {
mux.agent_records.lock().unwrap().get(terminal_id).cloned()
}

#[test]
fn stale_same_terminal_hook_sequence_cannot_overwrite_newer_state() {
let mux = test_mux();
Expand All @@ -24959,7 +24980,7 @@ mod tests {
.unwrap();
mux.apply_agent_hook_record(&newer, 2).unwrap();
mux.apply_agent_hook_record(&older, 1).unwrap();
assert_eq!(mux.list_agents(Some(surface.id), None)[0].state, AgentState::Blocked);
assert_eq!(hook_projected_agent(&mux, &terminal_id).unwrap().state, AgentState::Blocked);
}

#[test]
Expand All @@ -24985,8 +25006,11 @@ mod tests {
// must not try to acquire that guard again.
mux.apply_agent_hook_record(&hook, 1).unwrap();

assert_eq!(mux.list_agents(Some(surface.id), None)[0].state, AgentState::Working);
assert_eq!(mux.list_agents(Some(surface.id), None)[0].agent.as_deref(), Some("claude"));
assert_eq!(hook_projected_agent(&mux, &terminal_id).unwrap().state, AgentState::Working);
assert_eq!(
hook_projected_agent(&mux, &terminal_id).unwrap().agent.as_deref(),
Some("claude")
);
let snapshot = crate::resource_api::public_session_snapshot(&mux).unwrap();
assert_eq!(snapshot["agents"][0]["extra"]["agent"], serde_json::json!("claude"));
assert_eq!(
Expand Down Expand Up @@ -25803,9 +25827,28 @@ mod tests {
assert_eq!(mux.list_agents(Some(surface.id), None).len(), 1);
assert_eq!(mux.resource_agent_projection_count_for_test().unwrap(), 0);
mux.workspace_registry.lock().unwrap().set_resource_patch_failure(false).unwrap();

// Startup runs this reconciliation once restored surfaces exist.
// This test's terminal is a local PTY, which does not survive a
// daemon restart (only host-owned terminals are adopted), so run
// the startup repair against the live surface here.
mux.reconcile_agent_roster_projections();
let repaired = mux.list_agents(Some(surface.id), None);
assert_eq!(repaired.len(), 1);
assert_eq!(repaired[0].state, AgentState::Working);
assert_eq!(repaired[0].source, AgentSource::Plugin);
assert_eq!(repaired[0].agent.as_deref(), Some("codex"));
assert_eq!(repaired[0].session.as_deref(), Some("pid:42"));
assert_eq!(mux.resource_agent_projection_count_for_test().unwrap(), 1);
// A healthy repeat is a no-op.
let revision = mux.with_state(|state| state.resource_revision);
mux.reconcile_agent_roster_projections();
assert_eq!(mux.with_state(|state| state.resource_revision), revision);
mux.shutdown();
}

// The roster itself is durable: a restart restores the plugin entry
// from the reducer checkpoint and journal.
let registry = WorkspaceRegistry::open(&root, session).unwrap();
let reopened = Mux::from_workspace_registry(
session.into(),
Expand All @@ -25815,13 +25858,12 @@ mod tests {
true,
)
.unwrap();
let repaired = reopened.list_agents(None, None);
assert_eq!(repaired.len(), 1);
assert_eq!(repaired[0].state, AgentState::Working);
assert_eq!(repaired[0].source, AgentSource::Plugin);
assert_eq!(repaired[0].agent.as_deref(), Some("codex"));
assert_eq!(repaired[0].session.as_deref(), Some("pid:42"));
assert_eq!(reopened.resource_agent_projection_count_for_test().unwrap(), 1);
{
let host = reopened.agent_roster.lock().unwrap();
let entry = host.roster.entries.get(terminal_id.as_str()).expect("restored roster");
assert_eq!(entry.agent_source(), AgentSource::Plugin);
assert_eq!(entry.agent_state(), AgentState::Working);
}
reopened.shutdown();
drop(reopened);
std::fs::remove_dir_all(root).unwrap();
Expand Down Expand Up @@ -26424,9 +26466,8 @@ mod tests {
mux.apply_agent_hook_record(&ingress("SessionEnd", "old"), 1).unwrap();
mux.apply_agent_hook_record(&ingress("SessionStart", "new"), 2).unwrap();
mux.apply_agent_hook_record(&ingress("UserPromptSubmit", "old"), 3).unwrap();
let records = mux.list_agents(Some(surface.id), None);
assert_eq!(records.len(), 1);
assert_eq!(records[0].state, AgentState::Idle);
let record = hook_projected_agent(&mux, &terminal_id).expect("new session record");
assert_eq!(record.state, AgentState::Idle);
assert!(
!mux.agent_hook_fences
.lock()
Expand Down Expand Up @@ -26455,11 +26496,11 @@ mod tests {
// The old session can arrive after its end marker. Matching the ended
// identity is still a stale event, not permission to reopen it.
mux.apply_agent_hook_record(&ingress("SessionStart", "old"), 3).unwrap();
assert!(mux.list_agents(Some(surface.id), None).is_empty());
assert!(hook_projected_agent(&mux, &terminal_id).is_none());
assert!(mux.agent_hook_fences.lock().unwrap()[&terminal_id].ended);
mux.apply_agent_hook_record(&ingress("SessionStart", "new"), 4).unwrap();
mux.apply_agent_hook_record(&ingress("SessionStart", "old"), 5).unwrap();
assert_eq!(mux.list_agents(Some(surface.id), None)[0].state, AgentState::Idle);
assert_eq!(hook_projected_agent(&mux, &terminal_id).unwrap().state, AgentState::Idle);
assert_eq!(mux.agent_hook_fences.lock().unwrap()[&terminal_id].session_id, "new");
}

Expand Down Expand Up @@ -26493,7 +26534,7 @@ mod tests {
mux.apply_agent_hook_record(&sessionless("UserPromptSubmit"), 4).unwrap();

assert_eq!(mux.agent_hook_fences.lock().unwrap()[&terminal_id].session_id, "new");
assert_eq!(mux.list_agents(Some(surface.id), None)[0].state, AgentState::Idle);
assert_eq!(hook_projected_agent(&mux, &terminal_id).unwrap().state, AgentState::Idle);
}

#[test]
Expand All @@ -26520,7 +26561,7 @@ mod tests {
let second = mux.agent_hook_fences.lock().unwrap()[&terminal_id].session_id.clone();
assert_ne!(first, second);
mux.apply_agent_hook_record(&ingress("UserPromptSubmit"), 4).unwrap();
assert_eq!(mux.list_agents(Some(surface.id), None)[0].state, AgentState::Working);
assert_eq!(hook_projected_agent(&mux, &terminal_id).unwrap().state, AgentState::Working);
}

#[test]
Expand Down Expand Up @@ -26646,12 +26687,12 @@ mod tests {
mux.apply_agent_hook_record(&ingress, 7).unwrap();
assert_eq!(mux.workspace_registry.lock().unwrap().agent_hook_apply_cursor().unwrap(), 7);
assert_eq!(mux.agent_hook_fences.lock().unwrap()[&terminal_id].sequence, 7);
assert_eq!(mux.list_agents(Some(surface.id), None).len(), 1);
assert!(hook_projected_agent(&mux, &terminal_id).is_some());

// A replay of the committed sequence is an idempotent no-op.
mux.apply_agent_hook_record(&ingress, 7).unwrap();
assert_eq!(mux.workspace_registry.lock().unwrap().agent_hook_apply_cursor().unwrap(), 7);
assert_eq!(mux.list_agents(Some(surface.id), None).len(), 1);
assert!(hook_projected_agent(&mux, &terminal_id).is_some());
}

#[test]
Expand Down
6 changes: 4 additions & 2 deletions cmux-tui/crates/cmux-tui-core/src/resource_api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1102,7 +1102,9 @@ mod tests {
#[cfg(unix)]
#[test]
fn cloud_cwd_live_osc7_reaches_snapshot_and_event_feed() {
let mux = Mux::new_for_test(
// `Mux::new_for_test` surfaces never run their command, so these OSC 7
// tests need the real local PTY runtime.
let mux = Mux::new(
"cloud-cwd-osc",
SurfaceOptions {
command: Some(vec![
Expand Down Expand Up @@ -1139,7 +1141,7 @@ mod tests {
// A shell that reports a directory and later reports none (an empty
// OSC 7, as when it leaves the host it described) must clear the
// published cwd through the same incremental parser path.
let mux = Mux::new_for_test(
let mux = Mux::new(
"cloud-cwd-osc-clear",
SurfaceOptions {
command: Some(vec![
Expand Down
7 changes: 5 additions & 2 deletions cmux-tui/crates/cmux-tui-core/src/resource_router/topology.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1103,7 +1103,9 @@ mod tests {
),
)
.unwrap_err();
assert_eq!(error.code, "mutation.indeterminate");
// The close transaction rolled back before any effect ran, so its
// failure is committed as the key's durable outcome.
assert_eq!(error.code, "operation.failed");
mux.set_resource_patch_failure_for_test(false);

let replay = dispatch(
Expand All @@ -1116,7 +1118,8 @@ mod tests {
),
)
.unwrap_err();
assert_eq!(replay.code, "mutation.indeterminate");
assert_eq!(replay.code, error.code);
assert_eq!(replay.message, error.message);

assert_eq!(mux.with_state(|state| state.resource_revision), before_resource);
assert_eq!(mux.with_state(|state| state.workspace_revision), before_workspace);
Expand Down
20 changes: 9 additions & 11 deletions cmux-tui/crates/cmux-tui-core/src/surface.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4548,11 +4548,11 @@ impl Surface {
drop(runtime);
receipt.wait().map_err(ConfirmedInputFailure::Indeterminate)
}
// Same keep-on-exit contract as `write_bytes`: the child is gone,
// so there is no reader to deliver to and nothing to retry.
// Input to the final screen is a successful no-op.
#[cfg(unix)]
PtyRuntime::ExitedHosted => Err(ConfirmedInputFailure::Known(std::io::Error::new(
std::io::ErrorKind::NotConnected,
"terminal has no live PTY owner for receipted input",
))),
PtyRuntime::ExitedHosted => Ok(()),
}
}

Expand Down Expand Up @@ -7682,7 +7682,7 @@ mod tests {

#[cfg(unix)]
#[test]
fn receipted_input_rejects_an_exited_host_before_effect() {
fn receipted_input_to_an_exited_host_is_a_no_op() {
let mux = Mux::new_for_test("receipted-input-exited-host", SurfaceOptions::default());
let surface =
Surface::spawn_for_test(1, SurfaceOptions::default(), Arc::downgrade(&mux)).unwrap();
Expand All @@ -7692,12 +7692,10 @@ mod tests {
*runtime = PtyRuntime::ExitedHosted;
}

let error = surface.write_bytes_confirmed(b"must-not-drop").unwrap_err();
let ConfirmedInputFailure::Known(error) = error else {
panic!("exited-host rejection became indeterminate");
};
assert_eq!(error.kind(), std::io::ErrorKind::NotConnected);
assert!(error.to_string().contains("no live PTY owner"));
// A keep-on-exit terminal keeps its final screen after the child
// exits; typing there succeeds without an effect, as `write_bytes`
// does.
surface.write_bytes_confirmed(b"ignored").unwrap();
}

#[test]
Expand Down
2 changes: 1 addition & 1 deletion cmux-tui/crates/cmux-tui/tests/resource_cli_v2.rs
Original file line number Diff line number Diff line change
Expand Up @@ -824,7 +824,7 @@ fn local_plugin_jsonl_never_connects_to_the_session_socket() {
let invalid = parse_single_json(&invalid.stderr);
assert_eq!(invalid["code"], "validation.invalid");
assert_eq!(invalid["retryable"], false);
assert_eq!(invalid["details"]["field"], "sidebar_plugin");
assert_eq!(invalid["details"]["field"], "plugin");
assert!(invalid["details"]["reason"].is_string());
fs::remove_dir_all(dir).unwrap();
}
Expand Down
6 changes: 6 additions & 0 deletions tests/test-execution.toml
Original file line number Diff line number Diff line change
Expand Up @@ -144,10 +144,16 @@ reason = "Requires an isolated tagged cmux app and helper, two prepared split wo
[[test]]
path = "tests/test_hermes_wrapper_hooks.py"
lane = "macos-cli-no-socket"
# Exercises the installer's hard wall-clock deadline and a deliberately slow
# start; parallel lane load can turn the intended delay into a false hang.
serial = true

[[test]]
path = "tests/test_claude_wrapper_mutual_shim_loop.py"
lane = "macos-cli-no-socket"
# The finite-shim regression has a five-second hang guard. Run it without the
# parallel pool so host scheduling pressure cannot masquerade as a shim loop.
serial = true

[[test]]
path = "tests/test_claude_wrapper_shim_root_survives_tmpdir_change.py"
Expand Down
Loading