use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use nomoreide_core::remote::connector::EventSender;
use nomoreide_core::remote::protocol::device_bound::{
TerminalAttachRequest, TerminalDetach, TerminalInput, TerminalResize, TerminalSpawnRequest,
};
use nomoreide_core::remote::protocol::errors::{ErrorCode, ProtocolError};
use nomoreide_core::remote::protocol::limits;
use nomoreide_core::remote::protocol::platform_bound::{
TerminalAck, TerminalAttachAccepted, TerminalCloseReason, TerminalClosed, TerminalGeometry,
TerminalKilled, TerminalOutput, TerminalSessionsResponse, TerminalSpawned,
};
use nomoreide_core::remote::protocol::snapshot::RemoteTerminalSession;
use nomoreide_core::remote::protocol::PlatformBound;
use nomoreide_core::remote::protocol::TerminalBytes;
use nomoreide_core::terminal::TerminalManager;
use tokio::sync::broadcast::error::RecvError;
#[derive(Clone, Default)]
pub(crate) struct Mirrors {
open: Arc<Mutex<HashMap<String, Mirror>>>,
}
struct Mirror {
session_id: String,
_cancel: tokio::sync::oneshot::Sender<()>,
}
impl Mirrors {
pub(crate) fn attach(
&self,
terminal: &TerminalManager,
request: &TerminalAttachRequest,
events: EventSender,
) -> Result<PlatformBound, ProtocolError> {
if !terminal.is_mirrorable(&request.session_id, super::shell_allowed()) {
return Err(ProtocolError::new(
ErrorCode::CapabilityUnavailable,
"That is not a terminal this machine will mirror.",
)
.with_detail(request.session_id.clone()));
}
let mut open = self.open.lock().unwrap();
if open.len() >= limits::MAX_TERMINAL_STREAMS {
return Err(ProtocolError::new(
ErrorCode::CapabilityUnavailable,
"Too many terminals are already mirrored from this machine.",
));
}
let Some(mirror) = terminal.mirror_output(&request.session_id) else {
return Err(ProtocolError::new(
ErrorCode::CapabilityUnavailable,
"That terminal is no longer running.",
)
.with_detail(request.session_id.clone()));
};
let (cols, rows) = mirror.size;
let stream_id = format!("stream_{}", uuid::Uuid::new_v4());
let (cancel, cancelled) = tokio::sync::oneshot::channel();
open.insert(
stream_id.clone(),
Mirror {
session_id: request.session_id.clone(),
_cancel: cancel,
},
);
drop(open);
tokio::spawn(pump(
stream_id.clone(),
mirror,
cancelled,
events,
self.clone(),
));
Ok(PlatformBound::TerminalAttachAccepted(
TerminalAttachAccepted {
stream_id,
session_id: request.session_id.clone(),
cols,
rows,
},
))
}
pub(crate) fn input(
&self,
terminal: &TerminalManager,
request: &TerminalInput,
) -> Result<PlatformBound, ProtocolError> {
if request.data.len() > limits::MAX_TERMINAL_INPUT_BYTES {
return Err(ProtocolError::new(
ErrorCode::MalformedFrame,
"That is more input than one frame may carry.",
));
}
let session_id = self.session_for(&request.stream_id)?;
terminal
.write_input(&session_id, request.data.as_slice())
.map_err(|reason| {
ProtocolError::new(ErrorCode::CapabilityUnavailable, "That terminal is gone.")
.with_detail(reason)
})?;
Ok(PlatformBound::TerminalAck(TerminalAck {
stream_id: request.stream_id.clone(),
}))
}
pub(crate) fn resize(
&self,
terminal: &TerminalManager,
request: &TerminalResize,
) -> Result<PlatformBound, ProtocolError> {
let session_id = self.session_for(&request.stream_id)?;
let (cols, rows) = terminal.session_size(&session_id).unwrap_or((80, 24));
Ok(PlatformBound::TerminalAttachAccepted(
TerminalAttachAccepted {
stream_id: request.stream_id.clone(),
session_id,
cols,
rows,
},
))
}
pub(crate) fn detach(&self, request: &TerminalDetach) -> Result<PlatformBound, ProtocolError> {
self.close(&request.stream_id);
Ok(PlatformBound::TerminalClosed(TerminalClosed {
stream_id: request.stream_id.clone(),
reason: TerminalCloseReason::Detached,
}))
}
pub(crate) fn close_all(&self) {
self.open.lock().unwrap().clear();
}
fn close(&self, stream_id: &str) {
self.open.lock().unwrap().remove(stream_id);
}
fn session_for(&self, stream_id: &str) -> Result<String, ProtocolError> {
self.open
.lock()
.unwrap()
.get(stream_id)
.map(|mirror| mirror.session_id.clone())
.ok_or_else(|| {
ProtocolError::new(
ErrorCode::CapabilityUnavailable,
"That terminal is not mirrored.",
)
.with_detail(stream_id.to_string())
})
}
}
pub(crate) fn describe(
session: nomoreide_core::terminal::TerminalSession,
waiting: bool,
) -> RemoteTerminalSession {
RemoteTerminalSession {
id: session.id,
label: session.label,
provider: session.provider,
workspace: std::path::Path::new(&session.cwd)
.file_name()
.map(|name| name.to_string_lossy().into_owned()),
running: session.exit.is_none(),
started_at: session.started_at,
waiting,
}
}
pub(crate) fn spawned(session: nomoreide_core::terminal::TerminalSession) -> PlatformBound {
PlatformBound::TerminalSpawned(TerminalSpawned {
session: describe(session, false),
})
}
pub(crate) fn killed(session_id: String) -> PlatformBound {
PlatformBound::TerminalKilled(TerminalKilled { session_id })
}
pub(crate) fn check_prompt(request: &TerminalSpawnRequest) -> Result<(), ProtocolError> {
if request.prompt.trim().is_empty() {
return Err(ProtocolError::new(
ErrorCode::MalformedFrame,
"An agent needs something to work on.",
));
}
if request.prompt.len() > limits::MAX_AGENT_PROMPT_BYTES {
return Err(ProtocolError::new(
ErrorCode::MalformedFrame,
"That prompt is larger than one frame may carry.",
));
}
Ok(())
}
pub(crate) fn sessions(terminal: &TerminalManager) -> PlatformBound {
PlatformBound::TerminalSessions(TerminalSessionsResponse {
sessions: terminal
.mirrorable_sessions(super::shell_allowed())
.into_iter()
.map(|session| {
let waiting = terminal.awaiting_choice(&session.id);
describe(session, waiting)
})
.collect(),
})
}
async fn pump(
stream_id: String,
mirror: nomoreide_core::terminal::TerminalMirror,
mut cancelled: tokio::sync::oneshot::Receiver<()>,
events: EventSender,
mirrors: Mirrors,
) {
let nomoreide_core::terminal::TerminalMirror {
replay,
mut output,
size: _,
mut resized,
} = mirror;
let mut seq = 0u64;
let mut pending: Vec<u8> = replay;
let reason = loop {
while !pending.is_empty() {
let take = pending.len().min(limits::MAX_TERMINAL_CHUNK_BYTES);
let chunk: Vec<u8> = pending.drain(..take).collect();
let frame = PlatformBound::TerminalOutput(TerminalOutput {
stream_id: stream_id.clone(),
seq,
data: TerminalBytes::new(chunk),
});
seq += 1;
if events.send(frame).await.is_err() {
break;
}
}
tokio::select! {
_ = &mut cancelled => break TerminalCloseReason::Detached,
changed = resized.changed() => match changed {
Ok(()) => {
let (cols, rows) = *resized.borrow_and_update();
let frame = PlatformBound::TerminalGeometry(TerminalGeometry {
stream_id: stream_id.clone(),
cols,
rows,
});
if events.send(frame).await.is_err() {
break TerminalCloseReason::Detached;
}
}
Err(_) => break TerminalCloseReason::Exited,
},
received = output.recv() => match received {
Ok(data) => {
pending.extend_from_slice(&data);
tokio::time::sleep(limits::TERMINAL_COALESCE_INTERVAL).await;
while let Ok(more) = output.try_recv() {
pending.extend_from_slice(&more);
}
}
Err(RecvError::Closed) => break TerminalCloseReason::Exited,
Err(RecvError::Lagged(_)) => break TerminalCloseReason::Overrun,
},
}
};
mirrors.close(&stream_id);
let _ = events
.send(PlatformBound::TerminalClosed(TerminalClosed {
stream_id,
reason,
}))
.await;
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn a_resize_on_the_machine_is_sent_to_the_mirror() {
let terminal = TerminalManager::new();
let session = spawn_agent(&terminal, "geometry-agent");
let (events, mut received) = tokio::sync::mpsc::channel(16);
let mirrors = Mirrors::default();
let accepted = mirrors
.attach(
&terminal,
&TerminalAttachRequest {
session_id: session.clone(),
cols: 40,
rows: 20,
},
events,
)
.expect("attach");
let PlatformBound::TerminalAttachAccepted(accepted) = accepted else {
panic!("attach must answer with the geometry it will be drawing");
};
assert_eq!((accepted.cols, accepted.rows), (80, 24));
terminal.resize(&session, 132, 43).expect("resize");
let geometry = wait_for_geometry(&mut received).await;
assert_eq!(geometry.stream_id, accepted.stream_id);
assert_eq!((geometry.cols, geometry.rows), (132, 43));
terminal.close_session(&session).unwrap();
}
async fn wait_for_geometry(
received: &mut tokio::sync::mpsc::Receiver<PlatformBound>,
) -> TerminalGeometry {
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(5);
loop {
let frame = tokio::time::timeout_at(deadline, received.recv())
.await
.expect("a geometry frame within five seconds")
.expect("the pump is still running");
match frame {
PlatformBound::TerminalGeometry(geometry) => return geometry,
PlatformBound::TerminalOutput(_) => continue,
other => panic!("unexpected frame while waiting: {}", other.kind()),
}
}
}
fn spawn_agent(terminal: &TerminalManager, id: &str) -> String {
terminal
.create(
std::sync::Arc::new(SilentSink),
nomoreide_core::terminal::TerminalSpawnSpec {
id: id.to_string(),
service_name: None,
cwd: std::env::temp_dir().to_string_lossy().into_owned(),
shell: "/bin/sh".into(),
args: vec!["-c".to_string(), "sleep 30".to_string()],
env: Vec::new(),
label: None,
kind: Some("agent".to_string()),
provider: Some("claude".to_string()),
},
)
.expect("spawn")
.id
}
struct SilentSink;
impl nomoreide_core::event_sink::EventSink for SilentSink {
fn emit(
&self,
_event: &str,
_payload: serde_json::Value,
) -> Result<(), nomoreide_core::event_sink::EventSinkError> {
Ok(())
}
}
#[test]
fn the_session_listing_carries_the_start_time() {
let terminal = TerminalManager::new();
let id = spawn_agent(&terminal, "uptime-agent");
let PlatformBound::TerminalSessions(response) = sessions(&terminal) else {
panic!("expected a session listing");
};
let session = response
.sessions
.iter()
.find(|session| session.id == id)
.expect("the session just spawned");
assert!(
session.started_at.is_some(),
"the listing must carry when the session started"
);
assert!(!session.waiting);
}
}