diff --git a/crates/nebula-daemon/src/server.rs b/crates/nebula-daemon/src/server.rs index 27d5b918..e2a440e2 100644 --- a/crates/nebula-daemon/src/server.rs +++ b/crates/nebula-daemon/src/server.rs @@ -58,6 +58,8 @@ async fn handle_client(daemon: Arc, stream: UnixStream) -> Result<()> { // Per-connection attach state: forward-task handles keyed by session. let mut attached: HashMap> = HashMap::new(); + // The task forwarding broadcasts to a subscribed client. + let mut subscription: Option> = None; let mut handshaken = false; let result: Result<()> = async { @@ -91,6 +93,10 @@ async fn handle_client(daemon: Arc, stream: UnixStream) -> Result<()> { break; } ClientRequest::Subscribe => { + // Subscribed before the snapshot is taken, so a change + // made between the two is not lost: it arrives after the + // Snapshot, where a client folds it in by id. + let mut rx = daemon.events.subscribe(); let snapshot = daemon.snapshot().unwrap_or(ServerEvent::Snapshot { projects: vec![], worktrees: vec![], @@ -101,9 +107,8 @@ async fn handle_client(daemon: Arc, stream: UnixStream) -> Result<()> { ui_state: None, }); let _ = out_tx.send(snapshot).await; - let mut rx = daemon.events.subscribe(); let tx = out_tx.clone(); - tokio::spawn(async move { + let forward = tokio::spawn(async move { loop { match rx.recv().await { Ok(ev) => { @@ -118,6 +123,9 @@ async fn handle_client(daemon: Arc, stream: UnixStream) -> Result<()> { } } }); + if let Some(old) = subscription.replace(forward) { + old.abort(); + } } ClientRequest::Attach { session: sref, @@ -592,6 +600,12 @@ async fn handle_client(daemon: Arc, stream: UnixStream) -> Result<()> { h.abort(); daemon.note_detached(&sref); } + // The forwarder holds a sender and waits on the next broadcast: left + // running, it keeps the writer, and this client's socket, open until + // the daemon next has something to say. + if let Some(forward) = subscription { + forward.abort(); + } drop(out_tx); let _ = writer_task.await; result diff --git a/crates/nebula/tests/e2e_pty.rs b/crates/nebula/tests/e2e_pty.rs index 81033661..2d7964b5 100644 --- a/crates/nebula/tests/e2e_pty.rs +++ b/crates/nebula/tests/e2e_pty.rs @@ -358,14 +358,17 @@ async fn full_crud_attach_and_restart_persistence() { &mut c, &ClientRequest::Input { session: sref.clone(), - data: format!("echo {marker}; pwd\n").into_bytes(), + data: format!("pwd; echo {marker}\n").into_bytes(), }, ) .await .unwrap(); let events = read_events_until(&mut c, EVENT_TIMEOUT, |evs| { let text = String::from_utf8_lossy(&collected_output(evs)).into_owned(); - text.matches(marker).count() >= 2 + // The typed line is echoed, and echoed again when the shell's line + // editor redraws what was typed before its prompt. Only a marker at + // the start of a line is the command's own, printed after `pwd`'s. + text.contains(&format!("\n{marker}")) }) .await; // The shell runs in the worktree directory. @@ -3994,6 +3997,37 @@ async fn subscribe(c: &mut UnixStream) -> Vec { .await } +/// A subscribed client that hangs up is let go at once, with nothing +/// broadcast in between: the daemon closes its side of the connection as +/// soon as it reads the hang-up. It used to keep the socket open until the +/// next broadcast, so an idle daemon collected one per one-shot client. +#[tokio::test] +async fn a_subscribed_client_that_hangs_up_is_let_go_by_an_idle_daemon() { + use tokio::io::AsyncWriteExt; + let env = TestEnv::new(); + let mut daemon = env.spawn_daemon(); + let mut c = connect(&env.sock()).await; + handshake(&mut c).await; + subscribe(&mut c).await; + + // Hang up the writing half only: the daemon reads the end of the + // requests, and this half stays open to see the daemon close its own. + c.shutdown().await.unwrap(); + let closed = tokio::time::timeout(EVENT_TIMEOUT, async { + while let Ok(Some(_)) = read_frame::(&mut c).await {} + }) + .await; + assert!( + closed.is_ok(), + "the daemon kept a hung-up subscriber's connection open" + ); + + let mut c = connect(&env.sock()).await; + handshake(&mut c).await; + write_frame(&mut c, &ClientRequest::Shutdown).await.unwrap(); + wait_for_exit(&mut daemon); +} + /// `nebula spawn ""` from inside a session, end to end over real /// processes: the CLI (what the model runs) makes the daemon start a second /// agent in the caller's worktree — booted at once, on the default name so @@ -4018,7 +4052,13 @@ async fn nebula_spawn_cli_starts_a_sibling_session_in_the_same_worktree() { ) .unwrap(); make_executable(&script); - let mut daemon = env.spawn_daemon_with_agent_cmd(script.to_str().unwrap()); + // A Codex sibling installs Codex's managed hooks into Codex's home: + // this test's own, never the developer's `~/.codex`. + let codex_home = env.tmp.path().join("codex-home"); + let mut daemon = env.spawn_daemon_with( + script.to_str().unwrap(), + &[(env::CODEX_HOME, codex_home.to_str().unwrap())], + ); let mut c = connect(&env.sock()).await; handshake(&mut c).await;