use super::*;
use crate::controller::test_support::RefusingExecutor;
use tokio::io::AsyncWriteExt;
#[test]
fn newer_daemon_protocol_requires_updating_the_client() {
let metadata = |protocol_version| DaemonMetadata {
protocol_version,
pid: std::process::id(),
address: "127.0.0.1:1".parse().unwrap(),
token: "tok".into(),
started_at: "2026-09-01T07:48:14Z".into(),
build_version: "9.9.9".into(),
};
assert!(ensure_supported_daemon_protocol(&metadata(PROTOCOL_VERSION)).is_ok());
assert!(ensure_supported_daemon_protocol(&metadata(PROTOCOL_VERSION - 1)).is_ok());
let error = ensure_supported_daemon_protocol(&metadata(PROTOCOL_VERSION + 1)).unwrap_err();
let message = error.to_string();
assert!(message.contains("version 9.9.9"), "{message}");
assert!(message.contains("first on PATH"), "{message}");
}
#[test]
fn a_lifecycle_failure_carries_a_refusal_across_its_result_channel() {
let refused = LifecycleFailure::of(
&anyhow::Error::new(Refusal::precondition(
"repository \"app\" needs a network Git remote",
))
.context("provision the session target"),
);
let rebuilt = refused.clone().into_error();
assert_eq!(
Refusal::of(&rebuilt).map(|refusal| refusal.message().to_owned()),
Some("repository \"app\" needs a network Git remote".to_owned()),
"a waiter reading the channel must still see the reason, not a bare string"
);
let internal = LifecycleFailure::of(&anyhow::anyhow!("ssh host build-07 refused"));
assert!(internal.refusal.is_none());
assert!(internal.detail.contains("build-07"));
}
#[test]
fn graceful_close_retires_worker_polling_only_during_target_teardown() {
assert!(!lifecycle_owns_worker_target(
LifecycleKind::Suspend,
Some(SessionState::Running)
));
assert!(!lifecycle_owns_worker_target(
LifecycleKind::Suspend,
Some(SessionState::Checkpointing)
));
assert!(!lifecycle_owns_worker_target(
LifecycleKind::Suspend,
Some(SessionState::Closing)
));
assert!(lifecycle_owns_worker_target(
LifecycleKind::Suspend,
Some(SessionState::Destroying)
));
assert!(lifecycle_owns_worker_target(
LifecycleKind::ForceStop,
Some(SessionState::Running)
));
assert!(lifecycle_owns_worker_target(
LifecycleKind::ForceDestroy,
Some(SessionState::Running)
));
}
#[test]
fn a_close_past_its_verified_checkpoint_cannot_be_cancelled() {
assert!(lifecycle_cancellable(
LifecycleKind::Suspend,
Some(SessionState::Running)
));
assert!(lifecycle_cancellable(
LifecycleKind::Suspend,
Some(SessionState::Checkpointing)
));
assert!(lifecycle_cancellable(
LifecycleKind::Suspend,
Some(SessionState::Closing)
));
assert!(lifecycle_cancellable(LifecycleKind::Suspend, None));
assert!(!lifecycle_cancellable(
LifecycleKind::Suspend,
Some(SessionState::Destroying)
));
assert!(lifecycle_cancellable(
LifecycleKind::ForceStop,
Some(SessionState::Destroying)
));
assert!(lifecycle_cancellable(
LifecycleKind::Create,
Some(SessionState::Destroying)
));
assert!(lifecycle_cancellable(
LifecycleKind::Move,
Some(SessionState::Destroying)
));
}
#[test]
fn startup_reports_every_in_flight_state_that_no_operation_owns() {
let target = Some(mj_core::state::TargetLocator::LocalPodman {
borrowed_from: None,
container_id: "a".repeat(64),
workspace_storage: Default::default(),
});
let mut state = mj_core::state::State::default();
for (id, session_state, session_target) in [
("provisioning", SessionState::Provisioning, None),
("closing-orphan", SessionState::Closing, None),
("destroying-orphan", SessionState::Destroying, None),
("closing-owned", SessionState::Closing, target.clone()),
("destroying-owned", SessionState::Destroying, target.clone()),
("running", SessionState::Running, target.clone()),
("stopped", SessionState::Stopped, None),
("moving", SessionState::Provisioning, None),
] {
let mut session = runtime_test_session(id, "workspace", session_state);
session.target = session_target;
state.sessions.insert(id.into(), session);
}
let controller = Controller {
config: mj_core::config::Config::default(),
state,
};
let owned = ["moving".to_owned()].into_iter().collect();
let reconciled = unowned_interrupted_lifecycles(&controller, &owned);
let ids = reconciled
.iter()
.map(|(id, _)| id.as_str())
.collect::<Vec<_>>();
assert_eq!(ids, ["closing-orphan", "destroying-orphan", "provisioning"]);
let provisioning_cause = &reconciled
.iter()
.find(|(id, _)| id == "provisioning")
.expect("the interrupted provision is reported")
.1;
assert!(
provisioning_cause.contains("provisioning") && provisioning_cause.contains("recover scan"),
"the cause tells the user what happened and where the container went: {provisioning_cause}"
);
}
#[test]
fn a_stop_on_a_record_left_mid_close_routes_to_recovery() {
let target = Some(mj_core::state::TargetLocator::LocalPodman {
borrowed_from: None,
container_id: "a".repeat(64),
workspace_storage: Default::default(),
});
for state in [SessionState::Closing, SessionState::Destroying] {
let mut session = runtime_test_session("session", "workspace", state);
session.target = target.clone();
assert_eq!(close_route(Some(&session)), CloseRoute::RecoverInterrupted);
session.target = None;
assert_eq!(
close_route(Some(&session)),
CloseRoute::SettleWithoutCheckpoint
);
}
let mut running = runtime_test_session("session", "workspace", SessionState::Running);
running.target = target.clone();
assert_eq!(close_route(Some(&running)), CloseRoute::Graceful);
assert_eq!(close_route(None), CloseRoute::Graceful);
let mut provisioning = runtime_test_session("session", "workspace", SessionState::Provisioning);
assert_eq!(
close_route(Some(&provisioning)),
CloseRoute::SettleWithoutCheckpoint
);
provisioning.target = target.clone();
assert_eq!(
close_route(Some(&provisioning)),
CloseRoute::SettleWithoutCheckpoint
);
let failed = runtime_test_session("session", "workspace", SessionState::Error);
assert_eq!(
close_route(Some(&failed)),
CloseRoute::SettleWithoutCheckpoint
);
let mut failed_with_target = failed.clone();
failed_with_target.target = target.clone();
assert_eq!(close_route(Some(&failed_with_target)), CloseRoute::Graceful);
let mut stopped = runtime_test_session("session", "workspace", SessionState::Stopped);
assert_eq!(close_route(Some(&stopped)), CloseRoute::Done);
stopped.target = target;
assert_eq!(close_route(Some(&stopped)), CloseRoute::DeferredCleanup);
}
#[cfg(unix)]
#[test]
fn a_process_that_exited_but_was_not_reaped_counts_as_gone() {
let mut child = std::process::Command::new("true")
.spawn()
.expect("spawn a process that exits immediately");
let pid = child.id();
let deadline = std::time::Instant::now() + Duration::from_secs(5);
let gone = loop {
if !daemon_process_is_alive(pid) {
break true;
}
if std::time::Instant::now() >= deadline {
break false;
}
std::thread::sleep(Duration::from_millis(10));
};
assert!(
gone,
"an exited but unreaped process was reported as running"
);
let _ = child.wait();
}
#[cfg(unix)]
#[test]
fn attachment_liveness_probe_does_not_reap_children() {
use std::io::Read;
use std::process::Stdio;
let mut child = std::process::Command::new("true")
.stdout(Stdio::piped())
.spawn()
.expect("spawn a process that exits immediately");
let pid = child.id();
let mut output = Vec::new();
child
.stdout
.take()
.expect("capture child stdout")
.read_to_end(&mut output)
.expect("observe child exit");
let _ = process_is_alive(pid);
let status = child.wait().expect("attachment probe left child waitable");
assert!(status.success());
}
#[tokio::test]
async fn client_presence_is_global_and_detach_and_prune_remove_it() {
let state = test_runtime_state();
let metadata = test_metadata(SocketAddr::from((Ipv4Addr::LOCALHOST, 0)));
let cancellation = CancellationToken::new();
handle_action(
DaemonAction::Attach {
client_id: "client-a".into(),
pid: std::process::id(),
},
&metadata,
&state,
&cancellation,
)
.await
.expect("attach presence");
state
.attachments()
.insert("dead-client".into(), Attachment { pid: u32::MAX });
let DaemonReply::Status(status) =
handle_action(DaemonAction::Status, &metadata, &state, &cancellation)
.await
.expect("status")
else {
panic!("status action returned a different reply");
};
assert_eq!(status.attached_clients, 1);
handle_action(
DaemonAction::Detach {
client_id: "client-a".into(),
},
&metadata,
&state,
&cancellation,
)
.await
.expect("detach presence");
assert!(state.attachments().is_empty());
}
#[tokio::test]
async fn workspace_deletion_guard_ignores_global_client_presence() {
let state = test_runtime_state();
state.attachments().insert(
"client-a".into(),
Attachment {
pid: std::process::id(),
},
);
assert!(!state.workspace_has_active_resume("workspace-a"));
let (_completed, result) = tokio::sync::watch::channel(None);
state.owner().lifecycle.insert(
"session-a".into(),
ActiveLifecycle {
phase: LifecyclePhase::Executing,
operation_id: "resume-operation".into(),
create_control: None,
kind: LifecycleKind::Resume,
cancelled: Arc::new(AtomicBool::new(false)),
started_at_epoch_seconds: 1,
active_stages: BTreeMap::new(),
resume_workspace_id: Some("workspace-a".into()),
resume_destination: None,
notice: None,
request_key: None,
_move_guard: None,
result,
},
);
assert!(state.workspace_has_active_resume("workspace-a"));
}
#[tokio::test]
async fn checkpoint_lifecycle_guard_is_per_session() {
let state = test_runtime_state();
let (_completed, result) = tokio::sync::watch::channel(None);
state.owner().lifecycle.insert(
"session-b".into(),
ActiveLifecycle {
phase: LifecyclePhase::Executing,
operation_id: "resume-operation".into(),
create_control: None,
kind: LifecycleKind::Resume,
cancelled: Arc::new(AtomicBool::new(false)),
started_at_epoch_seconds: epoch_seconds().saturating_sub(45),
active_stages: BTreeMap::new(),
resume_workspace_id: None,
resume_destination: None,
notice: None,
request_key: None,
_move_guard: None,
result,
},
);
assert!(state.session_lifecycle_busy("session-a").is_none());
let busy = state
.session_lifecycle_busy("session-b")
.expect("session-b is running its own operation");
assert_eq!(busy.operation, "resume");
assert!(
(45..60).contains(&busy.age_seconds),
"the age is the operation's, not the clock's: {busy}"
);
let rename = ensure_no_active_lifecycle(&state).unwrap_err();
let message = format!("{rename:#}");
assert!(
message
.contains("cannot rename configuration while session session-b is busy with a resume"),
"the rename refusal names the operation: {message}"
);
}
#[cfg(target_os = "macos")]
#[test]
fn zombie_only_daemon_group_counts_as_gone() {
use std::io::Read;
use std::os::unix::process::CommandExt;
use std::process::Stdio;
let mut command = std::process::Command::new("true");
command.process_group(0).stdout(Stdio::piped());
let mut child = command
.spawn()
.expect("spawn process-group leader that exits immediately");
let pid = libc::pid_t::try_from(child.id()).expect("child PID fits pid_t");
let mut output = Vec::new();
child
.stdout
.take()
.expect("capture child stdout")
.read_to_end(&mut output)
.expect("observe child exit");
assert!(!owned_daemon_group_is_alive(pid));
child.wait().expect("reap process-group leader");
}
pub(super) fn test_runtime_state() -> Arc<RuntimeState> {
let remote = spawn_remote_session_manager().unwrap();
let recovery = crate::recovery::RecoveryCoordinator::spawn(remote.control.clone());
let upgrades = crate::worker_upgrade::WorkerUpgradeCoordinator::spawn(
remote.control.clone(),
&recovery.observer(),
);
Arc::new(RuntimeState::new_with_controller_loader(
remote.control,
Controller {
config: Config::default(),
state: mj_core::state::State::default(),
},
recovery.observer(),
upgrades.observer(),
Vec::new(),
|| {
Ok(Controller {
config: Config::default(),
state: mj_core::state::State::default(),
})
},
))
}
#[tokio::test]
async fn upgrade_blockers_names_the_daemon_owned_work_then_releases_it() {
const LABEL: &str = "test-only upgrade blocker";
let state = test_runtime_state();
let metadata = test_metadata("127.0.0.1:1".parse().unwrap());
let shutdown = CancellationToken::new();
let labels = |reply: DaemonReply| match reply {
DaemonReply::UpgradeBlockers(labels) => labels,
other => panic!("expected named blockers, got {other:?}"),
};
let held = crate::upgrade::activity(LABEL).unwrap();
let named = labels(
handle_action(DaemonAction::UpgradeBlockers, &metadata, &state, &shutdown)
.await
.unwrap(),
);
assert!(named.contains(&LABEL.to_owned()), "{named:?}");
drop(held);
let released = labels(
handle_action(DaemonAction::UpgradeBlockers, &metadata, &state, &shutdown)
.await
.unwrap(),
);
assert!(
!released.contains(&LABEL.to_owned()),
"released work is still named: {released:?}"
);
}
#[tokio::test]
async fn automatic_upgrade_drains_a_lifecycle_without_cancelling_it() {
const CHILD: &str = "MJ_UPGRADE_DRAIN_TEST_CHILD";
if std::env::var_os(CHILD).is_none() {
let root = tempfile::tempdir().unwrap();
let name = format!(
"{}::automatic_upgrade_drains_a_lifecycle_without_cancelling_it",
module_path!()
.strip_prefix("mj_controller::")
.unwrap_or(module_path!())
);
crate::controller::test_support::IsolatedTest::new(name)
.env(CHILD, "1")
.env("MJ_INSTANCE", "upgrade-in-flight-test")
.isolated_store(root.path())
.run();
return;
}
let _writer = crate::database::install_isolated_test_writer();
let state = test_runtime_state();
let metadata = test_metadata("127.0.0.1:1".parse().unwrap());
let shutdown = CancellationToken::new();
let release = Arc::new(tokio::sync::Notify::new());
let operation = state
.start_or_join_lifecycle("upgrade-provisioning".into(), LifecycleKind::Create, {
let release = release.clone();
move |_, _, cancelled| async move {
release.notified().await;
ensure!(
!cancelled.load(Ordering::Acquire),
"upgrade cancelled provisioning"
);
Ok(DaemonLifecycleResult::Done)
}
})
.unwrap();
tokio::time::pause();
for _ in 0..4 {
tokio::time::advance(Duration::from_secs(60 * 60 * 24)).await;
assert!(matches!(
handle_action(DaemonAction::PrepareUpgrade, &metadata, &state, &shutdown)
.await
.unwrap(),
DaemonReply::UpgradePending
));
assert!(!shutdown.is_cancelled());
assert!(matches!(
handle_action(DaemonAction::Ping, &metadata, &state, &shutdown)
.await
.unwrap(),
DaemonReply::Pong
));
assert!(
crate::upgrade::activity("another control request").is_ok(),
"waiting must leave controls available"
);
}
tokio::time::resume();
release.notify_one();
RuntimeState::wait_lifecycle_result(operation)
.await
.unwrap();
let (updates, mut receiver) = crate::session_manager::coalesced_update_channel();
updates.send(crate::session_manager::SessionManagerUpdate {
session_id: "completed-turn".into(),
view: ManagedSessionView::default(),
});
assert!(
matches!(
handle_action(DaemonAction::PrepareUpgrade, &metadata, &state, &shutdown)
.await
.unwrap(),
DaemonReply::UpgradePending
),
"a queued completion can still start a review or continuation"
);
receiver.recv().await.unwrap();
assert!(
matches!(
handle_action(DaemonAction::PrepareUpgrade, &metadata, &state, &shutdown)
.await
.unwrap(),
DaemonReply::UpgradePending
),
"the consumer still owns the completion until it asks for another update"
);
assert!(receiver.try_recv().is_err());
tokio::time::timeout(Duration::from_secs(5), async {
loop {
if matches!(
handle_action(DaemonAction::PrepareUpgrade, &metadata, &state, &shutdown)
.await
.unwrap(),
DaemonReply::Done
) {
break;
}
tokio::task::yield_now().await;
}
})
.await
.unwrap();
assert!(shutdown.is_cancelled());
assert!(crate::upgrade::activity("late operation").is_err());
}
struct TestRemoteManager {
control: SessionManagerControl,
requests: RemoteSessionRequests,
publisher: RemoteSessionPublisher,
_shutdown: SessionManagerShutdown,
_targets: tokio::sync::watch::Sender<Vec<RelaySessionTarget>>,
}
impl TestRemoteManager {
async fn new() -> Self {
let channels =
crate::session_manager::spawn_reply_fixture_session_manager().expect("remote manager");
let session_id = "session-1";
channels.targets.send_replace(vec![RelaySessionTarget {
session_id: session_id.to_owned(),
spec: CommandSpec::new("true", Vec::<String>::new()),
worker_recovery: None,
project_memory: None,
}]);
let manager = Self {
control: channels.control,
requests: channels.requests,
publisher: channels.publisher,
_shutdown: channels.shutdown,
_targets: channels.targets,
};
manager
.publisher
.publish(session_id.to_owned(), ManagedSessionView::default())
.await
.expect("publish test session view");
manager
.control
.wait_for_session(session_id, Duration::from_secs(5))
.await
.expect("remote manager creates test session");
manager
}
}
fn test_runtime_state_with_manager(manager: &TestRemoteManager) -> Arc<RuntimeState> {
let recovery = crate::recovery::RecoveryCoordinator::spawn(manager.control.clone());
let upgrades = crate::worker_upgrade::WorkerUpgradeCoordinator::spawn(
manager.control.clone(),
&recovery.observer(),
);
let mut runtime = RuntimeState::new_with_controller_loader(
manager.control.clone(),
Controller {
config: Config::default(),
state: mj_core::state::State::default(),
},
recovery.observer(),
upgrades.observer(),
Vec::new(),
|| {
Ok(Controller {
config: Config::default(),
state: mj_core::state::State::default(),
})
},
);
if std::env::var_os("MJ_TEST_DURABLE_STARTUP").is_some() {
runtime.committed = Some(crate::database::subscribe_committed_state().unwrap());
}
Arc::new(runtime)
}
fn test_metadata(address: SocketAddr) -> DaemonMetadata {
DaemonMetadata {
protocol_version: PROTOCOL_VERSION,
pid: 1,
address,
token: "right-token".into(),
started_at: "now".into(),
build_version: "test".into(),
}
}
#[tokio::test]
async fn in_process_reviewer_forwarding_stops_when_the_caller_goes_away() {
let mut manager = TestRemoteManager::new().await;
let (reply, response) = tokio::sync::oneshot::channel();
let forwarding = tokio::spawn(forward_in_process_session_request(
RemoteSessionRequest::Reviewer {
session_id: "session-1".into(),
role: None,
action: crate::session_manager::ReviewerAction::Status,
reply,
},
manager.control.clone(),
));
let request = tokio::time::timeout(Duration::from_secs(5), manager.requests.recv())
.await
.expect("reviewer forwarding did not reach the manager")
.expect("manager request stream ended");
let RemoteSessionRequest::Reviewer { mut reply, .. } = request else {
panic!("expected the forwarded reviewer request")
};
drop(response);
tokio::time::timeout(Duration::from_secs(5), reply.closed())
.await
.expect("in-process forwarding kept the actor reply alive");
forwarding.await.expect("forwarding task panicked");
}
#[tokio::test]
async fn daemon_drops_in_flight_reviewer_work_when_the_client_eof_arrives() {
let mut manager = TestRemoteManager::new().await;
let state = test_runtime_state_with_manager(&manager);
let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
let address = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
serve_client(
stream,
test_metadata(address),
state,
CancellationToken::new(),
)
.await
});
let mut stream = TcpStream::connect(address).await.unwrap();
write_frame(
&mut stream,
&RequestEnvelope {
protocol_version: PROTOCOL_VERSION,
request_id: 1,
token: "right-token".into(),
action: DaemonAction::ReviewerAction {
session_id: "session-1".into(),
role: None,
action: crate::session_manager::ReviewerAction::Status,
},
},
)
.await
.unwrap();
let request = tokio::time::timeout(Duration::from_secs(5), manager.requests.recv())
.await
.expect("reviewer request did not reach the manager")
.expect("manager request stream ended");
let RemoteSessionRequest::Reviewer { mut reply, .. } = request else {
panic!("expected a reviewer request")
};
drop(stream);
tokio::time::timeout(Duration::from_secs(5), reply.closed())
.await
.expect("daemon kept reviewer work alive after client EOF");
assert!(server.await.expect("daemon task panicked").is_ok());
}
#[tokio::test]
async fn daemon_peek_keeps_a_pipelined_request_for_the_next_loop() {
let mut manager = TestRemoteManager::new().await;
let state = test_runtime_state_with_manager(&manager);
let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
let address = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
serve_client(
stream,
test_metadata(address),
state,
CancellationToken::new(),
)
.await
});
let mut stream = TcpStream::connect(address).await.unwrap();
for (request_id, action) in [
(
1,
DaemonAction::ReviewerAction {
session_id: "session-1".into(),
role: None,
action: crate::session_manager::ReviewerAction::Pause,
},
),
(2, DaemonAction::Ping),
] {
write_frame(
&mut stream,
&RequestEnvelope {
protocol_version: PROTOCOL_VERSION,
request_id,
token: "right-token".into(),
action,
},
)
.await
.unwrap();
}
let request = tokio::time::timeout(Duration::from_secs(5), manager.requests.recv())
.await
.expect("reviewer request did not reach the manager")
.expect("manager request stream ended");
let RemoteSessionRequest::Reviewer { reply, .. } = request else {
panic!("expected a reviewer request")
};
reply
.send(Ok(crate::session_manager::ReviewerOutcome::Paused))
.expect("daemon still awaits the reviewer reply");
let first: ResponseEnvelope = read_frame(&mut stream).await.unwrap();
assert!(matches!(
first.result,
Ok(DaemonReply::Reviewer(outcome))
if matches!(*outcome, crate::session_manager::ReviewerOutcome::Paused)
));
let second: ResponseEnvelope = read_frame(&mut stream).await.unwrap();
assert!(matches!(second.result, Ok(DaemonReply::Pong)));
drop(stream);
let _ = server.await.expect("daemon task panicked");
}
#[tokio::test]
async fn daemon_client_eof_does_not_cancel_a_submitted_mutation() {
let mut manager = TestRemoteManager::new().await;
let state = test_runtime_state_with_manager(&manager);
let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
let address = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
serve_client(
stream,
test_metadata(address),
state,
CancellationToken::new(),
)
.await
});
let mut stream = TcpStream::connect(address).await.unwrap();
write_frame(
&mut stream,
&RequestEnvelope {
protocol_version: PROTOCOL_VERSION,
request_id: 1,
token: "right-token".into(),
action: DaemonAction::SubmitSessionCommand {
inherited_draft: None,
session_id: "session-1".into(),
command_id: "command-1".into(),
command: RelayCommand::ClearQueuedPrompts,
},
},
)
.await
.unwrap();
let request = tokio::time::timeout(Duration::from_secs(5), manager.requests.recv())
.await
.expect("submit request did not reach the manager")
.expect("manager request stream ended");
let RemoteSessionRequest::Submit { reply, .. } = request else {
panic!("expected a submitted mutation")
};
drop(stream);
reply
.send(Ok(7))
.expect("daemon incorrectly cancelled a non-reviewer mutation");
let _ = server.await.expect("daemon task panicked");
}
pub(super) fn runtime_test_session(
id: &str,
workspace_id: &str,
state: SessionState,
) -> SessionRecord {
SessionRecord {
target_runtime: None,
launch_base: None,
launch_branch: None,
checkout: None,
expected_runtime_identity: None,
publication: None,
build_cache: None,
container_workspace: None,
subagents: None,
create_managed_worktree: None,
id: id.into(),
workspace_id: workspace_id.into(),
title: id.into(),
harness_kind: mj_core::config::HarnessKind::Codex,
last_profile: "codex".into(),
bundle_id: "project".into(),
project_directory: None,
managed_worktree: None,
target_template_id: "local".into(),
resource_allocation: None,
additional_mounts: Vec::new(),
container_cpus: None,
container_memory: None,
state,
archived: false,
target: None,
native_session_id: None,
acp_session_title: None,
session_title_override: None,
created_at: "2026-09-03T00:00:00Z".into(),
updated_at: "2026-09-03T00:00:00Z".into(),
viewed_through_event_ordinal: 0,
draft_input: String::new(),
last_error: None,
last_checkpoint_error: None,
checkpoint: None,
}
}
#[test]
fn runtime_records_include_global_history_but_only_local_active_sessions() {
let local = runtime_test_session("local", "workspace-a", SessionState::Running);
let remote = runtime_test_session("remote", "workspace-b", SessionState::Running);
let history = runtime_test_session("history", "deleted-workspace", SessionState::Stopped);
let controller = Controller {
config: Config::default(),
state: mj_core::state::State {
sessions: [local, remote, history]
.into_iter()
.map(|session| (session.id.clone(), session))
.collect(),
..mj_core::state::State::default()
},
};
let records = runtime_records_for_workspace(&controller, &BTreeSet::from(["local".to_owned()]));
let ids = records
.iter()
.map(|session| session.id.as_str())
.collect::<BTreeSet<_>>();
assert_eq!(ids, BTreeSet::from(["history", "local"]));
}
pub(super) fn runtime_test_subagent(
child_session_id: &str,
parent_session_id: &str,
) -> SubagentRecord {
SubagentRecord {
child_session_id: child_session_id.into(),
parent_session_id: parent_session_id.into(),
task_name: "task".into(),
profile_id: "codex".into(),
model: None,
effort: None,
working_directory: Default::default(),
initial_prompt: "do the task".into(),
request_key: format!("request-{child_session_id}"),
created_at: "2026-09-03T00:00:00Z".into(),
noticed_turn: None,
handback_tool: false,
}
}
#[test]
fn runtime_subagents_include_only_relations_whose_child_is_in_the_returned_records() {
let parent_a = runtime_test_session("parent-a", "workspace-a", SessionState::Running);
let child_a1 = runtime_test_session("child-a1", "workspace-a", SessionState::Running);
let child_a2 = runtime_test_session("child-a2", "workspace-a", SessionState::Running);
let parent_b = runtime_test_session("parent-b", "workspace-b", SessionState::Running);
let child_b1 = runtime_test_session("child-b1", "workspace-b", SessionState::Running);
let mut state = mj_core::state::State {
sessions: [
parent_a.clone(),
child_a1.clone(),
child_a2.clone(),
parent_b.clone(),
child_b1.clone(),
]
.into_iter()
.map(|session| (session.id.clone(), session))
.collect(),
..mj_core::state::State::default()
};
state.subagents = [
runtime_test_subagent(&child_a1.id, &parent_a.id),
runtime_test_subagent(&child_a2.id, &parent_a.id),
runtime_test_subagent(&child_b1.id, &parent_b.id),
]
.into_iter()
.map(|subagent| (subagent.child_session_id.clone(), subagent))
.collect();
let controller = Controller {
config: Config::default(),
state,
};
let workspace_a_ids = BTreeSet::from([
"parent-a".to_owned(),
"child-a1".to_owned(),
"child-a2".to_owned(),
]);
let workspace_a_records = runtime_records_for_workspace(&controller, &workspace_a_ids);
let workspace_a_subagents = runtime_subagents_for_workspace(&controller, &workspace_a_records);
let mut workspace_a_child_ids = workspace_a_subagents
.iter()
.map(|subagent| subagent.child_session_id.as_str())
.collect::<Vec<_>>();
workspace_a_child_ids.sort_unstable();
assert_eq!(workspace_a_child_ids, ["child-a1", "child-a2"]);
let all_ids = BTreeSet::from([
"parent-a".to_owned(),
"child-a1".to_owned(),
"child-a2".to_owned(),
"parent-b".to_owned(),
"child-b1".to_owned(),
]);
let all_records = runtime_records_for_workspace(&controller, &all_ids);
let all_subagents = runtime_subagents_for_workspace(&controller, &all_records);
let mut all_child_ids = all_subagents
.iter()
.map(|subagent| subagent.child_session_id.as_str())
.collect::<Vec<_>>();
all_child_ids.sort_unstable();
assert_eq!(all_child_ids, ["child-a1", "child-a2", "child-b1"]);
}
#[test]
fn active_child_session_ids_skips_children_that_already_stopped() {
let parent = runtime_test_session("parent", "workspace", SessionState::Closing);
let running_child = runtime_test_session("running-child", "workspace", SessionState::Running);
let stopped_child = runtime_test_session("stopped-child", "workspace", SessionState::Stopped);
let unrelated = runtime_test_session("unrelated", "workspace", SessionState::Running);
let mut state = mj_core::state::State {
sessions: [
parent.clone(),
running_child.clone(),
stopped_child.clone(),
unrelated.clone(),
]
.into_iter()
.map(|session| (session.id.clone(), session))
.collect(),
..mj_core::state::State::default()
};
state.subagents = [
runtime_test_subagent(&running_child.id, &parent.id),
runtime_test_subagent(&stopped_child.id, &parent.id),
]
.into_iter()
.map(|subagent| (subagent.child_session_id.clone(), subagent))
.collect();
let children = active_child_session_ids(&state, &parent.id);
assert_eq!(children, vec!["running-child".to_owned()]);
}
#[tokio::test]
async fn daemon_records_definitive_missing_workspaces_without_an_attached_surface() {
let state = test_runtime_state();
let session = runtime_test_session("missing", "workspace", SessionState::Running);
state.owner().edit_sessions(|sessions| {
sessions.insert(session.id.clone(), session.clone());
});
assert!(state.attachments.lock().unwrap().is_empty());
let mut view = ManagedSessionView {
snapshot: None,
connected: false,
error: Some(ViewError::Unreachable(
"relay proxy disconnected during hello".into(),
)),
};
assert!(state.missing_target_record(&session.id, &view).is_none());
view.error = Some(ViewError::TargetMissing(
"working directory /missing is gone".into(),
));
assert_eq!(
state.missing_target_record(&session.id, &view),
Some((
"working directory /missing is gone".into(),
session.updated_at
)),
);
state.owner().edit_sessions(|sessions| {
sessions.get_mut(&session.id).unwrap().state = SessionState::Closing;
});
assert!(state.missing_target_record(&session.id, &view).is_none());
state.owner().edit_sessions(|sessions| {
sessions.get_mut(&session.id).unwrap().state = SessionState::Error;
});
assert!(state.missing_target_record(&session.id, &view).is_none());
}
#[tokio::test]
async fn review_host_notifier_wakes_runtime_revision_subscribers() {
let revisions = RuntimeRevisions::new(40);
let mut subscriber = revisions.subscribe();
let notify_review_publication = revisions.notifier();
notify_review_publication();
tokio::time::timeout(Duration::from_secs(1), subscriber.changed())
.await
.expect("review publication did not wake runtime subscribers")
.expect("runtime revision publisher stopped");
assert_eq!(*subscriber.borrow_and_update(), 41);
}
#[test]
fn late_runtime_revision_publication_cannot_move_cursor_backwards() {
let revisions = RuntimeRevisions::new(40);
let subscriber = revisions.subscribe();
revisions.publish_allocated(42);
revisions.publish_allocated(41);
assert_eq!(*subscriber.borrow(), 42);
}
#[tokio::test]
async fn workspace_publication_reaches_existing_phone_subscriber() {
let state = test_runtime_state();
let mut workspaces = state.workspaces();
let expected = WorkspaceRecord {
id: "workspace-1".into(),
name: "Reliability".into(),
created_at: "2026-08-30T00:00:00Z".into(),
last_opened_at: "2026-08-30T00:00:00Z".into(),
session_count: 0,
};
state.publish_workspaces(vec![expected.clone()]);
tokio::time::timeout(Duration::from_secs(1), workspaces.changed())
.await
.expect("workspace publication timed out")
.expect("workspace publisher stopped");
assert_eq!(workspaces.borrow_and_update().as_slice(), &[expected]);
}
#[tokio::test]
async fn framing_round_trips_payloads_larger_than_a_pipe_buffer() {
let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
let address = listener.local_addr().unwrap();
let sender = tokio::spawn(async move {
let mut stream = TcpStream::connect(address).await.unwrap();
write_frame(&mut stream, &"x".repeat(512 * 1024))
.await
.unwrap();
});
let (mut stream, _) = listener.accept().await.unwrap();
let received: String = read_frame(&mut stream).await.unwrap();
sender.await.unwrap();
assert_eq!(received.len(), 512 * 1024);
}
#[tokio::test]
async fn framing_rejects_an_oversized_frame_before_allocating_it() {
let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
let address = listener.local_addr().unwrap();
let sender = tokio::spawn(async move {
let mut stream = TcpStream::connect(address).await.unwrap();
stream
.write_u32((MAX_FRAME_BYTES + 1) as u32)
.await
.unwrap();
});
let (mut stream, _) = listener.accept().await.unwrap();
assert!(read_frame::<String>(&mut stream).await.is_err());
sender.await.unwrap();
}
#[tokio::test]
async fn daemon_stops_serving_data_actions_once_shutdown_begins() {
let state = test_runtime_state();
let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
let address = listener.local_addr().unwrap();
let metadata = DaemonMetadata {
protocol_version: PROTOCOL_VERSION,
pid: 1,
address,
token: "right-token".into(),
started_at: "now".into(),
build_version: "test".into(),
};
let cancellation = CancellationToken::new();
cancellation.cancel();
let server_metadata = metadata.clone();
let server_cancellation = cancellation.clone();
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
serve_client(stream, server_metadata, state, server_cancellation)
.await
.unwrap();
});
let mut stream = TcpStream::connect(address).await.unwrap();
let mut request_id = 0;
let mut ask = async |stream: &mut TcpStream, action: DaemonAction| {
request_id += 1;
write_frame(
stream,
&RequestEnvelope {
protocol_version: PROTOCOL_VERSION,
request_id,
token: "right-token".to_owned(),
action,
},
)
.await
.unwrap();
read_frame::<ResponseEnvelope>(stream).await.unwrap().result
};
let refused = ask(
&mut stream,
DaemonAction::Snapshot {
workspace_id: "workspace-a".into(),
},
)
.await;
assert_eq!(
refused.unwrap_err(),
"daemon is shutting down; retry to reach a fresh daemon"
);
assert!(matches!(
ask(&mut stream, DaemonAction::Ping).await,
Ok(DaemonReply::Pong)
));
assert!(matches!(
ask(&mut stream, DaemonAction::Status).await,
Ok(DaemonReply::Status(_))
));
assert!(matches!(
ask(&mut stream, DaemonAction::Stop).await,
Ok(DaemonReply::Done)
));
drop(stream);
server.await.unwrap();
}
#[test]
fn shutdown_force_exit_finishes_before_the_stop_deadline() {
assert!(SHUTDOWN_FORCE_EXIT_TIMEOUT < STOP_TIMEOUT);
}
#[tokio::test]
async fn management_stop_is_bounded_when_the_daemon_never_acknowledges() {
let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
let metadata = DaemonMetadata {
protocol_version: PROTOCOL_VERSION,
pid: 1,
address: listener.local_addr().unwrap(),
token: "test-token".into(),
started_at: "now".into(),
build_version: "test".into(),
};
let mut client = ManagementClient::new(DaemonClient::connect(metadata).await.unwrap());
let (mut peer, _) = listener.accept().await.unwrap();
let stop = tokio::spawn(async move { client.stop().await });
let request: RequestEnvelope = read_frame(&mut peer).await.unwrap();
assert!(matches!(request.action, DaemonAction::Stop));
tokio::time::pause();
tokio::time::advance(STOP_TIMEOUT).await;
let error = stop.await.unwrap().unwrap_err();
assert!(error.to_string().contains("did not acknowledge"));
}
#[tokio::test]
async fn daemon_rejects_a_request_with_the_wrong_owner_token() {
let state = test_runtime_state();
let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
let address = listener.local_addr().unwrap();
let metadata = DaemonMetadata {
protocol_version: PROTOCOL_VERSION,
pid: 1,
address,
token: "right-token".into(),
started_at: "now".into(),
build_version: "test".into(),
};
let server_metadata = metadata.clone();
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
serve_client(stream, server_metadata, state, CancellationToken::new())
.await
.unwrap();
});
let mut stream = TcpStream::connect(address).await.unwrap();
write_frame(
&mut stream,
&RequestEnvelope {
protocol_version: PROTOCOL_VERSION,
request_id: 42,
token: "wrong-token".into(),
action: DaemonAction::Ping,
},
)
.await
.unwrap();
let response: ResponseEnvelope = read_frame(&mut stream).await.unwrap();
assert_eq!(response.request_id, 42);
assert_eq!(response.result.unwrap_err(), "daemon authentication failed");
drop(stream);
server.await.unwrap();
}
#[tokio::test]
async fn the_daemon_rejects_a_client_one_protocol_behind_before_dispatch() {
let state = test_runtime_state();
let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
let address = listener.local_addr().unwrap();
let metadata = DaemonMetadata {
protocol_version: PROTOCOL_VERSION,
pid: 1,
address,
token: "right-token".into(),
started_at: "now".into(),
build_version: "test".into(),
};
let server_metadata = metadata.clone();
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
serve_client(stream, server_metadata, state, CancellationToken::new())
.await
.unwrap();
});
let mut stream = TcpStream::connect(address).await.unwrap();
write_frame(
&mut stream,
&RequestEnvelope {
protocol_version: PROTOCOL_VERSION - 1,
request_id: 43,
token: metadata.token,
action: DaemonAction::PersistReadReceipt {
client_id: "client-a".into(),
workspace_id: "workspace-a".into(),
session_id: "session-a".into(),
through: 7,
},
},
)
.await
.unwrap();
let response: ResponseEnvelope = read_frame(&mut stream).await.unwrap();
assert_eq!(response.request_id, 43);
assert!(response.result.unwrap_err().contains(&format!(
"incompatible daemon protocol {}; expected {PROTOCOL_VERSION}",
PROTOCOL_VERSION - 1
)));
drop(stream);
server.await.unwrap();
}
#[test]
fn management_wire_shapes_stay_frozen_across_protocol_versions() {
for (action, expected) in [
(DaemonAction::Ping, serde_json::json!({"action": "ping"})),
(
DaemonAction::Status,
serde_json::json!({"action": "status"}),
),
(DaemonAction::Stop, serde_json::json!({"action": "stop"})),
] {
let request = RequestEnvelope {
protocol_version: 3,
request_id: 7,
token: "tok".into(),
action,
};
assert_eq!(
serde_json::to_value(&request).unwrap(),
serde_json::json!({
"protocol_version": 3,
"request_id": 7,
"token": "tok",
"action": expected,
})
);
}
let response: ResponseEnvelope = serde_json::from_value(serde_json::json!({
"protocol_version": 3,
"request_id": 7,
"result": {"Ok": {"reply": "status", "value": {
"pid": 4242,
"started_at": "2026-09-01T07:48:14Z",
"build_version": "0.3.1",
"attached_clients": 1,
"phone_status": {"state": "disabled"},
}}}
}))
.unwrap();
match response.result.unwrap() {
DaemonReply::Status(status) => {
assert_eq!(status.pid, 4242);
assert_eq!(status.build_version, "0.3.1");
}
reply => panic!("unexpected reply {reply:?}"),
}
}
#[test]
fn client_presence_and_workspace_listing_use_global_wire_shapes() {
let attach = serde_json::to_value(DaemonAction::Attach {
client_id: "client-a".into(),
pid: 4242,
})
.unwrap();
assert_eq!(
attach,
serde_json::json!({
"action": "attach",
"arguments": {"client_id": "client-a", "pid": 4242},
})
);
let listing = serde_json::to_value(WorkspaceListing {
workspace: WorkspaceRecord {
id: "workspace-a".into(),
name: "Workspace A".into(),
created_at: "2026-09-01T00:00:00Z".into(),
last_opened_at: "2026-09-01T00:00:00Z".into(),
session_count: 0,
},
})
.unwrap();
assert!(listing.get("attached_pids").is_none());
}
#[tokio::test]
async fn daemon_serves_management_actions_for_any_protocol_version() {
for version in [3_u32, 5] {
let state = test_runtime_state();
let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
let address = listener.local_addr().unwrap();
let metadata = DaemonMetadata {
protocol_version: PROTOCOL_VERSION,
pid: 1,
address,
token: "right-token".into(),
started_at: "now".into(),
build_version: "test".into(),
};
let server_metadata = metadata.clone();
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
serve_client(stream, server_metadata, state, CancellationToken::new())
.await
.unwrap();
});
let mut stream = TcpStream::connect(address).await.unwrap();
write_frame(
&mut stream,
&RequestEnvelope {
protocol_version: version,
request_id: 44,
token: "right-token".into(),
action: DaemonAction::Status,
},
)
.await
.unwrap();
let response: ResponseEnvelope = read_frame(&mut stream).await.unwrap();
assert_eq!(
response.protocol_version, version,
"reply must use the caller's dialect"
);
assert_eq!(response.request_id, 44);
match response.result.unwrap() {
DaemonReply::Status(status) => assert_eq!(status.build_version, "test"),
reply => panic!("unexpected reply {reply:?}"),
}
drop(stream);
server.await.unwrap();
}
}
struct ProtocolTranscript {
protocol_version: u32,
daemon_build: &'static str,
expected_requests: [serde_json::Value; 2],
responses: [&'static str; 2],
}
fn released_protocol_transcripts() -> Vec<ProtocolTranscript> {
let requests = |version: u32| {
[
serde_json::json!({
"protocol_version": version,
"request_id": 1,
"token": "tok",
"action": {"action": "status"},
}),
serde_json::json!({
"protocol_version": version,
"request_id": 2,
"token": "tok",
"action": {"action": "stop"},
}),
]
};
vec![
ProtocolTranscript {
protocol_version: 3,
daemon_build: "0.3.1",
expected_requests: requests(3),
responses: [
r#"{"protocol_version":3,"request_id":1,"result":{"Ok":{"reply":"status","value":{"pid":4242,"started_at":"2026-09-01T07:48:14Z","build_version":"0.3.1","attached_clients":1,"phone_status":{"state":"ready","viewer_url":"https://example.test:1","viewer_code":"690451","qr_login_url":null,"fallback_reason":null}}}}}"#,
r#"{"protocol_version":3,"request_id":2,"result":{"Ok":{"reply":"done"}}}"#,
],
},
ProtocolTranscript {
protocol_version: 4,
daemon_build: "0.4.1",
expected_requests: requests(4),
responses: [
r#"{"protocol_version":4,"request_id":1,"result":{"Ok":{"reply":"status","value":{"pid":4242,"started_at":"2026-09-01T07:48:14Z","build_version":"0.4.1","attached_clients":1,"phone_status":{"state":"disabled"}}}}}"#,
r#"{"protocol_version":4,"request_id":2,"result":{"Ok":{"reply":"done"}}}"#,
],
},
ProtocolTranscript {
protocol_version: 11,
daemon_build: "2.1.0",
expected_requests: requests(11),
responses: [
r#"{"protocol_version":11,"request_id":1,"result":{"Ok":{"reply":"status","value":{"pid":4242,"started_at":"2026-09-01T07:48:14Z","build_version":"2.1.0","attached_clients":1,"phone_status":{"state":"disabled"}}}}}"#,
r#"{"protocol_version":11,"request_id":2,"result":{"Ok":{"reply":"done"}}}"#,
],
},
ProtocolTranscript {
protocol_version: 14,
daemon_build: "2.1.4",
expected_requests: requests(14),
responses: [
r#"{"protocol_version":14,"request_id":1,"result":{"Ok":{"reply":"status","value":{"pid":4242,"started_at":"2026-09-01T07:48:14Z","build_version":"2.1.4","attached_clients":1,"phone_status":{"state":"disabled"}}}}}"#,
r#"{"protocol_version":14,"request_id":2,"result":{"Ok":{"reply":"done"}}}"#,
],
},
ProtocolTranscript {
protocol_version: 15,
daemon_build: "2.2.0",
expected_requests: requests(15),
responses: [
r#"{"protocol_version":15,"request_id":1,"result":{"Ok":{"reply":"status","value":{"pid":4242,"started_at":"2026-09-01T07:48:14Z","build_version":"2.2.0","attached_clients":1,"phone_status":{"state":"disabled"}}}}}"#,
r#"{"protocol_version":15,"request_id":2,"result":{"Ok":{"reply":"done"}}}"#,
],
},
ProtocolTranscript {
protocol_version: 16,
daemon_build: "2.4.0",
expected_requests: requests(16),
responses: [
r#"{"protocol_version":16,"request_id":1,"result":{"Ok":{"reply":"status","value":{"pid":4242,"started_at":"2026-09-01T07:48:14Z","build_version":"2.4.0","attached_clients":1,"phone_status":{"state":"disabled"}}}}}"#,
r#"{"protocol_version":16,"request_id":2,"result":{"Ok":{"reply":"done"}}}"#,
],
},
]
}
#[tokio::test]
async fn management_client_talks_to_every_released_protocol_version() {
for transcript in released_protocol_transcripts() {
let protocol_version = transcript.protocol_version;
let daemon_build = transcript.daemon_build;
let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
let address = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut stream, _) = listener.accept().await.unwrap();
for (expected, response) in transcript
.expected_requests
.iter()
.zip(transcript.responses)
{
let request: serde_json::Value = read_frame(&mut stream).await.unwrap();
assert_eq!(
&request, expected,
"protocol {} daemon would reject this frame",
transcript.protocol_version
);
stream.write_u32(response.len() as u32).await.unwrap();
stream.write_all(response.as_bytes()).await.unwrap();
stream.flush().await.unwrap();
}
});
let metadata = DaemonMetadata {
protocol_version,
pid: 4242,
address,
token: "tok".into(),
started_at: "2026-09-01T07:48:14Z".into(),
build_version: daemon_build.into(),
};
let mut client = ManagementClient::new(DaemonClient::connect(metadata).await.unwrap());
let status = client.status().await.unwrap();
assert_eq!(status.build_version, daemon_build);
assert_eq!(status.attached_clients, 1);
assert_eq!(client.protocol_version(), protocol_version);
client.stop().await.unwrap();
server.await.unwrap();
}
}
#[tokio::test]
async fn equivalent_lifecycle_requests_join_one_daemon_operation() {
let state = test_runtime_state();
let starts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let release = Arc::new(tokio::sync::Notify::new());
let first = state
.start_or_join_lifecycle("session-1".into(), LifecycleKind::Suspend, {
let starts = starts.clone();
let release = release.clone();
move |_state, _session_id, _cancelled| async move {
starts.fetch_add(1, Ordering::AcqRel);
release.notified().await;
Ok(DaemonLifecycleResult::Done)
}
})
.unwrap();
tokio::task::yield_now().await;
let second = state
.start_or_join_lifecycle(
"session-1".into(),
LifecycleKind::Suspend,
|_state, _session_id, _cancelled| async move {
panic!("joined lifecycle request started duplicate work")
},
)
.unwrap();
assert_eq!(starts.load(Ordering::Acquire), 1);
assert!(
state
.start_or_join_lifecycle(
"session-1".into(),
LifecycleKind::Resume,
|_state, _session_id, _cancelled| async move { Ok(DaemonLifecycleResult::Done) },
)
.is_err()
);
drop(first);
release.notify_one();
assert!(matches!(
RuntimeState::wait_lifecycle_result(second).await.unwrap(),
DaemonLifecycleResult::Done
));
assert_eq!(starts.load(Ordering::Acquire), 1);
}
#[tokio::test]
async fn close_keeps_worker_target_available_for_checkpoint_lease() {
let state = test_runtime_state();
let release = Arc::new(tokio::sync::Notify::new());
let result = state
.start_or_join_lifecycle("session-1".into(), LifecycleKind::Suspend, {
let release = release.clone();
move |_state, _session_id, _cancelled| async move {
release.notified().await;
Ok(DaemonLifecycleResult::Done)
}
})
.unwrap();
assert!(state.worker_poll_exclusion_session_ids().is_empty());
release.notify_one();
RuntimeState::wait_lifecycle_result(result).await.unwrap();
}
#[tokio::test]
async fn move_joins_only_matching_selections_and_runs_other_sessions_concurrently() {
let state = test_runtime_state();
let release = Arc::new(tokio::sync::Notify::new());
let first = state
.start_or_join_lifecycle_with_key(
"move-join-one".into(),
LifecycleKind::Move,
None,
Some("profile-a/target-a/discard".into()),
{
let release = release.clone();
move |_, _, _| async move {
release.notified().await;
Ok(DaemonLifecycleResult::Done)
}
},
)
.unwrap();
assert!(crate::controller::move_session::move_owns_session(
"move-join-one"
));
let duplicate = state
.start_or_join_lifecycle_with_key(
"move-join-one".into(),
LifecycleKind::Move,
None,
Some("profile-a/target-a/discard".into()),
|_, _, _| async move { panic!("duplicate move launched a second writer") },
)
.unwrap();
assert!(
state
.start_or_join_lifecycle_with_key(
"move-join-one".into(),
LifecycleKind::Move,
None,
Some("profile-b/target-a/discard".into()),
|_, _, _| async move { Ok(DaemonLifecycleResult::Done) },
)
.is_err()
);
let unrelated = state
.start_or_join_lifecycle_with_key(
"move-join-two".into(),
LifecycleKind::Move,
None,
Some("profile-b/target-b/start".into()),
|_, _, _| async move { Ok(DaemonLifecycleResult::Done) },
)
.unwrap();
let unrelated_channel = unrelated.clone();
tokio::time::timeout(
Duration::from_secs(2),
RuntimeState::wait_lifecycle_result(unrelated),
)
.await
.unwrap()
.unwrap();
state.remove_completed_lifecycle(&unrelated_channel);
assert!(!crate::controller::move_session::move_owns_session(
"move-join-two"
));
drop(first); release.notify_one();
let channel = duplicate.clone();
RuntimeState::wait_lifecycle_result(duplicate)
.await
.unwrap();
state.remove_completed_lifecycle(&channel);
assert!(!crate::controller::move_session::move_owns_session(
"move-join-one"
));
}
#[tokio::test]
async fn move_task_panic_reports_failure_and_releases_mutation_hold() {
let state = test_runtime_state();
let result = state
.start_or_join_lifecycle_with_key(
"move-panics".into(),
LifecycleKind::Move,
None,
Some("destination".into()),
|_, _, _| async move { panic!("injected move task panic") },
)
.unwrap();
let channel = result.clone();
let error = tokio::time::timeout(
Duration::from_secs(2),
RuntimeState::wait_lifecycle_result(result),
)
.await
.unwrap()
.unwrap_err();
assert!(error.to_string().contains("daemon lifecycle task failed"));
state.remove_completed_lifecycle(&channel);
assert!(!crate::controller::move_session::move_owns_session(
"move-panics"
));
}
#[tokio::test]
async fn daemon_lifecycle_reports_balanced_concurrent_stages() {
let state = test_runtime_state();
let release = Arc::new(tokio::sync::Notify::new());
let result = state
.start_or_join_lifecycle("session-1".into(), LifecycleKind::Create, {
let release = release.clone();
move |_state, _session_id, _cancelled| async move {
release.notified().await;
Ok(DaemonLifecycleResult::Done)
}
})
.unwrap();
assert_eq!(
state.worker_poll_exclusion_session_ids(),
BTreeSet::from(["session-1".to_owned()])
);
let executor = DaemonStageReportingExecutor::new(
RefusingExecutor("a stage notification"),
state.clone(),
"session-1".into(),
);
executor.stage_started(ProvisionStage::Cloning);
executor.stage_started(ProvisionStage::Cloning);
executor.stage_started(ProvisionStage::Syncing);
executor.stage_finished(ProvisionStage::Cloning);
{
let lifecycle_owner = state.owner();
let lifecycle = &lifecycle_owner.lifecycle;
let stages = &lifecycle.get("session-1").unwrap().active_stages;
assert_eq!(stages.get(&ProvisionStage::Cloning).unwrap().0, 1);
assert_eq!(stages.get(&ProvisionStage::Syncing).unwrap().0, 1);
}
executor.stage_finished(ProvisionStage::Cloning);
executor.stage_finished(ProvisionStage::Syncing);
assert!(
state
.owner()
.lifecycle
.get("session-1")
.unwrap()
.active_stages
.is_empty()
);
release.notify_one();
assert!(matches!(
RuntimeState::wait_lifecycle_result(result).await.unwrap(),
DaemonLifecycleResult::Done
));
}
#[tokio::test]
async fn deferred_cleanup_is_visible_and_drains_before_shutdown_cancellation() {
let state = test_runtime_state();
let saw_early_cancellation = Arc::new(AtomicBool::new(false));
let result = state
.start_or_join_lifecycle("session-1".into(), LifecycleKind::Cleanup, {
let saw_early_cancellation = saw_early_cancellation.clone();
move |state, session_id, cancelled| async move {
let executor = DaemonStageReportingExecutor::new(
crate::targets::ProcessExecutor,
state,
session_id,
);
executor.stage_started(ProvisionStage::RemovingStorage);
tokio::time::sleep(Duration::from_millis(20)).await;
saw_early_cancellation.store(cancelled.load(Ordering::Acquire), Ordering::Release);
executor.stage_finished(ProvisionStage::RemovingStorage);
Ok(DaemonLifecycleResult::Done)
}
})
.unwrap();
tokio::task::yield_now().await;
let visible = state.active_lifecycles();
assert_eq!(visible.len(), 1);
assert_eq!(visible[0].kind, RuntimeLifecycleKind::Cleanup);
assert_eq!(
visible[0].active_stages[0].0,
ProvisionStage::RemovingStorage
);
state.cancel_and_wait_lifecycles().await.unwrap();
assert!(!saw_early_cancellation.load(Ordering::Acquire));
assert!(matches!(
RuntimeState::wait_lifecycle_result(result).await.unwrap(),
DaemonLifecycleResult::Done
));
}
#[test]
fn force_destroy_serializes_as_its_own_lifecycle_kind() {
assert_eq!(
serde_json::to_string(&RuntimeLifecycleKind::ForceDestroy).unwrap(),
"\"force_destroy\""
);
}
#[test]
fn create_cancellation_and_commit_have_one_winner() {
for _ in 0..32 {
let control = CreateSessionControl::default();
let barrier = Arc::new(std::sync::Barrier::new(2));
let canceller = {
let control = control.clone();
let barrier = barrier.clone();
std::thread::spawn(move || {
barrier.wait();
control.request_cancel()
})
};
barrier.wait();
let committed = control.grant_commit();
let cancelled = canceller.join().expect("canceller panicked");
assert_ne!(committed, cancelled);
assert_eq!(control.cancelled.load(Ordering::Acquire), cancelled);
assert!(!control.is_cancellable());
assert!(!control.request_cancel());
assert!(!control.grant_commit());
}
}
#[tokio::test]
async fn completed_stop_stays_visible_until_cleanup_takes_ownership() {
let state = test_runtime_state();
let (complete, result) = tokio::sync::watch::channel(None);
state.owner().lifecycle.insert(
"cleanup-gap".into(),
ActiveLifecycle {
phase: LifecyclePhase::Executing,
operation_id: "closing-operation".into(),
create_control: None,
kind: LifecycleKind::Suspend,
cancelled: Arc::new(AtomicBool::new(false)),
started_at_epoch_seconds: 1,
active_stages: BTreeMap::new(),
resume_workspace_id: None,
resume_destination: None,
notice: None,
request_key: None,
_move_guard: None,
result: result.clone(),
},
);
state
.owner()
.lifecycle
.get_mut("cleanup-gap")
.unwrap()
.phase = LifecyclePhase::Completed(Ok(DaemonLifecycleResult::DeferredCleanup));
complete.send_replace(Some(Ok(DaemonLifecycleResult::DeferredCleanup)));
state.remove_completed_lifecycle(&result);
let view = state.active_lifecycles();
assert_eq!(view.len(), 1);
assert_eq!(view[0].operation_id, "closing-operation");
assert!(!view[0].cancellable);
let release = Arc::new(tokio::sync::Notify::new());
let cleanup = state
.start_or_join_lifecycle("cleanup-gap".into(), LifecycleKind::Cleanup, {
let release = release.clone();
move |_, _, _| async move {
release.notified().await;
Ok(DaemonLifecycleResult::Done)
}
})
.unwrap();
state.remove_completed_lifecycle(&result);
let view = state.active_lifecycles();
assert_eq!(view.len(), 1);
assert_eq!(view[0].kind, RuntimeLifecycleKind::Cleanup);
assert_ne!(view[0].operation_id, "closing-operation");
release.notify_one();
RuntimeState::wait_lifecycle_result(cleanup).await.unwrap();
assert!(state.active_lifecycles().is_empty());
}
#[tokio::test]
async fn session_projection_reads_lifecycles_without_relocking_the_controller() {
let state = test_runtime_state();
let release = Arc::new(tokio::sync::Notify::new());
let running = state
.start_or_join_lifecycle("session-1".into(), LifecycleKind::Suspend, {
let release = release.clone();
move |_, _, _| async move {
release.notified().await;
Ok(DaemonLifecycleResult::Done)
}
})
.unwrap();
let (sender, receiver) = std::sync::mpsc::channel();
let projecting = state.clone();
std::thread::spawn(move || {
let (_, lifecycles) = projecting.session_projection();
let _ = sender.send(lifecycles.len());
});
let visible = receiver
.recv_timeout(Duration::from_secs(5))
.expect("session projection must not deadlock on the controller lock");
assert_eq!(visible, 1);
release.notify_one();
RuntimeState::wait_lifecycle_result(running).await.unwrap();
}
#[tokio::test]
async fn lifecycle_identity_survives_join_but_changes_for_next_operation() {
let state = test_runtime_state();
let release = Arc::new(tokio::sync::Notify::new());
let first = state
.start_or_join_lifecycle("identity".into(), LifecycleKind::Resume, {
let release = release.clone();
move |_, _, _| async move {
release.notified().await;
Ok(DaemonLifecycleResult::Done)
}
})
.unwrap();
let first_id = state.active_lifecycles()[0].operation_id.clone();
let joined = state
.start_or_join_lifecycle("identity".into(), LifecycleKind::Resume, |_, _, _| async {
panic!("joined operation must not run twice")
})
.unwrap();
assert_eq!(state.active_lifecycles()[0].operation_id, first_id);
release.notify_one();
RuntimeState::wait_lifecycle_result(first.clone())
.await
.unwrap();
RuntimeState::wait_lifecycle_result(joined).await.unwrap();
state.remove_completed_lifecycle(&first);
let second = state
.start_or_join_lifecycle("identity".into(), LifecycleKind::Resume, {
let release = release.clone();
move |_, _, _| async move {
release.notified().await;
Ok(DaemonLifecycleResult::Done)
}
})
.unwrap();
let second_id = state.active_lifecycles()[0].operation_id.clone();
assert_ne!(second_id, first_id);
state.remove_completed_lifecycle(&first);
assert_eq!(state.active_lifecycles()[0].operation_id, second_id);
release.notify_one();
RuntimeState::wait_lifecycle_result(second).await.unwrap();
}
#[tokio::test]
async fn close_waits_for_cancelled_or_committed_provisioning_to_release_ownership() {
for committed in [false, true] {
let state = test_runtime_state();
let control = CreateSessionControl::default();
let release = Arc::new(tokio::sync::Notify::new());
state
.start_or_join_lifecycle_controlled(
"close-race".into(),
LifecycleKind::Create,
None,
None,
Some(control.clone()),
{
let release = release.clone();
move |_, _, _| async move {
release.notified().await;
Ok(DaemonLifecycleResult::Done)
}
},
)
.unwrap();
if committed {
assert!(control.grant_commit());
}
state.request_close("close-race");
let waiter = {
let state = state.clone();
tokio::spawn(async move { state.wait_before_close("close-race").await })
};
tokio::task::yield_now().await;
assert!(
!waiter.is_finished(),
"cleanup must wait for the owning operation"
);
assert_eq!(control.cancelled.load(Ordering::Acquire), !committed);
assert_eq!(
state.session_state("close-race"),
Some(SessionState::Closing)
);
release.notify_one();
tokio::time::timeout(Duration::from_secs(2), waiter)
.await
.unwrap()
.unwrap()
.unwrap();
assert!(!state.owner().lifecycle.contains_key("close-race"));
}
}
#[tokio::test]
async fn committed_creation_cannot_be_cancelled_by_another_surface() {
let state = test_runtime_state();
let control = CreateSessionControl::default();
let release = Arc::new(tokio::sync::Notify::new());
let result = state
.start_or_join_lifecycle_controlled(
"committed".into(),
LifecycleKind::Create,
None,
None,
Some(control.clone()),
{
let release = release.clone();
move |_, _, _| async move {
release.notified().await;
Ok(DaemonLifecycleResult::Done)
}
},
)
.unwrap();
assert!(state.active_lifecycles()[0].cancellable);
assert!(control.grant_commit());
assert!(!state.active_lifecycles()[0].cancellable);
assert!(state.cancel_lifecycle("committed").is_err());
state.cancel_lifecycle_if_active("committed");
assert!(!control.cancelled.load(Ordering::Acquire));
release.notify_one();
RuntimeState::wait_lifecycle_result(result).await.unwrap();
}
#[tokio::test]
async fn a_close_removing_the_target_stops_offering_cancellation() {
let state = test_runtime_state();
let mut session = runtime_test_session("destroying", "workspace", SessionState::Closing);
session.target = Some(mj_core::state::TargetLocator::LocalPodman {
borrowed_from: None,
container_id: "a".repeat(64),
workspace_storage: Default::default(),
});
state.owner().edit_sessions(|sessions| {
sessions.insert(session.id.clone(), session.clone());
});
let release = Arc::new(tokio::sync::Notify::new());
let result = state
.start_or_join_lifecycle("destroying".into(), LifecycleKind::Suspend, {
let release = release.clone();
move |_state, _session_id, _cancelled| async move {
release.notified().await;
Ok(DaemonLifecycleResult::Done)
}
})
.unwrap();
assert!(state.active_lifecycles()[0].cancellable);
session.state = SessionState::Destroying;
state.owner().edit_sessions(|sessions| {
sessions.insert(session.id.clone(), session);
});
assert!(!state.active_lifecycles()[0].cancellable);
let error = state.cancel_lifecycle("destroying").unwrap_err();
assert!(
error.to_string().contains("cannot be cancelled"),
"unexpected error: {error:#}"
);
assert!(
!state
.owner()
.lifecycle
.get("destroying")
.expect("lifecycle entry")
.cancelled
.load(Ordering::Acquire),
"a refused cancel must not reach the running teardown"
);
release.notify_one();
RuntimeState::wait_lifecycle_result(result).await.unwrap();
}
#[tokio::test]
async fn force_destruction_preempts_a_running_lifecycle_and_waits_for_it() {
let state = test_runtime_state();
let release = Arc::new(tokio::sync::Notify::new());
state
.start_or_join_lifecycle("session-1".into(), LifecycleKind::Create, {
let release = release.clone();
move |_state, _session_id, _cancelled| async move {
release.notified().await;
Ok(DaemonLifecycleResult::Done)
}
})
.unwrap();
tokio::task::yield_now().await;
let preempt_state = state.clone();
let preempted =
tokio::spawn(async move { preempt_state.preempt_active_lifecycle("session-1").await });
tokio::task::yield_now().await;
{
let lifecycle_owner = state.owner();
let lifecycle = &lifecycle_owner.lifecycle;
assert!(
lifecycle
.get("session-1")
.expect("lifecycle entry")
.cancelled
.load(Ordering::Acquire),
"preemption must cancel the running operation"
);
}
release.notify_one();
preempted
.await
.expect("preempt task")
.expect("a cancelled-and-finished lifecycle lets force destruction proceed");
}
#[tokio::test(start_paused = true)]
async fn force_destruction_preemption_times_out_without_destroying() {
let state = test_runtime_state();
state
.start_or_join_lifecycle("session-1".into(), LifecycleKind::Create, {
|_state, _session_id, _cancelled| async move {
std::future::pending::<()>().await;
#[allow(unreachable_code)]
Ok(DaemonLifecycleResult::Done)
}
})
.unwrap();
tokio::task::yield_now().await;
let error = state
.preempt_active_lifecycle("session-1")
.await
.unwrap_err();
assert!(
error
.to_string()
.contains("did not stop after cancellation"),
"{error:#}"
);
}
#[test]
fn a_daemon_owned_notice_reaches_every_workspace_snapshot() {
let session_ids = BTreeSet::from(["018f9dd2-a3b4".to_owned()]);
let daemon_notice = RuntimeNotice {
id: 1,
session_id: String::new(),
text: "Downloading image ghcr.io/example/dev:latest for local podman\u{2026}".to_owned(),
};
let own_session = RuntimeNotice {
id: 2,
session_id: "018f9dd2-a3b4".to_owned(),
text: "Mounted /data read-only.".to_owned(),
};
let other_session = RuntimeNotice {
id: 3,
session_id: "018f9dd2-cccc".to_owned(),
text: "Mounted /data read-only.".to_owned(),
};
assert!(snapshot::notice_reaches_workspace(
&daemon_notice,
&session_ids
));
assert!(snapshot::notice_reaches_workspace(
&own_session,
&session_ids
));
assert!(!snapshot::notice_reaches_workspace(
&other_session,
&session_ids
));
}
fn ready_startup_view() -> ManagedSessionView {
let materialized = mj_core::state::MaterializedSession::empty("session-1");
let operational = mj_core::relay::RelayOperationalState {
assessment: None,
assessment_context: None,
turn_completion: None,
continuation: Default::default(),
relay_protocol_version: Some(mj_core::relay::RELAY_PROTOCOL_VERSION),
native_agents: Vec::new(),
steering: None,
cancelling_prompt_id: None,
clear_context: false,
clear_context_started_at_ms: None,
native_agent_count: 0,
expected_continuation: None,
inferred_idle_since_ms: None,
goal: Default::default(),
capacity_retry: None,
retry_assessment_pending: false,
activity_turn_started_at_ms: None,
idle_since_ms: None,
store_id: None,
session_id: "session-1".into(),
execution: mj_core::relay::RelayExecutionState::Idle,
latest_ordinal: 0,
latest_digest: mj_core::relay::RELAY_EVENT_GENESIS_DIGEST.into(),
acknowledged_through: 0,
acknowledged_digest: mj_core::relay::RELAY_EVENT_GENESIS_DIGEST.into(),
recovery_floor_ordinal: 0,
recovery_floor_digest: mj_core::relay::RELAY_EVENT_GENESIS_DIGEST.into(),
native_session_id: Some("native-1".into()),
native_continuity_lost: false,
replaced_unused_native_session_id: None,
checkpoint_only: false,
acp_ready: Some(true),
agent_capabilities: None,
agent_info: None,
runtime: None,
steering_supported: None,
config_options: Vec::new(),
modes: None,
available_commands: Vec::new(),
config: BTreeMap::new(),
active_prompt: None,
queued_prompts: Vec::new(),
active_user_shells: Vec::new(),
active_agent_terminals: Vec::new(),
command_ledger_seal: None,
checkpoint_barrier: None,
checkpoint_ready: None,
last_acp_activity_at_ms: None,
current_step_started_at_ms: None,
foreground_tool_started_at_ms: None,
harness_turn: None,
last_harness_turn_started_ordinal: None,
background_commands: Vec::new(),
background_work_known: None,
tools_in_flight: Vec::new(),
activity: None,
};
ManagedSessionView {
snapshot: Some(mj_core::state::ManagedSessionSnapshot {
subagent_requests: Vec::new(),
subagent_results: Vec::new(),
window: mj_core::state::ProjectionWindow::of(&materialized),
materialized,
operational,
latest_credential_sync_signal: None,
worker_build: None,
}),
connected: true,
error: None,
}
}
#[tokio::test]
async fn quiet_background_work_is_reobserved_without_publishing_a_new_view() {
let mut state = test_runtime_state();
let (observations, mut observed) = crate::recovery_gate::observation_channel();
Arc::get_mut(&mut state).unwrap().recovery_observer = RecoveryObserver {
observations,
gate: Arc::new(crate::recovery_gate::RecoveryGate::default()),
};
let mut session = runtime_test_session("session-1", "workspace", SessionState::Running);
session.harness_kind = mj_core::config::HarnessKind::Claude;
state.owner().edit_sessions(|sessions| {
sessions.insert(session.id.clone(), session);
});
let mut view = ready_startup_view();
let snapshot = view.snapshot.as_mut().unwrap();
snapshot.materialized.execution = mj_core::state::MaterializedExecutionState::Idle;
snapshot.window.latest_turn_start_position = Some(9);
state
.publish_session("session-1".into(), view)
.await
.unwrap();
let first = observed.try_recv().unwrap().observation;
assert_eq!(first.checkpoint_wait, None);
assert_eq!(first.latest_completed_turn_ordinal, Some(9));
let revision = state.revisions.current();
state.owner().edit_sessions(|sessions| {
sessions.get_mut("session-1").unwrap().title = "renamed".into();
});
state.refresh_background_policies();
let retry = observed.try_recv().unwrap().observation;
assert_eq!(retry.session.title, "renamed");
assert_eq!(retry.latest_completed_turn_ordinal, Some(9));
assert_eq!(retry.checkpoint_wait, None);
assert_eq!(state.revisions.current(), revision);
state.owner().edit_sessions(|sessions| {
sessions.get_mut("session-1").unwrap().state = SessionState::Parked;
});
state.refresh_background_policies();
assert!(observed.try_recv().is_none());
state.owner().edit_sessions(|sessions| {
sessions.get_mut("session-1").unwrap().state = SessionState::Running;
});
state.refresh_background_policies();
assert!(
observed.try_recv().is_none(),
"resumed sessions need a fresh worker view"
);
}
#[tokio::test]
async fn a_session_with_background_commands_is_not_ready_for_a_recovery_copy() {
let mut state = test_runtime_state();
let (observations, mut observed) = crate::recovery_gate::observation_channel();
Arc::get_mut(&mut state).unwrap().recovery_observer = RecoveryObserver {
observations,
gate: Arc::new(crate::recovery_gate::RecoveryGate::default()),
};
let mut session = runtime_test_session("session-1", "workspace", SessionState::Running);
session.harness_kind = mj_core::config::HarnessKind::Claude;
state.owner().edit_sessions(|sessions| {
sessions.insert(session.id.clone(), session);
});
let mut view = ready_startup_view();
let snapshot = view.snapshot.as_mut().unwrap();
snapshot.materialized.execution = mj_core::state::MaterializedExecutionState::Idle;
snapshot.window.latest_turn_start_position = Some(9);
snapshot.operational.background_commands = vec![mj_core::relay::BackgroundCommand {
id: "shell-1".into(),
started_at_ms: 1,
command: "cargo test".into(),
can_stop: true,
}];
assert!(snapshot.operational.has_work_in_flight());
state
.publish_session("session-1".into(), view)
.await
.unwrap();
let observation = observed.try_recv().unwrap().observation;
assert_eq!(
observation.checkpoint_wait,
Some(mj_core::activity::CheckpointWait::WorkInFlight),
"a session the checkpoint barrier would defer was observed as ready for a copy"
);
state.refresh_background_policies();
assert_eq!(
observed.try_recv().unwrap().observation.checkpoint_wait,
Some(mj_core::activity::CheckpointWait::WorkInFlight)
);
}
#[tokio::test]
async fn disconnected_and_removed_sessions_stop_background_retries() {
let mut state = test_runtime_state();
let (observations, mut observed) = crate::recovery_gate::observation_channel();
Arc::get_mut(&mut state).unwrap().recovery_observer = RecoveryObserver {
observations,
gate: Arc::new(crate::recovery_gate::RecoveryGate::default()),
};
let session = runtime_test_session("session-1", "workspace", SessionState::Running);
state.owner().edit_sessions(|sessions| {
sessions.insert(session.id.clone(), session);
});
state
.publish_session("session-1".into(), ready_startup_view())
.await
.unwrap();
observed.try_recv().unwrap();
let mut disconnected = ready_startup_view();
disconnected.connected = false;
state
.publish_session("session-1".into(), disconnected)
.await
.unwrap();
state.refresh_background_policies();
assert!(observed.try_recv().is_none());
state
.publish_session("session-1".into(), ready_startup_view())
.await
.unwrap();
observed.try_recv().unwrap();
state.owner().edit_sessions(|sessions| {
sessions.remove("session-1");
});
state.refresh_background_policies();
assert!(observed.try_recv().is_none());
assert!(state.owner().background_policies.is_empty());
}
fn insert_starting_session(state: &Arc<RuntimeState>, draft: &str) {
let mut session = runtime_test_session(
"session-1",
mj_core::workspace::DEFAULT_WORKSPACE_ID,
SessionState::Provisioning,
);
draft.clone_into(&mut session.draft_input);
crate::database::save_session(&session).unwrap();
assert!(state.session_record(&session.id).is_some());
}
fn startup_prompt_test_store(name: &str) -> Option<crate::database::DatabaseWriterOwner> {
if std::env::var_os("MJ_TEST_DURABLE_STARTUP").is_none() {
let root = tempfile::tempdir().unwrap();
crate::controller::test_support::IsolatedTest::new(
crate::controller::test_support::test_name(module_path!(), name),
)
.env("MJ_TEST_DURABLE_STARTUP", "1")
.isolated_store(root.path())
.run();
return None;
}
Some(crate::database::install_isolated_test_writer())
}
fn empty_archive_snapshot() -> mj_core::archive::CanonicalSessionSnapshot {
mj_core::archive::CanonicalSessionSnapshot {
command_ledger: None,
assessment_state: None,
event_frontier: 0,
event_frontier_digest: mj_core::relay::RELAY_EVENT_GENESIS_DIGEST.into(),
session: mj_core::archive::CanonicalSessionState {
execution: mj_core::archive::CanonicalExecutionState::Idle,
last_activity_at_ms: None,
session_title: None,
configuration: BTreeMap::new(),
},
transcript: Vec::new(),
queued_prompts: Vec::new(),
}
}
fn queued_prompt_action(text: &str) -> DaemonAction {
DaemonAction::QueueStartupPrompt {
session_id: "session-1".into(),
text: text.into(),
inherited_draft: None,
}
}
fn submitted_prompt_text(command: &RelayCommand) -> String {
let RelayCommand::Prompt { prompt } = command else {
panic!("expected a prompt command, got {command:?}")
};
prompt
.iter()
.map(|block| match block {
agent_client_protocol::schema::v1::ContentBlock::Text(text) => text.text.clone(),
other => panic!("expected text content, got {other:?}"),
})
.collect()
}
fn session_draft_input(state: &RuntimeState) -> String {
state
.session_record("session-1")
.expect("the test session is still in memory")
.draft_input
}
fn notice_texts(state: &RuntimeState) -> Vec<String> {
state
.notices
.lock()
.unwrap()
.iter()
.map(|notice| notice.text.clone())
.collect()
}
async fn next_submit(
manager: &mut TestRemoteManager,
) -> (
String,
tokio::sync::oneshot::Sender<std::result::Result<u64, mj_client::session::SubmitFailure>>,
) {
let request = tokio::time::timeout(Duration::from_secs(10), manager.requests.recv())
.await
.expect("a queued prompt did not reach the session manager")
.expect("manager request stream ended");
let RemoteSessionRequest::Submit { command, reply, .. } = request else {
panic!("expected a submitted prompt")
};
(submitted_prompt_text(&command), reply)
}
#[tokio::test]
async fn queued_startup_prompts_are_delivered_in_order_once_the_harness_is_ready() {
let Some(_writer) = startup_prompt_test_store(
"queued_startup_prompts_are_delivered_in_order_once_the_harness_is_ready",
) else {
return;
};
let mut manager = TestRemoteManager::new().await;
let state = test_runtime_state_with_manager(&manager);
insert_starting_session(&state, "");
let metadata = test_metadata(SocketAddr::from((Ipv4Addr::LOCALHOST, 0)));
let cancellation = CancellationToken::new();
handle_action(
queued_prompt_action("first prompt"),
&metadata,
&state,
&cancellation,
)
.await
.expect("queue the first startup prompt");
assert!(
tokio::time::timeout(Duration::from_millis(600), manager.requests.recv())
.await
.is_err(),
"a prompt was submitted before the harness was ready"
);
manager
.publisher
.publish("session-1".into(), ready_startup_view())
.await
.expect("publish the ready view");
let (text, reply) = next_submit(&mut manager).await;
assert_eq!(text, "first prompt");
reply
.send(Ok(1))
.expect("the drain awaits the submit reply");
handle_action(
queued_prompt_action("second prompt"),
&metadata,
&state,
&cancellation,
)
.await
.expect("queue the second startup prompt");
let (text, reply) = next_submit(&mut manager).await;
assert_eq!(text, "second prompt");
reply
.send(Ok(2))
.expect("the drain awaits the submit reply");
assert!(notice_texts(&state).is_empty(), "delivery posted a notice");
}
#[tokio::test]
async fn a_queued_hand_off_runs_before_the_prompt_behind_it() {
let Some(_writer) =
startup_prompt_test_store("a_queued_hand_off_runs_before_the_prompt_behind_it")
else {
return;
};
let mut manager = TestRemoteManager::new().await;
let state = test_runtime_state_with_manager(&manager);
insert_starting_session(&state, "");
let cancellation = CancellationToken::new();
state
.queue_startup_step(
"session-1",
StartupStep::InstallHandoff(Box::new(empty_archive_snapshot())),
&cancellation,
)
.await
.expect("queue the hand-off");
state
.queue_startup_step(
"session-1",
StartupStep::Prompt {
text: "after the hand-off".into(),
inherited_draft: None,
},
&cancellation,
)
.await
.expect("queue the prompt behind the hand-off");
manager
.publisher
.publish("session-1".into(), ready_startup_view())
.await
.expect("publish the ready view");
let request = tokio::time::timeout(Duration::from_secs(10), manager.requests.recv())
.await
.unwrap()
.unwrap();
let RemoteSessionRequest::Submit { command, reply, .. } = request else {
panic!("expected handoff command")
};
assert!(matches!(command, RelayCommand::InstallPromptContext { .. }));
reply.send(Err("handoff rejected".into())).unwrap();
wait_for_startup_pause(&state).await;
assert_eq!(session_draft_input(&state), "");
assert_eq!(crate::database::load_startup_deliveries().unwrap().len(), 2);
assert!(
tokio::time::timeout(Duration::from_millis(200), manager.requests.recv())
.await
.is_err()
);
}
#[tokio::test]
async fn a_session_that_stops_preserves_its_durable_startup_queue() {
let Some(_writer) =
startup_prompt_test_store("a_session_that_stops_preserves_its_durable_startup_queue")
else {
return;
};
let manager = TestRemoteManager::new().await;
let state = test_runtime_state_with_manager(&manager);
insert_starting_session(&state, "half-written note");
let metadata = test_metadata(SocketAddr::from((Ipv4Addr::LOCALHOST, 0)));
let cancellation = CancellationToken::new();
handle_action(
queued_prompt_action("undelivered prompt"),
&metadata,
&state,
&cancellation,
)
.await
.expect("queue the startup prompt");
state.owner().edit_sessions(|sessions| {
sessions.get_mut("session-1").unwrap().state = SessionState::Stopped;
});
wait_for_startup_pause(&state).await;
assert_eq!(session_draft_input(&state), "half-written note");
let saved = crate::database::load_startup_deliveries().unwrap();
assert_eq!(saved.len(), 1);
assert!(saved[0].step_json.contains("undelivered prompt"));
}
#[tokio::test]
async fn a_refused_submit_preserves_order_without_creating_duplicate_drafts() {
let Some(_writer) = startup_prompt_test_store(
"a_refused_submit_preserves_order_without_creating_duplicate_drafts",
) else {
return;
};
let mut manager = TestRemoteManager::new().await;
let state = test_runtime_state_with_manager(&manager);
insert_starting_session(&state, "");
let metadata = test_metadata(SocketAddr::from((Ipv4Addr::LOCALHOST, 0)));
let cancellation = CancellationToken::new();
for text in ["refused prompt", "prompt behind it"] {
handle_action(queued_prompt_action(text), &metadata, &state, &cancellation)
.await
.expect("queue a startup prompt");
}
manager
.publisher
.publish("session-1".into(), ready_startup_view())
.await
.expect("publish the ready view");
let (text, reply) = next_submit(&mut manager).await;
assert_eq!(text, "refused prompt");
reply
.send(Err("the harness refused the prompt".into()))
.expect("the drain awaits the submit reply");
wait_for_startup_pause(&state).await;
assert_eq!(session_draft_input(&state), "");
let saved = crate::database::load_startup_deliveries().unwrap();
assert_eq!(saved.len(), 2);
assert_eq!(saved[0].phase, "delivering");
assert_eq!(saved[1].phase, "pending");
}
#[tokio::test]
async fn cancelling_startup_prompts_returns_while_a_drain_is_waiting() {
let Some(_writer) =
startup_prompt_test_store("cancelling_startup_prompts_returns_while_a_drain_is_waiting")
else {
return;
};
let manager = TestRemoteManager::new().await;
let state = test_runtime_state_with_manager(&manager);
insert_starting_session(&state, "");
let metadata = test_metadata(SocketAddr::from((Ipv4Addr::LOCALHOST, 0)));
let cancellation = CancellationToken::new();
handle_action(
queued_prompt_action("never delivered"),
&metadata,
&state,
&cancellation,
)
.await
.expect("queue the startup prompt");
let started = tokio::time::Instant::now();
state
.cancel_and_join_startup_prompts()
.await
.expect("the waiting drain stopped cleanly");
assert!(
started.elapsed() < Duration::from_secs(2),
"cancelling the startup queues took {:?}",
started.elapsed()
);
assert_eq!(session_draft_input(&state), "");
assert_eq!(crate::database::load_startup_deliveries().unwrap().len(), 1);
}
#[tokio::test]
async fn startup_restart_reuses_the_durable_command_identity_after_a_lost_reply() {
let Some(_writer) = startup_prompt_test_store(
"startup_restart_reuses_the_durable_command_identity_after_a_lost_reply",
) else {
return;
};
let mut manager = TestRemoteManager::new().await;
let state = test_runtime_state_with_manager(&manager);
insert_starting_session(&state, "");
let cancellation = CancellationToken::new();
state
.queue_startup_step(
"session-1",
StartupStep::Prompt {
text: "once".into(),
inherited_draft: None,
},
&cancellation,
)
.await
.unwrap();
let saved_id = crate::database::load_startup_deliveries().unwrap()[0]
.command_id
.clone();
manager
.publisher
.publish("session-1".into(), ready_startup_view())
.await
.unwrap();
let request = tokio::time::timeout(Duration::from_secs(10), manager.requests.recv())
.await
.unwrap()
.unwrap();
let RemoteSessionRequest::Submit {
command_id, reply, ..
} = request
else {
panic!("expected prompt")
};
assert_eq!(command_id, saved_id);
reply
.send(Err("connection lost after acceptance".into()))
.unwrap();
wait_for_startup_pause(&state).await;
state.cancel_and_join_startup_prompts().await.unwrap();
assert_eq!(session_draft_input(&state), "");
let recovered = test_runtime_state_with_manager(&manager);
recovered
.restore_startup_deliveries(&CancellationToken::new())
.await
.unwrap();
let request = tokio::time::timeout(Duration::from_secs(10), manager.requests.recv())
.await
.unwrap()
.unwrap();
let RemoteSessionRequest::Submit {
command_id, reply, ..
} = request
else {
panic!("expected recovered prompt")
};
assert_eq!(
command_id, saved_id,
"restart assigned a new execution identity"
);
reply.send(Ok(7)).unwrap();
tokio::time::timeout(Duration::from_secs(10), async {
while !crate::database::load_startup_deliveries()
.unwrap()
.is_empty()
{
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.unwrap();
}
#[tokio::test]
async fn startup_batch_identity_rejects_changed_payload_and_keeps_order() {
let Some(_writer) =
startup_prompt_test_store("startup_batch_identity_rejects_changed_payload_and_keeps_order")
else {
return;
};
let manager = TestRemoteManager::new().await;
let state = test_runtime_state_with_manager(&manager);
insert_starting_session(&state, "");
let cancellation = CancellationToken::new();
let steps = || {
vec![
(
"request-first".into(),
StartupStep::Prompt {
text: "first".into(),
inherited_draft: None,
},
),
(
"request-second".into(),
StartupStep::Prompt {
text: "second".into(),
inherited_draft: None,
},
),
]
};
state
.queue_startup_steps_with_ids("session-1", steps(), Some("request".into()), &cancellation)
.await
.unwrap();
state
.queue_startup_steps_with_ids("session-1", steps(), Some("request".into()), &cancellation)
.await
.unwrap();
let mut changed = steps();
changed.pop();
assert!(
state
.queue_startup_steps_with_ids(
"session-1",
changed,
Some("request".into()),
&cancellation
)
.await
.is_err()
);
let rows = crate::database::load_startup_deliveries().unwrap();
assert_eq!(rows.len(), 2);
assert_eq!(rows[0].command_id, "request-first");
assert_eq!(rows[1].command_id, "request-second");
state.cancel_and_join_startup_prompts().await.unwrap();
crate::database::cancel_startup_groups("session-1").unwrap();
assert!(
crate::database::set_startup_delivery_phase("request-first", "delivering", None).is_err()
);
assert!(crate::database::set_startup_delivery_accepted("request-first", Some(9)).is_err());
assert_eq!(
crate::database::next_startup_delivery("session-1")
.unwrap()
.unwrap()
.phase,
"cancelling"
);
}
#[tokio::test]
async fn queueing_a_startup_prompt_is_refused_for_blank_text_and_unusable_sessions() {
let Some(_writer) = startup_prompt_test_store(
"queueing_a_startup_prompt_is_refused_for_blank_text_and_unusable_sessions",
) else {
return;
};
let manager = TestRemoteManager::new().await;
let state = test_runtime_state_with_manager(&manager);
insert_starting_session(&state, "");
let metadata = test_metadata(SocketAddr::from((Ipv4Addr::LOCALHOST, 0)));
let cancellation = CancellationToken::new();
let blank = handle_action(
queued_prompt_action(" "),
&metadata,
&state,
&cancellation,
)
.await
.expect_err("a blank startup prompt was accepted");
assert!(blank.to_string().contains("needs text"), "{blank:#}");
let unknown = handle_action(
DaemonAction::QueueStartupPrompt {
session_id: "no-such-session".into(),
text: "hello".into(),
inherited_draft: None,
},
&metadata,
&state,
&cancellation,
)
.await
.expect_err("a prompt for an unknown session was accepted");
assert!(
unknown.to_string().contains("unknown session"),
"{unknown:#}"
);
state.owner().edit_sessions(|sessions| {
sessions.get_mut("session-1").unwrap().state = SessionState::Stopped;
});
let stopped = handle_action(
queued_prompt_action("hello"),
&metadata,
&state,
&cancellation,
)
.await
.expect_err("a prompt for a stopped session was accepted");
assert!(
stopped.to_string().contains("cannot take a queued prompt"),
"{stopped:#}"
);
}
async fn wait_for_startup_pause(state: &Arc<RuntimeState>) {
tokio::time::timeout(Duration::from_secs(10), async {
while !notice_texts(state)
.iter()
.any(|notice| notice.contains("remains saved"))
{
tokio::time::sleep(Duration::from_millis(25)).await;
}
})
.await
.expect("startup delivery did not report its saved work");
}
#[cfg(unix)]
const DISCARD_LOST_TEST_CHILD: &str = "MJ_TEST_DISCARD_LOST_CHILD";
#[cfg(unix)]
#[tokio::test]
async fn discarding_a_lost_session_removes_its_record_but_keeps_a_dirty_checkout() {
const TEST: &str = "discarding_a_lost_session_removes_its_record_but_keeps_a_dirty_checkout";
if std::env::var_os(DISCARD_LOST_TEST_CHILD).is_none() {
let directory = tempfile::tempdir().unwrap();
crate::controller::test_support::IsolatedTest::new(
crate::controller::test_support::test_name(module_path!(), TEST),
)
.env(DISCARD_LOST_TEST_CHILD, "1")
.isolated_store(directory.path())
.run();
return;
}
let _writer = crate::database::install_isolated_test_writer();
let repository = crate::controller::test_support::committed_repository();
let clean_id = "0123456789abcdef0123456789abcdef";
let dirty_id = "fedcba9876543210fedcba9876543210";
let mut lost = Vec::new();
for session_id in [clean_id, dirty_id] {
let mut session = crate::controller::test_support::managed_worktree_session(
repository.path(),
session_id,
);
session.state = SessionState::Lost;
session.checkpoint = None;
session.last_error = Some("working directory is gone".into());
crate::database::save_session(&session).unwrap();
lost.push(session);
}
let dirty_checkout = repository.path().join(".mj/worktrees").join(dirty_id);
std::fs::write(dirty_checkout.join("unsaved.txt"), "work in progress\n").unwrap();
let state = test_runtime_state_loading_the_store();
for session in &lost {
state.discard_lost_session(session.id.clone()).await;
}
let stored = crate::database::load_state().unwrap().sessions;
assert!(
!stored.contains_key(clean_id) && !stored.contains_key(dirty_id),
"a lost session must not leave a record behind: {:?}",
stored.keys().collect::<Vec<_>>()
);
assert!(
!repository
.path()
.join(".mj/worktrees")
.join(clean_id)
.exists(),
"a clean checkout holds nothing, so it goes with the session"
);
assert!(
dirty_checkout.join("unsaved.txt").exists(),
"uncommitted work must survive a discard nobody confirmed"
);
for session_id in [clean_id, dirty_id] {
assert_eq!(
crate::controller::test_support::test_git(
repository.path(),
&["branch", "--list", &format!("mj/{session_id}")]
)
.trim()
.trim_start_matches(['*', '+', ' ']),
format!("mj/{session_id}"),
"a discard keeps every session branch"
);
}
let notices = notice_texts(&state);
for session_id in [clean_id, dirty_id] {
let expected = format!(
"Session {} was lost because its managed target no longer exists; its record was removed.",
mj_core::state::short_id(session_id)
);
assert_eq!(
notices
.iter()
.filter(|notice| notice.starts_with(&expected))
.count(),
1,
"each discard reports itself exactly once: {notices:?}"
);
}
assert!(
notices
.iter()
.any(|notice| notice.contains(&dirty_checkout.display().to_string())),
"the retained checkout must be named where the person will see it: {notices:?}"
);
}
#[cfg(unix)]
#[tokio::test]
async fn a_destroyed_session_is_still_found_by_its_id() {
const DESTROY_INDEX_TEST_CHILD: &str = "MJ_TEST_DESTROY_INDEX_CHILD";
const TEST: &str = "a_destroyed_session_is_still_found_by_its_id";
if std::env::var_os(DESTROY_INDEX_TEST_CHILD).is_none() {
let directory = tempfile::tempdir().unwrap();
let home = directory.path().join("home");
std::fs::create_dir_all(&home).unwrap();
crate::controller::test_support::IsolatedTest::new(
crate::controller::test_support::test_name(module_path!(), TEST),
)
.env(DESTROY_INDEX_TEST_CHILD, "1")
.isolated_store(directory.path())
.env(
mj_core::config::SESSION_INDEX_ENV,
directory.path().join("sessionwiki"),
)
.env("HOME", &home)
.env("XDG_DATA_HOME", home.join(".local/share"))
.env("XDG_CONFIG_HOME", home.join(".config"))
.run();
return;
}
let _writer = crate::database::install_isolated_test_writer();
let session_id = "0123456789abcdef0123456789abcdef";
let mut session = runtime_test_session(
session_id,
mj_core::workspace::DEFAULT_WORKSPACE_ID,
SessionState::Running,
);
session.title = "destroyed before any sync".into();
crate::database::save_session(&session).unwrap();
let mut conversation = mj_core::state::MaterializedSession::empty(session_id);
conversation.applied_event_ordinal = 2;
conversation.applied_event_digest = format!("{:064x}", 2);
conversation.last_activity_at_ms = Some(1_700_000_000_002);
for (position, body) in [
mj_core::transcript::TranscriptBody::User {
content: vec![serde_json::json!({"type": "text", "text": "find me after destroy"})],
},
mj_core::transcript::TranscriptBody::Agent {
chunks: vec![serde_json::json!({"content": {"type": "text", "text": "noted"}})],
streaming: false,
},
]
.into_iter()
.enumerate()
{
let position = position as u64 + 1;
let streamed = matches!(body, mj_core::transcript::TranscriptBody::Agent { .. });
conversation
.transcript
.push(Arc::new(mj_core::transcript::TranscriptItem {
stable_id: format!("item-{position}"),
position,
latest_content_event_ordinal: streamed.then_some(position),
created_at_ms: 1_700_000_000_000 + position as i64,
last_changed_at_ms: 1_700_000_000_000 + position as i64,
body,
}));
}
crate::database::save_materialized_session(&conversation).unwrap();
sessionwiki::index::open().unwrap();
let state = test_runtime_state_loading_the_store();
state
.force_destroy_session(session_id.to_owned(), BranchDisposition::Keep)
.await
.unwrap();
assert!(
!crate::database::load_state()
.unwrap()
.sessions
.contains_key(session_id),
"destroy removes the record"
);
let found = state
.wiki_session(session_id.to_owned())
.await
.unwrap()
.expect("the destroyed session is found by its id");
assert_eq!(found.status, mj_client::daemon::WikiSessionStatus::Archived);
assert_eq!(found.mjolnir_session_id.as_deref(), Some(session_id));
assert!(
!found.nothing_to_restore,
"the indexed conversation keeps the prompt a restore starts from"
);
}
#[cfg(unix)]
fn test_runtime_state_loading_the_store() -> Arc<RuntimeState> {
let remote = spawn_remote_session_manager().unwrap();
let recovery = crate::recovery::RecoveryCoordinator::spawn(remote.control.clone());
let upgrades = crate::worker_upgrade::WorkerUpgradeCoordinator::spawn(
remote.control.clone(),
&recovery.observer(),
);
Arc::new(RuntimeState::new_with_controller_loader(
remote.control,
Controller::load().unwrap(),
recovery.observer(),
upgrades.observer(),
Vec::new(),
Controller::load,
))
}
#[test]
fn startup_selects_every_record_that_is_only_a_tombstone() {
let sessions = [
runtime_test_session("lost", "workspace", SessionState::Lost),
runtime_test_session(
"data-loss",
"workspace",
SessionState::DestroyedWithDataLoss,
),
runtime_test_session("stopped", "workspace", SessionState::Stopped),
runtime_test_session("running", "workspace", SessionState::Running),
runtime_test_session("failed", "workspace", SessionState::Error),
];
let controller = Controller {
config: Config::default(),
state: mj_core::state::State {
sessions: sessions
.into_iter()
.map(|session| (session.id.clone(), session))
.collect(),
..mj_core::state::State::default()
},
};
assert_eq!(
tombstone_session_ids(&controller),
vec!["data-loss".to_owned(), "lost".to_owned()]
);
}
#[cfg(unix)]
#[tokio::test]
async fn workspace_close_retains_history_discards_drafts_and_refuses_resume_races() {
const TEST: &str = "workspace_close_retains_history_discards_drafts_and_refuses_resume_races";
const CHILD: &str = "MJ_TEST_WORKSPACE_CLOSE_CHILD";
if std::env::var_os(CHILD).is_none() {
let directory = tempfile::tempdir().unwrap();
crate::controller::test_support::IsolatedTest::new(
crate::controller::test_support::test_name(module_path!(), TEST),
)
.env(CHILD, "1")
.isolated_store(directory.path())
.run();
return;
}
let _writer = crate::database::install_isolated_test_writer();
let workspace = crate::database::create_workspace("Close me").unwrap();
let mut history = runtime_test_session("history-close", &workspace.id, SessionState::Stopped);
history.draft_input = "unsent".into();
crate::database::save_session(&history).unwrap();
crate::database::save_detached_draft(
&workspace.id,
Some(&history.id),
"terminal",
Some(42),
"saved draft",
)
.unwrap();
let state = test_runtime_state_loading_the_store();
state.refresh_workspaces().await.unwrap();
assert!(
state
.workspaces()
.borrow()
.iter()
.any(|w| w.id == workspace.id)
);
let admission = state
.workspace_resume_gate(&workspace.id)
.read_owned()
.await;
let error = state
.close_workspace(workspace.id.clone())
.await
.unwrap_err();
assert!(
format!("{error:#}").contains("resume is in progress"),
"{error:#}"
);
assert_eq!(
crate::database::list_detached_drafts(&workspace.id)
.unwrap()
.len(),
1
);
assert!(
state.workspace_closes.lock().unwrap().is_empty(),
"failed close is retryable"
);
drop(admission);
let _other_resume = state.workspace_resume_gate("other").read_owned().await;
let metadata = test_metadata(SocketAddr::from((Ipv4Addr::LOCALHOST, 0)));
assert!(matches!(
handle_action(
DaemonAction::DeleteWorkspace {
workspace_id: workspace.id.clone(),
},
&metadata,
&state,
&CancellationToken::new(),
)
.await
.unwrap(),
DaemonReply::Done
));
let stored = crate::database::load_state().unwrap();
assert_eq!(stored.sessions[&history.id].state, SessionState::Stopped);
assert!(stored.sessions[&history.id].draft_input.is_empty());
assert!(
crate::database::list_detached_drafts(&workspace.id)
.unwrap()
.is_empty()
);
assert!(
crate::database::list_workspaces()
.unwrap()
.iter()
.all(|w| w.id != workspace.id)
);
assert!(
state
.workspaces()
.borrow()
.iter()
.all(|w| w.id != workspace.id)
);
assert!(
!state
.runtime_snapshot("", 0, true)
.await
.unwrap()
.workspace_names
.contains_key(&workspace.id)
);
}
#[cfg(unix)]
#[tokio::test]
async fn workspace_feed_tracks_names_and_a_delayed_refresh_cannot_restore_a_deleted_tab() {
const TEST: &str =
"workspace_feed_tracks_names_and_a_delayed_refresh_cannot_restore_a_deleted_tab";
const CHILD: &str = "MJ_TEST_WORKSPACE_FEED_CHILD";
if std::env::var_os(CHILD).is_none() {
let directory = tempfile::tempdir().unwrap();
crate::controller::test_support::IsolatedTest::new(
crate::controller::test_support::test_name(module_path!(), TEST),
)
.env(CHILD, "1")
.isolated_store(directory.path())
.run();
return;
}
let _writer = crate::database::install_isolated_test_writer();
let state = test_runtime_state_loading_the_store();
let metadata = test_metadata(SocketAddr::from((Ipv4Addr::LOCALHOST, 0)));
let cancellation = CancellationToken::new();
let DaemonReply::Workspace(workspace) = handle_action(
DaemonAction::CreateWorkspace {
name: "Before".into(),
},
&metadata,
&state,
&cancellation,
)
.await
.unwrap() else {
panic!("create did not return the workspace");
};
assert_eq!(state.workspaces().borrow()[0].id, workspace.id);
handle_action(
DaemonAction::RenameWorkspace {
workspace_id: workspace.id.clone(),
name: "After".into(),
},
&metadata,
&state,
&cancellation,
)
.await
.unwrap();
assert_eq!(state.workspaces().borrow()[0].name, "After");
assert_eq!(
state
.runtime_snapshot("", 0, true)
.await
.unwrap()
.workspace_names[&workspace.id],
"After"
);
let refresh_guard = state.workspace_refresh.lock().await;
let delayed = tokio::spawn({
let state = state.clone();
async move { state.refresh_workspaces().await }
});
tokio::task::yield_now().await;
crate::database::close_workspace(&workspace.id).unwrap();
drop(refresh_guard);
delayed.await.unwrap().unwrap();
assert!(
state
.workspaces()
.borrow()
.iter()
.all(|w| w.id != workspace.id)
);
}
#[tokio::test]
async fn workspace_close_cancellation_marks_only_the_requested_operation() {
let state = test_runtime_state();
let first = Arc::new(AtomicBool::new(false));
let second = Arc::new(AtomicBool::new(false));
state.workspace_closes.lock().unwrap().extend([
("first".into(), first.clone()),
("second".into(), second.clone()),
]);
state.cancel_workspace_close("first").unwrap();
assert!(first.load(Ordering::Acquire));
assert!(!second.load(Ordering::Acquire));
assert!(state.cancel_workspace_close("missing").is_err());
}
#[cfg(unix)]
#[tokio::test]
async fn workspace_close_stops_independent_sessions_concurrently_and_retries_after_a_new_create() {
const TEST: &str =
"workspace_close_stops_independent_sessions_concurrently_and_retries_after_a_new_create";
const CHILD: &str = "MJ_TEST_WORKSPACE_CLOSE_CONCURRENCY_CHILD";
if std::env::var_os(CHILD).is_none() {
let directory = tempfile::tempdir().unwrap();
crate::controller::test_support::IsolatedTest::new(
crate::controller::test_support::test_name(module_path!(), TEST),
)
.env(CHILD, "1")
.isolated_store(directory.path())
.run();
return;
}
let _writer = crate::database::install_isolated_test_writer();
let workspace = crate::database::create_workspace("Concurrent close").unwrap();
for id in ["close-first", "close-second"] {
crate::database::save_session(&runtime_test_session(
id,
&workspace.id,
SessionState::Provisioning,
))
.unwrap();
}
crate::database::save_detached_draft(
&workspace.id,
None,
"terminal",
Some(42),
"do not discard on failure",
)
.unwrap();
let state = test_runtime_state_loading_the_store();
let barrier = Arc::new(tokio::sync::Barrier::new(3));
let finish = Arc::new(tokio::sync::Semaphore::new(0));
for id in ["close-first", "close-second"] {
let barrier = barrier.clone();
let finish = finish.clone();
state
.start_or_join_lifecycle(
id.into(),
LifecycleKind::Create,
move |_, _, cancelled| async move {
while !cancelled.load(Ordering::Acquire) {
tokio::time::sleep(Duration::from_millis(5)).await;
}
barrier.wait().await;
finish.acquire().await.unwrap().forget();
Ok(DaemonLifecycleResult::Done)
},
)
.unwrap();
}
let closing = tokio::spawn({
let state = state.clone();
let id = workspace.id.clone();
async move { state.close_workspace(id).await }
});
tokio::time::timeout(Duration::from_secs(5), barrier.wait())
.await
.expect("both independent closes must run concurrently");
crate::database::save_session(&runtime_test_session(
"close-late",
&workspace.id,
SessionState::Provisioning,
))
.unwrap();
finish.add_permits(2);
let error = tokio::time::timeout(Duration::from_secs(10), closing)
.await
.unwrap()
.unwrap()
.unwrap_err();
assert!(
format!("{error:#}").contains("active sessions remain"),
"{error:#}"
);
assert_eq!(
crate::database::list_detached_drafts(&workspace.id)
.unwrap()
.len(),
1
);
let stored = crate::database::load_state().unwrap();
for id in ["close-first", "close-second"] {
assert_eq!(stored.sessions[id].state, SessionState::Stopped);
}
assert_eq!(
stored.sessions["close-late"].state,
SessionState::Provisioning
);
state.close_workspace(workspace.id.clone()).await.unwrap();
let stored = crate::database::load_state().unwrap();
for id in ["close-first", "close-second", "close-late"] {
assert_eq!(stored.sessions[id].state, SessionState::Stopped);
}
assert!(
crate::database::list_detached_drafts(&workspace.id)
.unwrap()
.is_empty()
);
}
#[tokio::test]
async fn oversized_response_reports_its_size_and_keeps_connection_usable() {
let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
let address = listener.local_addr().unwrap();
let sender = tokio::spawn(async move {
let mut stream = TcpStream::connect(address).await.unwrap();
for (request_id, result) in [
(7, Ok(DaemonReply::Text("x".repeat(MAX_FRAME_BYTES)))),
(8, Ok(DaemonReply::Pong)),
] {
super::serve::write_response(
&mut stream,
ResponseEnvelope {
protocol_version: PROTOCOL_VERSION,
request_id,
result,
},
)
.await
.unwrap();
}
});
let (mut stream, _) = listener.accept().await.unwrap();
let response: ResponseEnvelope = read_frame(&mut stream).await.unwrap();
assert_eq!(response.request_id, 7);
let error = response.result.unwrap_err();
assert!(error.contains("text response for request 7"), "{error}");
assert!(error.contains(&MAX_FRAME_BYTES.to_string()), "{error}");
let response: ResponseEnvelope = read_frame(&mut stream).await.unwrap();
assert_eq!(response.request_id, 8);
assert!(matches!(response.result, Ok(DaemonReply::Pong)));
sender.await.unwrap();
}
#[cfg(unix)]
#[tokio::test]
async fn suspension_intent_survives_restart_and_missing_worker_reports_failure() {
const NAME: &str = "suspension_intent_survives_restart_and_missing_worker_reports_failure";
const CHILD: &str = "MJ_TEST_SUSPENSION_RESTART_CHILD";
if std::env::var_os(CHILD).is_none() {
let root = tempfile::tempdir().unwrap();
crate::controller::test_support::IsolatedTest::new(
crate::controller::test_support::test_name(module_path!(), NAME),
)
.env(CHILD, "1")
.isolated_store(root.path())
.run();
return;
}
let _writer = crate::database::install_isolated_test_writer();
let workspace = crate::database::create_workspace("Suspend restart").unwrap();
let root = tempfile::tempdir().unwrap();
let mut session = runtime_test_session("restart-suspend", &workspace.id, SessionState::Running);
session.target = Some(mj_core::state::TargetLocator::LocalBare {
worker_root: root.path().join(&session.id),
});
session.last_error = Some("old connection failure".into());
crate::database::save_session(&session).unwrap();
let state = test_runtime_state_loading_the_store();
state.prepare_suspension(&session.id).await.unwrap();
let stored = crate::database::load_state().unwrap();
assert_eq!(stored.sessions[&session.id].state, SessionState::Closing);
assert!(interrupted_suspend_session_ids(&Controller::load().unwrap()).contains(&session.id));
let restarted = test_runtime_state_loading_the_store();
let failure = tokio::time::timeout(
Duration::from_secs(15),
restarted.suspend_session(session.id.clone()),
)
.await
.unwrap();
assert!(failure.is_err());
let restored = crate::database::load_state().unwrap();
let retained = &restored.sessions[&session.id];
assert_eq!(retained.target, session.target);
assert!(
retained
.public_error()
.unwrap()
.starts_with(mj_core::state::CLOSE_FAILURE_PREFIX)
);
assert!(!restarted.close_is_requested(&session.id));
}
#[cfg(unix)]
#[tokio::test]
async fn suspending_a_parent_stops_its_sub_agents_and_lists_them_on_the_parent() {
const NAME: &str = "suspending_a_parent_stops_its_sub_agents_and_lists_them_on_the_parent";
const CHILD: &str = "MJ_TEST_SUSPEND_STOPS_SUBAGENTS";
if std::env::var_os(CHILD).is_none() {
let directory = tempfile::tempdir().unwrap();
let home = directory.path().join("home");
std::fs::create_dir_all(&home).unwrap();
crate::controller::test_support::IsolatedTest::new(
crate::controller::test_support::test_name(module_path!(), NAME),
)
.env(CHILD, "1")
.isolated_store(directory.path())
.env(
mj_core::config::SESSION_INDEX_ENV,
directory.path().join("sessionwiki"),
)
.env("HOME", &home)
.env("XDG_DATA_HOME", home.join(".local/share"))
.env("XDG_CONFIG_HOME", home.join(".config"))
.run();
return;
}
use mj_core::subagent::StoppedSubagent;
let _writer = crate::database::install_isolated_test_writer();
let workspace = crate::database::create_workspace("Parent suspension").unwrap();
let root = tempfile::tempdir().unwrap();
let parent_id = "0123456789abcdef0123456789abcdef";
let finished_id = "11111111111111111111111111111111";
let stuck_id = "22222222222222222222222222222222";
let mut parent = runtime_test_session(parent_id, &workspace.id, SessionState::Provisioning);
parent.target = Some(mj_core::state::TargetLocator::LocalBare {
worker_root: root.path().join(parent_id),
});
crate::database::save_session(&parent).unwrap();
let mut finished = runtime_test_session(finished_id, &workspace.id, SessionState::Running);
finished.target = Some(mj_core::state::TargetLocator::LocalBare {
worker_root: root.path().join(finished_id),
});
finished.session_title_override = Some("Fix the parser".into());
let mut relation = runtime_test_subagent(finished_id, parent_id);
relation.initial_prompt = "Fix the off-by-one in the parser.".into();
relation.handback_tool = true;
crate::database::save_subagent_session(&finished, &relation).unwrap();
let mut conversation = mj_core::state::MaterializedSession::empty(finished_id);
conversation.applied_event_ordinal = 3;
conversation.applied_event_digest = format!("{:064x}", 3);
conversation.last_activity_at_ms = Some(1_700_000_000_003);
for (position, body) in [
mj_core::transcript::TranscriptBody::User {
content: vec![serde_json::json!({"type": "text", "text": "fix the parser"})],
},
mj_core::transcript::TranscriptBody::Agent {
chunks: vec![serde_json::json!({"content": {"type": "text", "text": "fixed"}})],
streaming: false,
},
]
.into_iter()
.enumerate()
{
let position = position as u64 + 1;
let streamed = matches!(body, mj_core::transcript::TranscriptBody::Agent { .. });
conversation
.transcript
.push(Arc::new(mj_core::transcript::TranscriptItem {
stable_id: format!("item-{position}"),
position,
latest_content_event_ordinal: streamed.then_some(position),
created_at_ms: 1_700_000_000_000 + position as i64,
last_changed_at_ms: 1_700_000_000_000 + position as i64,
body,
}));
}
conversation.last_turn_outcome = Some(mj_core::state::MaterializedTurnOutcome {
diagnostic: None,
usage: None,
command_id: "task-1".into(),
accepted_ordinal: Some(1),
turn_start_position: Some(1),
completed_ordinal: 3,
completed_at_ms: 1_700_000_000_003,
outcome: mj_core::state::TurnOutcomeKind::Completed {
stop_reason: "end_turn".into(),
},
});
crate::database::save_materialized_session(&conversation).unwrap();
assert!(
crate::database::record_subagent_handback(
finished_id,
&mj_core::subagent::SubagentHandback {
command_id: "task-1".into(),
message: "Fixed the parser.".into(),
recorded_at_ms: 1_700_000_000_002,
},
)
.unwrap()
);
let mut stuck = runtime_test_session(stuck_id, &workspace.id, SessionState::Running);
stuck.title = "Review the docs".into();
stuck.target_template_id = "removed-target".into();
stuck.target = Some(mj_core::state::TargetLocator::SshBare {
host: "builder.invalid".into(),
workspace: PathBuf::from(format!(".local/share/hel/workspaces/{stuck_id}")),
worker_id: Some(stuck_id.into()),
});
crate::database::save_subagent_session(&stuck, &runtime_test_subagent(stuck_id, parent_id))
.unwrap();
sessionwiki::index::open().unwrap();
let state = test_runtime_state_loading_the_store();
tokio::time::timeout(
Duration::from_secs(60),
state.suspend_session(parent_id.to_owned()),
)
.await
.unwrap()
.unwrap();
let stored = crate::database::load_state().unwrap();
assert_eq!(stored.sessions[parent_id].state, SessionState::Stopped);
for child in [finished_id, stuck_id] {
assert!(
!stored.sessions.contains_key(child) && !stored.subagents.contains_key(child),
"sub-agent {child} is removed, not suspended"
);
}
assert_eq!(
crate::database::load_stopped_subagents(parent_id).unwrap(),
[
StoppedSubagent {
child_session_id: finished_id.into(),
title: "Fix the parser".into(),
task: Some("Fix the off-by-one in the parser.".into()),
handed_back: true,
},
StoppedSubagent {
child_session_id: stuck_id.into(),
title: "Review the docs".into(),
task: Some("do the task".into()),
handed_back: false,
},
]
);
let found = state
.wiki_session(finished_id.to_owned())
.await
.unwrap()
.expect("the stopped sub-agent is found by its id");
assert_eq!(found.status, mj_client::daemon::WikiSessionStatus::Archived);
}
#[cfg(unix)]
const LIVE_PARENT_CHILD: &str = "MJ_TEST_LIVE_PARENT_CHILD";
#[cfg(unix)]
const WORKING_CHILD: &str = "33333333333333333333333333333333";
#[cfg(unix)]
fn in_isolated_live_parent_test(name: &str) -> bool {
if std::env::var_os(LIVE_PARENT_CHILD).is_some() {
return true;
}
let directory = tempfile::tempdir().unwrap();
let home = directory.path().join("home");
std::fs::create_dir_all(&home).unwrap();
crate::controller::test_support::IsolatedTest::new(crate::controller::test_support::test_name(
module_path!(),
name,
))
.env(LIVE_PARENT_CHILD, "1")
.env(
crate::controller::checkpoint::tests::LATCH_CHECKPOINT_ONLY,
"1",
)
.isolated_store(directory.path())
.env(
mj_core::config::SESSION_INDEX_ENV,
directory.path().join("sessionwiki"),
)
.env("HOME", &home)
.env("XDG_DATA_HOME", home.join(".local/share"))
.env("XDG_CONFIG_HOME", home.join(".config"))
.run();
false
}
#[cfg(unix)]
#[derive(Clone, Copy, PartialEq, Eq)]
enum ParentRecoveryCopy {
Installed,
Missing,
}
#[cfg(unix)]
struct LiveParent {
_directory: tempfile::TempDir,
relay_root: PathBuf,
channels: crate::session_manager::SessionManagerChannels,
state: Arc<RuntimeState>,
}
#[cfg(unix)]
impl LiveParent {
fn relay_state(&self) -> String {
std::fs::read_to_string(self.relay_root.join(mj_core::relay::RELAY_STATE_FILE)).unwrap()
}
fn journal(&self) -> String {
std::fs::read_to_string(
self.relay_root
.join(mj_core::relay::RELAY_JOURNAL_DIR)
.join("active.jsonl"),
)
.unwrap()
}
}
#[cfg(unix)]
async fn live_parent_with_a_working_sub_agent(recovery_copy: ParentRecoveryCopy) -> LiveParent {
use crate::controller::checkpoint::tests::LATCH_RELAY_SESSION;
use std::os::unix::fs::PermissionsExt;
let directory = tempfile::tempdir().unwrap();
let worker_root = directory.path().join(LATCH_RELAY_SESSION);
let relay_root = directory.path().join("relay");
let archives = directory.path().join("archives");
let checkout = directory.path().join("checkout");
for path in [&worker_root, &relay_root, &archives, &checkout] {
std::fs::create_dir_all(path).unwrap();
}
seed_live_session(directory.path(), &relay_root);
let hel = worker_root.join("hel");
std::fs::write(
&hel,
"#!/bin/sh\ncat >/dev/null\necho 'No space left on device' >&2\nexit 1\n",
)
.unwrap();
std::fs::set_permissions(&hel, std::fs::Permissions::from_mode(0o755)).unwrap();
let workspace = crate::database::create_workspace("Parent suspension").unwrap();
let mut parent = crate::controller::test_support::checkpoint_test_session(LATCH_RELAY_SESSION);
parent.workspace_id = workspace.id.clone();
parent.target_template_id = "removed-local".into();
parent.target_runtime = Some((&mj_core::config::TargetTemplate::LocalBare).into());
parent.target = Some(mj_core::state::TargetLocator::LocalBare {
worker_root: worker_root.clone(),
});
parent.project_directory = Some(checkout);
if recovery_copy == ParentRecoveryCopy::Installed {
parent.checkpoint = Some(
crate::controller::test_support::write_checkpoint_gate_archive(
&archives,
LATCH_RELAY_SESSION,
2,
),
);
}
crate::database::save_session(&parent).unwrap();
crate::database::save_materialized_session(&mj_core::state::MaterializedSession::empty(
LATCH_RELAY_SESSION,
))
.unwrap();
let mut child = runtime_test_session(WORKING_CHILD, &workspace.id, SessionState::Running);
child.title = "Review the docs".into();
child.target = Some(mj_core::state::TargetLocator::LocalBare {
worker_root: directory.path().join(WORKING_CHILD),
});
crate::database::save_subagent_session(
&child,
&runtime_test_subagent(WORKING_CHILD, LATCH_RELAY_SESSION),
)
.unwrap();
let (channels, state) = serve_live_session(&relay_root).await;
LiveParent {
_directory: directory,
relay_root,
channels,
state,
}
}
#[cfg(unix)]
fn seed_live_session(directory: &Path, relay_root: &Path) {
use crate::controller::checkpoint::tests::LATCH_RELAY_SESSION;
let mut seed =
mj_worker::relay::DurableRelay::open(relay_root, LATCH_RELAY_SESSION, "1.0.0").unwrap();
seed.record_observation(mj_core::relay::RelayObservation::SessionOpened {
native_session_id: "native-session".into(),
native_continuity_lost: false,
replaced_unused_native_session_id: None,
resumed: true,
})
.unwrap();
seed.record_observation(mj_core::relay::RelayObservation::SessionConfigured {
config_options: Vec::new(),
})
.unwrap();
drop(seed);
let profile_home = directory.join("profile");
std::fs::create_dir_all(&profile_home).unwrap();
Config::update(|config| {
config.profiles.insert(
"codex".into(),
mj_core::config::HarnessProfile {
enabled: true,
kind: mj_core::config::HarnessKind::Codex,
home: profile_home,
environment: BTreeMap::new(),
context_window_bytes: None,
guardian_review_model: None,
},
);
Ok(())
})
.unwrap();
}
#[cfg(unix)]
async fn serve_live_session(
relay_root: &Path,
) -> (
crate::session_manager::SessionManagerChannels,
Arc<RuntimeState>,
) {
use crate::controller::checkpoint::tests::{
LATCH_RELAY_SESSION, ReleaseSupport, latch_relay_target,
};
let channels = crate::session_manager::spawn_session_manager().unwrap();
channels
.targets
.send(vec![latch_relay_target(
relay_root,
None,
ReleaseSupport::Supported,
false,
)])
.unwrap();
channels
.control
.wait_for_session(LATCH_RELAY_SESSION, Duration::from_secs(10))
.await
.unwrap();
let recovery = crate::recovery::RecoveryCoordinator::spawn(channels.control.clone());
let upgrades = crate::worker_upgrade::WorkerUpgradeCoordinator::spawn(
channels.control.clone(),
&recovery.observer(),
);
let state = Arc::new(RuntimeState::new_with_controller_loader(
channels.control.clone(),
Controller::load().unwrap(),
recovery.observer(),
upgrades.observer(),
Vec::new(),
Controller::load,
));
(channels, state)
}
#[cfg(unix)]
const STAND_IN_EXPORT: &str = "MJ_TEST_STAND_IN_EXPORT";
#[cfg(unix)]
#[test]
fn stand_in_worker_exports_a_checkpoint() {
if std::env::var_os(STAND_IN_EXPORT).is_none() {
return;
}
println!();
let exported =
mj_worker::checkpoint::export_from_spec_reader(&mut std::io::stdin().lock()).unwrap();
println!("{}", serde_json::to_string(&exported).unwrap());
}
#[cfg(unix)]
struct LiveClone {
_source: tempfile::TempDir,
_directory: tempfile::TempDir,
checkout: PathBuf,
worker_root: PathBuf,
unpushed_commit: String,
relay_root: PathBuf,
channels: crate::session_manager::SessionManagerChannels,
state: Arc<RuntimeState>,
}
#[cfg(unix)]
impl LiveClone {
fn sealed(&self) -> bool {
let journal = std::fs::read_to_string(
self.relay_root
.join(mj_core::relay::RELAY_JOURNAL_DIR)
.join("active.jsonl"),
)
.unwrap();
journal.lines().any(|line| {
let event: mj_core::relay::RelayEvent = serde_json::from_str(line).unwrap();
matches!(
event.observation,
mj_core::relay::RelayObservation::CommandQueued {
command: RelayCommand::Close { .. },
..
}
)
})
}
}
#[cfg(unix)]
async fn live_clone_with_an_unpushed_commit() -> LiveClone {
use crate::controller::checkpoint::tests::LATCH_RELAY_SESSION;
use crate::controller::test_support::{committed_repository, test_git};
use std::os::unix::fs::PermissionsExt;
let directory = tempfile::tempdir().unwrap();
let worker_root = directory.path().join(LATCH_RELAY_SESSION);
let relay_root = directory.path().join("relay");
for path in [&worker_root.join("profile"), &relay_root] {
std::fs::create_dir_all(path).unwrap();
}
seed_live_session(directory.path(), &relay_root);
let hel = worker_root.join("hel");
std::fs::write(
&hel,
format!(
"#!/bin/sh\n\
[ \"$1 $2\" = 'worker export-checkpoint' ] || \
{{ echo \"the stand-in worker only exports: $*\" >&2; exit 2; }}\n\
{STAND_IN_EXPORT}=1 '{program}' --exact '{test}' --nocapture | grep '^{{'\n",
program = std::env::current_exe().unwrap().display(),
test = crate::controller::test_support::test_name(
module_path!(),
"stand_in_worker_exports_a_checkpoint",
),
),
)
.unwrap();
std::fs::set_permissions(&hel, std::fs::Permissions::from_mode(0o755)).unwrap();
let source = committed_repository();
let origin = directory.path().join("origin.git");
test_git(
directory.path(),
&["init", "--bare", "--initial-branch=master", "origin.git"],
);
test_git(
source.path(),
&["remote", "add", "origin", &origin.to_string_lossy()],
);
test_git(source.path(), &["push", "origin", "master"]);
let mut session =
crate::controller::test_support::managed_clone_session(source.path(), LATCH_RELAY_SESSION);
let checkout = session.project_directory.clone().unwrap();
test_git(&checkout, &["config", "user.name", "Hel Tests"]);
test_git(&checkout, &["config", "user.email", "hel@example.invalid"]);
std::fs::write(checkout.join("session.txt"), "work\n").unwrap();
test_git(&checkout, &["add", "."]);
test_git(&checkout, &["commit", "-m", "session work"]);
let unpushed_commit = test_git(&checkout, &["rev-parse", "HEAD"]);
let workspace = crate::database::create_workspace("Clone suspension").unwrap();
session.workspace_id = workspace.id.clone();
session.state = SessionState::Running;
session.target_template_id = "removed-local".into();
session.target_runtime = Some((&mj_core::config::TargetTemplate::LocalBare).into());
session.target = Some(mj_core::state::TargetLocator::LocalBare {
worker_root: worker_root.clone(),
});
crate::database::save_session(&session).unwrap();
crate::database::save_materialized_session(&mj_core::state::MaterializedSession::empty(
LATCH_RELAY_SESSION,
))
.unwrap();
let (channels, state) = serve_live_session(&relay_root).await;
LiveClone {
_source: source,
_directory: directory,
checkout,
worker_root,
unpushed_commit,
relay_root,
channels,
state,
}
}
#[cfg(unix)]
#[tokio::test]
async fn a_suspend_without_the_acknowledgement_refuses_a_live_clone_with_unpushed_work() {
use crate::controller::checkpoint::tests::LATCH_RELAY_SESSION;
if !in_isolated_live_parent_test(
"a_suspend_without_the_acknowledgement_refuses_a_live_clone_with_unpushed_work",
) {
return;
}
let _writer = crate::database::install_isolated_test_writer();
let clone = live_clone_with_an_unpushed_commit().await;
let refusal = tokio::time::timeout(
Duration::from_secs(60),
clone
.state
.suspend_session_with_ack(LATCH_RELAY_SESSION.to_owned(), false),
)
.await
.unwrap()
.expect_err("a suspend without the acknowledgement must be refused");
assert!(
format!("{refusal:#}").contains("acknowledge_unpublished_work"),
"the refusal names the acknowledgement: {refusal:#}"
);
let stored = crate::database::load_state().unwrap();
let record = &stored.sessions[LATCH_RELAY_SESSION];
assert_eq!(record.state, SessionState::Running, "{record:?}");
assert!(record.target.is_some(), "{record:?}");
assert!(!clone.sealed(), "a refused suspend must not seal the relay");
assert!(clone.worker_root.is_dir());
assert_eq!(
crate::controller::test_support::test_git(&clone.checkout, &["rev-parse", "HEAD"]),
clone.unpushed_commit
);
tokio::time::timeout(
Duration::from_secs(60),
clone.state.suspend_session(LATCH_RELAY_SESSION.to_owned()),
)
.await
.unwrap()
.unwrap();
let stored = crate::database::load_state().unwrap();
let record = &stored.sessions[LATCH_RELAY_SESSION];
assert_eq!(record.state, SessionState::Stopped, "{record:?}");
assert!(clone.sealed());
clone.channels.shutdown.shutdown().await.unwrap();
}
#[cfg(unix)]
#[tokio::test]
async fn a_suspend_whose_checkpoint_fails_leaves_the_sub_agents_running() {
use crate::controller::checkpoint::tests::LATCH_RELAY_SESSION;
if !in_isolated_live_parent_test(
"a_suspend_whose_checkpoint_fails_leaves_the_sub_agents_running",
) {
return;
}
let _writer = crate::database::install_isolated_test_writer();
let parent = live_parent_with_a_working_sub_agent(ParentRecoveryCopy::Missing).await;
let failure = tokio::time::timeout(
Duration::from_secs(60),
parent.state.suspend_session(LATCH_RELAY_SESSION.to_owned()),
)
.await
.unwrap();
assert!(failure.is_err());
let stored = crate::database::load_state().unwrap();
let record = &stored.sessions[LATCH_RELAY_SESSION];
assert_eq!(record.state, SessionState::Running, "{record:?}");
assert!(
record
.last_checkpoint_error
.as_deref()
.is_some_and(|error| error.contains("No space left on device")),
"the checkpoint's export is what failed: {record:?}"
);
assert_eq!(stored.sessions[WORKING_CHILD].state, SessionState::Running);
assert!(stored.subagents.contains_key(WORKING_CHILD));
assert!(
crate::database::load_stopped_subagents(LATCH_RELAY_SESSION)
.unwrap()
.is_empty()
);
assert!(!parent.relay_state().contains("<mj-stopped-subagents>"));
assert!(!parent.journal().contains("Suspend stopped"));
parent.channels.shutdown.shutdown().await.unwrap();
}
#[cfg(unix)]
#[tokio::test]
async fn a_suspend_stops_the_sub_agents_once_the_parents_checkpoint_is_verified() {
use crate::controller::checkpoint::tests::LATCH_RELAY_SESSION;
use mj_core::subagent::StoppedSubagent;
if !in_isolated_live_parent_test(
"a_suspend_stops_the_sub_agents_once_the_parents_checkpoint_is_verified",
) {
return;
}
let _writer = crate::database::install_isolated_test_writer();
let parent = live_parent_with_a_working_sub_agent(ParentRecoveryCopy::Installed).await;
tokio::time::timeout(
Duration::from_secs(60),
parent.state.suspend_session(LATCH_RELAY_SESSION.to_owned()),
)
.await
.unwrap()
.unwrap();
let stored = crate::database::load_state().unwrap();
assert_eq!(
stored.sessions[LATCH_RELAY_SESSION].state,
SessionState::Stopped
);
assert!(
!stored.sessions.contains_key(WORKING_CHILD)
&& !stored.subagents.contains_key(WORKING_CHILD),
"the sub-agent is removed, not suspended"
);
assert_eq!(
crate::database::load_stopped_subagents(LATCH_RELAY_SESSION).unwrap(),
[StoppedSubagent {
child_session_id: WORKING_CHILD.into(),
title: "Review the docs".into(),
task: Some("do the task".into()),
handed_back: false,
}]
);
let sealed = parent.journal().lines().any(|line| {
let event: mj_core::relay::RelayEvent = serde_json::from_str(line).unwrap();
matches!(
event.observation,
mj_core::relay::RelayObservation::CommandQueued {
command: RelayCommand::Close { .. },
..
}
)
});
assert!(sealed, "{}", parent.journal());
parent.channels.shutdown.shutdown().await.unwrap();
}
#[cfg(unix)]
#[tokio::test]
async fn discarding_changes_since_a_checkpoint_stops_the_sub_agents() {
use crate::controller::checkpoint::tests::LATCH_RELAY_SESSION;
if !in_isolated_live_parent_test("discarding_changes_since_a_checkpoint_stops_the_sub_agents") {
return;
}
let _writer = crate::database::install_isolated_test_writer();
let parent = live_parent_with_a_working_sub_agent(ParentRecoveryCopy::Installed).await;
let checkpoint = crate::database::load_state().unwrap().sessions[LATCH_RELAY_SESSION]
.checkpoint
.clone()
.unwrap();
tokio::time::timeout(
Duration::from_secs(60),
parent
.state
.discard_since_checkpoint(LATCH_RELAY_SESSION.to_owned(), checkpoint),
)
.await
.unwrap()
.unwrap();
let stored = crate::database::load_state().unwrap();
assert_eq!(
stored.sessions[LATCH_RELAY_SESSION].state,
SessionState::Stopped
);
assert!(
!stored.sessions.contains_key(WORKING_CHILD)
&& !stored.subagents.contains_key(WORKING_CHILD)
);
let stopped = crate::database::load_stopped_subagents(LATCH_RELAY_SESSION).unwrap();
assert_eq!(
stopped
.iter()
.map(|child| child.child_session_id.as_str())
.collect::<Vec<_>>(),
[WORKING_CHILD]
);
parent.channels.shutdown.shutdown().await.unwrap();
}
#[cfg(unix)]
#[tokio::test]
async fn a_failed_suspend_tells_a_live_parent_at_once_which_sub_agents_were_stopped() {
use crate::controller::checkpoint::tests::LATCH_RELAY_SESSION;
if !in_isolated_live_parent_test(
"a_failed_suspend_tells_a_live_parent_at_once_which_sub_agents_were_stopped",
) {
return;
}
let _writer = crate::database::install_isolated_test_writer();
let parent = live_parent_with_a_working_sub_agent(ParentRecoveryCopy::Missing).await;
crate::database::record_stopped_subagents(
LATCH_RELAY_SESSION,
&[mj_core::subagent::StoppedSubagent {
child_session_id: "44444444444444444444444444444444".into(),
title: "Fix the parser".into(),
task: Some("Fix the off-by-one in the parser.".into()),
handed_back: false,
}],
)
.unwrap();
let failure = tokio::time::timeout(
Duration::from_secs(60),
parent.state.suspend_session(LATCH_RELAY_SESSION.to_owned()),
)
.await
.unwrap();
assert!(failure.is_err());
let relay_state = parent.relay_state();
assert!(
relay_state.contains("<mj-stopped-subagents>") && relay_state.contains("Fix the parser"),
"{relay_state}"
);
assert!(
parent
.journal()
.contains("Suspend stopped 1 sub-agent: \\\"Fix the parser\\\" (had not handed back)."),
"{}",
parent.journal()
);
assert!(
crate::database::load_stopped_subagents(LATCH_RELAY_SESSION)
.unwrap()
.is_empty()
);
parent.channels.shutdown.shutdown().await.unwrap();
}
#[cfg(unix)]
#[tokio::test]
async fn a_sub_agent_its_parents_suspend_stops_is_shown_stopping_not_destroying() {
const NAME: &str = "a_sub_agent_its_parents_suspend_stops_is_shown_stopping_not_destroying";
const CHILD: &str = "MJ_TEST_SUSPEND_SHOWS_STOPPING";
if std::env::var_os(CHILD).is_none() {
let directory = tempfile::tempdir().unwrap();
crate::controller::test_support::IsolatedTest::new(
crate::controller::test_support::test_name(module_path!(), NAME),
)
.env(CHILD, "1")
.isolated_store(directory.path())
.run();
return;
}
let _writer = crate::database::install_isolated_test_writer();
let workspace = crate::database::create_workspace("Parent suspension").unwrap();
let root = tempfile::tempdir().unwrap();
let parent_id = "0123456789abcdef0123456789abcdef";
let mut parent = runtime_test_session(parent_id, &workspace.id, SessionState::Provisioning);
parent.target = Some(mj_core::state::TargetLocator::LocalBare {
worker_root: root.path().join(parent_id),
});
crate::database::save_session(&parent).unwrap();
let mut child = runtime_test_session(WORKING_CHILD, &workspace.id, SessionState::Running);
child.target = Some(mj_core::state::TargetLocator::LocalBare {
worker_root: root.path().join(WORKING_CHILD),
});
crate::database::save_subagent_session(
&child,
&runtime_test_subagent(WORKING_CHILD, parent_id),
)
.unwrap();
let mut state = test_runtime_state_loading_the_store();
let mut recovery = crate::recovery::RecoveryCoordinator::spawn(state.session_manager.clone());
Arc::get_mut(&mut state).unwrap().recovery_observer = recovery.observer();
let recovery_attempt = state
.recovery_observer
.gate
.try_start(WORKING_CHILD)
.unwrap();
let suspend = tokio::spawn({
let state = state.clone();
async move { state.suspend_session(parent_id.to_owned()).await }
});
let stopping = tokio::time::timeout(Duration::from_secs(30), async {
loop {
if let Some(view) = state
.active_lifecycles()
.into_iter()
.find(|view| view.session_id == WORKING_CHILD)
{
return view;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.unwrap();
drop(recovery_attempt);
assert_eq!(
serde_json::to_value(stopping.kind).unwrap(),
"stop_subagent",
"{stopping:?}"
);
assert!(!stopping.cancellable, "{stopping:?}");
tokio::time::timeout(Duration::from_secs(60), suspend)
.await
.unwrap()
.unwrap()
.unwrap();
assert!(
!crate::database::load_state()
.unwrap()
.sessions
.contains_key(WORKING_CHILD)
);
recovery.shutdown().await.unwrap();
}
#[cfg(unix)]
fn in_isolated_parked_test(name: &str) -> bool {
const PARKED_CHILD: &str = "MJ_TEST_PARKED_SUBAGENT_DAEMON_CHILD";
if std::env::var_os(PARKED_CHILD).is_some() {
return true;
}
let directory = tempfile::tempdir().unwrap();
let home = directory.path().join("home");
std::fs::create_dir_all(&home).unwrap();
crate::controller::test_support::IsolatedTest::new(crate::controller::test_support::test_name(
module_path!(),
name,
))
.env(PARKED_CHILD, "1")
.isolated_store(directory.path())
.env(
mj_core::config::SESSION_INDEX_ENV,
directory.path().join("sessionwiki"),
)
.env("HOME", &home)
.env("XDG_DATA_HOME", home.join(".local/share"))
.env("XDG_CONFIG_HOME", home.join(".config"))
.run();
false
}
#[cfg(unix)]
fn parked_child(root: &Path, child_id: &str, parent_id: &str, workspace_id: &str) -> SessionRecord {
use std::os::unix::fs::PermissionsExt;
let worker_root = root.join(child_id);
std::fs::create_dir_all(&worker_root).unwrap();
let hel = worker_root.join("hel");
std::fs::write(
&hel,
format!(
"#!/bin/sh\necho started >> {}\n",
root.join("started").display()
),
)
.unwrap();
std::fs::set_permissions(&hel, std::fs::Permissions::from_mode(0o755)).unwrap();
let mut child = runtime_test_session(child_id, workspace_id, SessionState::Parked);
child.target = Some(mj_core::state::TargetLocator::LocalBare { worker_root });
child.session_title_override = Some("Map the parser".into());
crate::database::save_subagent_session(&child, &runtime_test_subagent(child_id, parent_id))
.unwrap();
child
}
#[cfg(unix)]
#[tokio::test]
async fn suspending_a_parent_removes_its_parked_sub_agents_without_starting_or_warning() {
if !in_isolated_parked_test(
"suspending_a_parent_removes_its_parked_sub_agents_without_starting_or_warning",
) {
return;
}
let _writer = crate::database::install_isolated_test_writer();
let workspace = crate::database::create_workspace("Parked children").unwrap();
let root = tempfile::tempdir().unwrap();
let parent_id = "0123456789abcdef0123456789abcdef";
let child_id = "44444444444444444444444444444444";
let mut parent = runtime_test_session(parent_id, &workspace.id, SessionState::Provisioning);
parent.target = Some(mj_core::state::TargetLocator::LocalBare {
worker_root: root.path().join(parent_id),
});
crate::database::save_session(&parent).unwrap();
parked_child(root.path(), child_id, parent_id, &workspace.id);
assert!(crate::controller::subagent_has_handed_back(child_id).unwrap());
sessionwiki::index::open().unwrap();
let state = test_runtime_state_loading_the_store();
tokio::time::timeout(
Duration::from_secs(60),
state.suspend_session(parent_id.to_owned()),
)
.await
.unwrap()
.unwrap();
let stored = crate::database::load_state().unwrap();
assert_eq!(stored.sessions[parent_id].state, SessionState::Stopped);
assert!(
!stored.sessions.contains_key(child_id) && !stored.subagents.contains_key(child_id),
"the parked child is removed with its parent's suspend"
);
let stopped = crate::database::load_stopped_subagents(parent_id).unwrap();
assert_eq!(stopped.len(), 1);
assert!(stopped[0].handed_back, "a parked child is not warned about");
assert!(
!root.path().join("started").exists(),
"the parked child's worker was not started"
);
}
#[cfg(unix)]
#[tokio::test]
async fn closing_a_parked_sub_agent_stops_it_without_starting_its_worker() {
if !in_isolated_parked_test("closing_a_parked_sub_agent_stops_it_without_starting_its_worker") {
return;
}
let _writer = crate::database::install_isolated_test_writer();
let workspace = crate::database::create_workspace("Parked close").unwrap();
let root = tempfile::tempdir().unwrap();
let parent_id = "0123456789abcdef0123456789abcdef";
let child_id = "55555555555555555555555555555555";
let parent = runtime_test_session(parent_id, &workspace.id, SessionState::Running);
crate::database::save_session(&parent).unwrap();
let child = parked_child(root.path(), child_id, parent_id, &workspace.id);
let worker_root = match &child.target {
Some(mj_core::state::TargetLocator::LocalBare { worker_root }) => worker_root.clone(),
_ => unreachable!(),
};
let state = test_runtime_state_loading_the_store();
tokio::time::timeout(
Duration::from_secs(60),
state.suspend_session(child_id.to_owned()),
)
.await
.unwrap()
.unwrap();
let stored = crate::database::load_state().unwrap();
let closed = &stored.sessions[child_id];
assert_eq!(closed.state, SessionState::Stopped);
assert!(
closed.target.is_none(),
"the child's private target state is gone"
);
assert!(
stored.subagents.contains_key(child_id),
"the relation stays"
);
assert!(
!worker_root.exists(),
"only the child's own worker root is removed"
);
assert!(
!root.path().join("started").exists(),
"the parked child's worker was not started"
);
assert_eq!(stored.sessions[parent_id].state, SessionState::Running);
}
#[tokio::test]
async fn late_stage_callbacks_cannot_change_a_replacement_lifecycle() {
let state = test_runtime_state();
let first_release = Arc::new(tokio::sync::Notify::new());
let first = state
.start_or_join_lifecycle("same-session".into(), LifecycleKind::Create, {
let release = first_release.clone();
move |_, _, _| async move {
release.notified().await;
Ok(DaemonLifecycleResult::Done)
}
})
.unwrap();
let old = DaemonStageReportingExecutor::new(
RefusingExecutor("stage-only test"),
state.clone(),
"same-session".into(),
);
first_release.notify_one();
RuntimeState::wait_lifecycle_result(first).await.unwrap();
let second_release = Arc::new(tokio::sync::Notify::new());
let second = state
.start_or_join_lifecycle("same-session".into(), LifecycleKind::Resume, {
let release = second_release.clone();
move |_, _, _| async move {
release.notified().await;
Ok(DaemonLifecycleResult::Done)
}
})
.unwrap();
let current = DaemonStageReportingExecutor::new(
RefusingExecutor("stage-only test"),
state.clone(),
"same-session".into(),
);
old.stage_started(ProvisionStage::Cloning);
old.notify_notice("obsolete operation");
assert!(state.active_lifecycles()[0].active_stages.is_empty());
assert!(state.active_lifecycles()[0].notice.is_none());
current.stage_started(ProvisionStage::Cloning);
old.stage_finished(ProvisionStage::Cloning);
assert_eq!(state.active_lifecycles()[0].active_stages.len(), 1);
second_release.notify_one();
RuntimeState::wait_lifecycle_result(second).await.unwrap();
}
#[tokio::test]
async fn a_late_worker_view_cannot_recreate_a_deleted_session() {
let state = test_runtime_state();
state.owner().edit_sessions(|sessions| {
sessions.insert(
"session-1".into(),
runtime_test_session("session-1", "workspace", SessionState::Running),
);
});
state
.publish_session("session-1".into(), ready_startup_view())
.await
.unwrap();
assert!(state.owner().sessions.contains_key("session-1"));
state.owner().edit_sessions(|sessions| {
sessions.remove("session-1");
});
assert!(!state.owner().sessions.contains_key("session-1"));
state
.publish_session("session-1".into(), ready_startup_view())
.await
.unwrap();
assert!(!state.owner().sessions.contains_key("session-1"));
}
#[tokio::test]
async fn completed_lifecycle_retention_does_not_grow_with_history() {
let state = test_runtime_state();
for index in 0..300 {
let result = state
.start_or_join_lifecycle(
format!("session-{index}"),
LifecycleKind::Create,
|_, _, _| async { Ok(DaemonLifecycleResult::Done) },
)
.unwrap();
RuntimeState::wait_lifecycle_result(result).await.unwrap();
}
assert_eq!(state.owner().lifecycle.len(), 256);
assert!(state.active_lifecycles().is_empty());
assert!(state.owner().lifecycle.contains_key("session-299"));
assert!(!state.owner().lifecycle.contains_key("session-0"));
}
#[tokio::test]
async fn move_destination_ownership_is_typed_and_survives_cancellation() {
let state = test_runtime_state();
state.owner().edit_sessions(|sessions| {
sessions.insert(
"moving".into(),
runtime_test_session("moving", "workspace", SessionState::Running),
);
});
let (_, result) = tokio::sync::watch::channel(None);
state.owner().lifecycle.insert(
"moving".into(),
ActiveLifecycle {
phase: LifecyclePhase::Executing,
operation_id: "current".into(),
create_control: None,
kind: LifecycleKind::Move,
cancelled: Arc::new(AtomicBool::new(false)),
started_at_epoch_seconds: 1,
active_stages: BTreeMap::new(),
resume_workspace_id: None,
resume_destination: None,
notice: None,
request_key: None,
_move_guard: None,
result,
},
);
state.set_lifecycle_notice("moving", "current", "Preparing destination");
assert!(!state.owner().worker_is_owned("moving"));
state.reserve_move_destination("moving", "obsolete");
assert!(!state.owner().worker_is_owned("moving"));
state.reserve_move_destination("moving", "current");
assert!(state.owner().worker_is_owned("moving"));
state
.owner()
.lifecycle
.get_mut("moving")
.unwrap()
.request_cancel();
assert!(state.owner().worker_is_owned("moving"));
state.owner().lifecycle.get_mut("moving").unwrap().phase =
LifecyclePhase::Completed(Ok(DaemonLifecycleResult::Done));
assert!(!state.owner().worker_is_owned("moving"));
}
#[tokio::test]
async fn runtime_publication_crosses_frame_limit_atomically_and_keeps_connection_usable() {
use mj_client::runtime_feed::{RuntimeCursor, RuntimeFrame, RuntimeProjection};
let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
let address = listener.local_addr().unwrap();
let mut projection = RuntimeProjection::default();
projection
.metadata
.workspace_names
.insert("large".into(), "x".repeat(MAX_FRAME_BYTES + 1024));
let expected = projection.clone();
let sender = tokio::spawn(async move {
let mut stream = TcpStream::connect(address).await.unwrap();
for (request_id, result) in [
(
7,
Ok(DaemonReply::RuntimeChanges(Box::new(
RuntimeFrame::Snapshot {
cursor: RuntimeCursor {
incarnation: "test".into(),
sequence: 1,
},
projection: Box::new(projection),
},
))),
),
(8, Ok(DaemonReply::Pong)),
] {
super::serve::write_response(
&mut stream,
ResponseEnvelope {
protocol_version: PROTOCOL_VERSION,
request_id,
result,
},
)
.await
.unwrap();
}
});
let (mut stream, _) = listener.accept().await.unwrap();
let response = mj_client::daemon::read_response(&mut stream, true)
.await
.unwrap();
assert_eq!(response.request_id, 7);
let Ok(DaemonReply::RuntimeChanges(frame)) = response.result else {
panic!("expected complete publication")
};
let RuntimeFrame::Snapshot { projection, .. } = *frame else {
panic!("expected snapshot")
};
assert_eq!(*projection, expected);
let response = mj_client::daemon::read_response(&mut stream, false)
.await
.unwrap();
assert_eq!(response.request_id, 8);
assert!(matches!(response.result, Ok(DaemonReply::Pong)));
sender.await.unwrap();
}
#[tokio::test]
async fn interrupted_runtime_fragments_never_publish_partial_state() {
let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
let address = listener.local_addr().unwrap();
let sender = tokio::spawn(async move {
let mut stream = TcpStream::connect(address).await.unwrap();
write_frame(
&mut stream,
&ResponseEnvelope {
protocol_version: PROTOCOL_VERSION,
request_id: 7,
result: Ok(DaemonReply::RuntimeChunk {
bytes: b"{\"Snapshot\":".to_vec(),
finished: false,
}),
},
)
.await
.unwrap();
});
let (mut stream, _) = listener.accept().await.unwrap();
assert!(
mj_client::daemon::read_response(&mut stream, true)
.await
.is_err()
);
sender.await.unwrap();
}
#[tokio::test]
async fn recovered_drafts_append_to_committed_edits_and_failed_writes_do_not_publish() {
let Some(_writer) = startup_prompt_test_store(
"recovered_drafts_append_to_committed_edits_and_failed_writes_do_not_publish",
) else {
return;
};
let manager = TestRemoteManager::new().await;
let state = test_runtime_state_with_manager(&manager);
insert_starting_session(&state, "original");
crate::database::set_session_draft_input("session-1", "newer client edit").unwrap();
let (first, second) = tokio::join!(
state.append_draft_input("session-1", "first"),
state.append_draft_input("session-1", "second")
);
first.unwrap();
second.unwrap();
let saved = state.session_record("session-1").unwrap().draft_input;
assert!(saved.starts_with("newer client edit\n\n"));
assert!(saved.contains("first") && saved.contains("second"));
assert!(
state
.append_draft_input("missing", "lost input")
.await
.is_err()
);
assert!(state.session_record("missing").is_none());
}
#[tokio::test]
async fn retry_admission_reserves_only_the_matching_move_destination() {
let Some(_writer) =
startup_prompt_test_store("retry_admission_reserves_only_the_matching_move_destination")
else {
return;
};
let manager = TestRemoteManager::new().await;
let state = test_runtime_state_with_manager(&manager);
let session = runtime_test_session(
"session-1",
mj_core::workspace::DEFAULT_WORKSPACE_ID,
SessionState::Running,
);
crate::database::save_session(&session).unwrap();
let operation = mj_core::state::MoveOperation {
source_checkpoint_only: false,
in_place: false,
operation_id: "retained-move".into(),
selection: MoveSelection {
session_id: session.id.clone(),
profile_id: Some(session.last_profile.clone()),
target_template_id: Some(session.target_template_id.clone()),
additional_mounts: None,
resource_allocation: None,
clear_resource_allocation: false,
},
source_profile_id: session.last_profile.clone(),
source_target_template_id: session.target_template_id.clone(),
source_target: None,
source_native_session_id: None,
source_additional_mounts: Vec::new(),
source_resource_allocation: None,
destination_target: None,
destination_native_session_id: None,
destination_store_id: None,
configuration_fingerprint: "test".into(),
checkpoint: None,
recovery_session: Some(session.clone()),
queue: mj_core::state::ResumeQueueDisposition::Start,
phase: mj_core::state::MovePhase::StartingQueue,
queue_admission_started: true,
queue_admission_finished: false,
cancellation_requested: false,
created_at: session.created_at.clone(),
updated_at: session.updated_at.clone(),
error: None,
};
crate::database::save_move_operation(&operation).unwrap();
crate::controller::move_session::restore_move_queue_hold(&operation);
assert!(
state
.start_or_join_lifecycle(session.id.clone(), LifecycleKind::Resume, |_, _, _| async {
Ok(DaemonLifecycleResult::Done)
})
.is_err()
);
for (id, owned) in [("retained-move", true), ("different-move", false)] {
let release = Arc::new(tokio::sync::Notify::new());
let released = release.clone();
let watch = state
.admit_lifecycle(
session.id.clone(),
LifecycleKind::Move,
super::lifecycle::LifecycleStart {
resume_workspace_id: None,
request_key: Some(id.into()),
create_control: None,
phase: LifecyclePhase::Executing,
move_operation_id: Some(id.into()),
},
move |_, _, _| async move {
released.notified().await;
Ok(DaemonLifecycleResult::Done)
},
)
.unwrap();
assert_eq!(state.owner().worker_is_owned(&session.id), owned);
release.notify_one();
RuntimeState::wait_lifecycle_result(watch.clone())
.await
.unwrap();
state.remove_completed_lifecycle(&watch);
}
}
#[tokio::test]
async fn api_startup_persists_the_entire_ordered_followup_before_acknowledging() {
let Some(_writer) = startup_prompt_test_store(
"api_startup_persists_the_entire_ordered_followup_before_acknowledging",
) else {
return;
};
let mut manager = TestRemoteManager::new().await;
let state = test_runtime_state_with_manager(&manager);
insert_starting_session(&state, "");
let followup = crate::server::api::StartFollowup {
model: Some("model-a".into()),
effort: Some("high".into()),
fast_mode: true,
prompt: Some("accepted before readiness".into()),
};
state
.queue_api_followup("session-1", "api-test".into(), followup.clone())
.await
.unwrap();
state
.queue_api_followup("session-1", "api-test".into(), followup.clone())
.await
.unwrap();
let rows = crate::database::load_latest_startup_group("session-1").unwrap();
assert_eq!(
rows.len(),
4,
"retry must reuse the entire existing request"
);
let commands: Vec<_> = rows.iter().map(|row| row.command_id.as_str()).collect();
assert_eq!(
commands,
[
"api-test:model",
"api-test:effort",
"api-test:fast-mode",
"api-test:prompt"
]
);
assert!(
matches!(serde_json::from_str::<StartupStep>(&rows.last().unwrap().step_json).unwrap(), StartupStep::ApiPrompt { text } if text == "accepted before readiness")
);
let mut changed = followup;
changed.prompt = None;
assert!(
state
.queue_api_followup("session-1", "api-test".into(), changed)
.await
.is_err(),
"a request ID cannot change the saved workflow"
);
assert!(
tokio::time::timeout(Duration::from_millis(10), manager.requests.recv())
.await
.is_err(),
"nothing is submitted before readiness"
);
assert!(
!crate::upgrade::active_labels()
.iter()
.any(|label| label.contains("API startup followup"))
);
state.cancel_and_join_startup_prompts().await.unwrap();
assert_eq!(
crate::database::load_latest_startup_group("session-1")
.unwrap()
.len(),
4
);
}
#[tokio::test]
async fn cancelled_api_startup_is_not_submitted_after_daemon_reconstruction() {
let Some(_writer) = startup_prompt_test_store(
"cancelled_api_startup_is_not_submitted_after_daemon_reconstruction",
) else {
return;
};
let mut manager = TestRemoteManager::new().await;
let state = test_runtime_state_with_manager(&manager);
insert_starting_session(&state, "");
state
.queue_api_followup(
"session-1",
"cancelled-api".into(),
crate::server::api::StartFollowup {
prompt: Some("must never start".into()),
..Default::default()
},
)
.await
.unwrap();
state.cancel_api_followup("session-1").await.unwrap();
state.cancel_and_join_startup_prompts().await.unwrap();
let recovered = test_runtime_state_with_manager(&manager);
manager
.publisher
.publish("session-1".into(), ready_startup_view())
.await
.unwrap();
recovered
.restore_startup_deliveries(&CancellationToken::new())
.await
.unwrap();
tokio::time::timeout(Duration::from_secs(5), async {
loop {
if crate::database::next_startup_delivery("session-1")
.unwrap()
.is_none()
{
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("cancelled work should finish receipt cleanup");
assert!(
tokio::time::timeout(Duration::from_millis(10), manager.requests.recv())
.await
.is_err(),
"cancelled initial input must not become a new prompt on restart"
);
assert_eq!(
crate::database::load_latest_startup_group("session-1").unwrap()[0].phase,
"dismissed"
);
}
#[tokio::test]
async fn startup_retry_keeps_earlier_uncertain_input_ahead_of_new_input() {
let Some(_writer) =
startup_prompt_test_store("startup_retry_keeps_earlier_uncertain_input_ahead_of_new_input")
else {
return;
};
let mut manager = TestRemoteManager::new().await;
let state = test_runtime_state_with_manager(&manager);
insert_starting_session(&state, "");
let cancellation = CancellationToken::new();
state
.queue_startup_step(
"session-1",
StartupStep::Prompt {
text: "first".into(),
inherited_draft: None,
},
&cancellation,
)
.await
.unwrap();
manager
.publisher
.publish("session-1".into(), ready_startup_view())
.await
.unwrap();
let (text, reply) = next_submit(&mut manager).await;
assert_eq!(text, "first");
reply.send(Err("transient disconnect".into())).unwrap();
wait_for_startup_pause(&state).await;
state
.queue_startup_step(
"session-1",
StartupStep::Prompt {
text: "second".into(),
inherited_draft: None,
},
&cancellation,
)
.await
.unwrap();
let (text, reply) = next_submit(&mut manager).await;
assert_eq!(
text, "first",
"new input overtook the uncertain earlier command"
);
reply.send(Ok(4)).unwrap();
let (text, reply) = next_submit(&mut manager).await;
assert_eq!(text, "second");
reply.send(Ok(8)).unwrap();
state.cancel_and_join_startup_prompts().await.unwrap();
}
#[cfg(unix)]
#[tokio::test]
async fn replayed_child_close_cannot_stop_or_mark_a_resumed_incarnation() {
if !in_isolated_parked_test("replayed_child_close_cannot_stop_or_mark_a_resumed_incarnation") {
return;
}
let _writer = crate::database::install_isolated_test_writer();
let workspace = crate::database::create_workspace("Close replay").unwrap();
let child_id = "eeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee";
let mut child = runtime_test_session(child_id, &workspace.id, SessionState::Running);
crate::database::save_session(&child).unwrap();
let previous = crate::database::session_incarnation(child_id)
.unwrap()
.unwrap();
child.state = SessionState::Stopped;
crate::database::save_lifecycle_session(&child).unwrap();
child.state = SessionState::Provisioning;
crate::database::save_lifecycle_session(&child).unwrap();
child.state = SessionState::Running;
child.last_error = Some("new incarnation diagnostic".into());
crate::database::save_lifecycle_session(&child).unwrap();
let current = crate::database::session_incarnation(child_id)
.unwrap()
.unwrap();
assert_ne!(previous, current);
let state = test_runtime_state_loading_the_store();
let error = state
.close_subagent_request(
child_id.into(),
"parent".into(),
"old-close".into(),
previous,
)
.await
.unwrap_err();
assert!(error.to_string().contains("earlier child incarnation"));
let stored = crate::database::read_durable_session_record(child_id)
.unwrap()
.unwrap();
assert_eq!(stored.state, SessionState::Running);
assert_eq!(stored.last_error, child.last_error);
assert_eq!(
crate::database::session_incarnation(child_id)
.unwrap()
.as_deref(),
Some(current.as_str())
);
assert!(!state.owner().lifecycle.contains_key(child_id));
}