Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 16 additions & 2 deletions crates/nebula-daemon/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,8 @@ async fn handle_client(daemon: Arc<Daemon>, stream: UnixStream) -> Result<()> {

// Per-connection attach state: forward-task handles keyed by session.
let mut attached: HashMap<SessionRef, tokio::task::JoinHandle<()>> = HashMap::new();
// The task forwarding broadcasts to a subscribed client.
let mut subscription: Option<tokio::task::JoinHandle<()>> = None;
let mut handshaken = false;

let result: Result<()> = async {
Expand Down Expand Up @@ -91,6 +93,10 @@ async fn handle_client(daemon: Arc<Daemon>, 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![],
Expand All @@ -101,9 +107,8 @@ async fn handle_client(daemon: Arc<Daemon>, 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) => {
Expand All @@ -118,6 +123,9 @@ async fn handle_client(daemon: Arc<Daemon>, stream: UnixStream) -> Result<()> {
}
}
});
if let Some(old) = subscription.replace(forward) {
old.abort();
}
}
ClientRequest::Attach {
session: sref,
Expand Down Expand Up @@ -592,6 +600,12 @@ async fn handle_client(daemon: Arc<Daemon>, 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
Expand Down
46 changes: 43 additions & 3 deletions crates/nebula/tests/e2e_pty.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -3994,6 +3997,37 @@ async fn subscribe(c: &mut UnixStream) -> Vec<ServerEvent> {
.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::<ServerEvent, _>(&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 "<task>"` 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
Expand All @@ -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;
Expand Down