//! The live session registry. Every session mutation -- spawn, delete, //! token changes -- funnels through [`SessionManager`] under one lock, so //! in-memory state and `config.ron` can't come apart (the same pattern as //! dev-updater's `registry.rs`). //! //! A live session is a driver plus one event pump: the driver reports //! [`Event`]s into an mpsc channel; the pump assigns each a sequence //! number, appends it to the session's transcript file, and fans it out to //! SSE subscribers. The transcript is the source of truth -- subscribers //! that fall behind or reconnect catch up from the file by cursor. pub mod claude; pub mod driver; pub mod echo; pub mod import; pub mod llama; pub mod process; pub mod transcript; pub mod transport; use std::collections::HashMap; use std::path::{Path, PathBuf}; use std::sync::{Arc, Mutex, RwLock}; use std::time::{SystemTime, UNIX_EPOCH}; use anyhow::{Context, Result, bail}; use serde::Serialize; use tokio::sync::{broadcast, mpsc}; use crate::config::{ Config, DriverKind, ProviderConfig, SessionConfig, SetupConfig, SshConfig, TokenEntry, }; use claude::ClaudeDriver; use driver::{Driver, Event, EventSink, ImageRef, SessionStatus}; use echo::EchoDriver; use llama::LlamaDriver; use transcript::{SeqEvent, Transcript}; use transport::Transport; /// Fan-out buffer per session. A subscriber that falls further behind than /// this is caught up from the transcript file instead (see `routes`), so /// the size only bounds memory, not correctness. const EVENT_BUFFER: usize = 256; pub fn now() -> f64 { SystemTime::now() .duration_since(UNIX_EPOCH) .unwrap_or_default() .as_secs_f64() } /// What the phone needs to spawn a session -- the spawn screen's fields. pub struct SpawnSpec { /// Which machine, and which of its providers. pub setup: String, pub provider: String, pub title: Option, pub model: Option, pub cwd: Option, pub permission_mode: Option, /// Driver-interpreted settings; see `SessionConfig::params`. pub params: std::collections::BTreeMap, } /// One row of `GET /sessions`. #[derive(Debug, Clone, Serialize)] #[serde(rename_all = "camelCase")] pub struct SessionInfo { pub id: String, pub provider: String, /// Id of the machine it runs on, which is what the session stored. pub setup: String, /// That machine's current label, resolved when this row is built -- /// so renaming a setup renames it everywhere it appears, rather than /// leaving old sessions showing the old name. pub setup_name: String, pub title: String, #[serde(skip_serializing_if = "Option::is_none")] pub model: Option, /// How much this session asks before acting. Reported so the phone can /// *show* the current mode rather than assume one -- a picker that /// guesses its own value is how you end up changing something you /// thought you were confirming. #[serde(skip_serializing_if = "Option::is_none")] pub permission_mode: Option, /// Whether this session continues one the machine already had. /// /// Reported because it changes what deleting *means*: an imported /// session's real transcript belongs to the CLI and survives, so /// removing it here is undoing a view. A session started here has no /// copy anywhere else, and removing it ends the conversation. Saying /// "this cannot be undone" of both makes the warning worthless on the /// one where it is true. pub imported: bool, #[serde(skip_serializing_if = "Option::is_none")] pub cwd: Option, pub status: SessionStatus, pub last_activity: f64, pub created: f64, } /// A running session: its driver plus the shared state the event pump /// keeps current. Cheap to clone-by-`Arc` into request handlers. pub struct LiveSession { meta: SessionConfig, driver: Box, /// The same channel the driver reports into; the manager injects /// `UserMessage`/`Answered` here so they take a sequence number in /// order with everything else. sink: mpsc::UnboundedSender, events: broadcast::Sender, transcript_path: PathBuf, shared: Arc, } /// The pump-maintained view of a session, read by the list endpoint. /// `model` also lives here (not in the immutable meta) because it can /// change mid-session via `set_model`. struct Shared { status: Mutex, last_activity: Mutex, model: Mutex>, /// Beside the model and for the same reason: `meta` is the shape the /// session was *launched* with, so reporting from it would show the /// mode a change had already replaced. permission_mode: Mutex>, /// How many events this session has ever recorded. /// /// Only the import sync reads it, and only to answer one question: /// "did *we* write anything since I last looked?" A session and a /// terminal append to the same file, so that is the whole of what /// separates lines worth replaying from lines already shown. Status /// cannot answer it -- a turn that starts and finishes between two /// polls is idle at both, and its output then gets replayed on top of /// itself. written: Mutex, } impl LiveSession { /// Hands the user's message to the driver, which records it in the /// transcript by reporting that it has taken it -- see `MessageTaken`. /// /// The message is deliberately not recorded here. Sent into a running /// turn it waits, and writing it down on the way past would put it /// above output that happened before the session ever saw it. pub fn send_message(&self, text: String, images: Vec) { // Attachments are the exception, recorded on the way past: they // are uploaded whether or not the message waits, and the phone // fetches them by the same ref the files route serves. So a queued // message's picture appears a little before its text. for image in &images { let _ = self.sink.send(Event::Image { image: image.clone(), }); } self.driver.send_user_message(text, images); } pub fn answer_question(&self, question_id: &str, answer: &str) { let _ = self.sink.send(Event::Answered { id: question_id.to_string(), answer: answer.to_string(), }); self.driver.answer_question(question_id, answer); } pub fn interrupt(&self) { self.driver.interrupt(); } /// Leaves this session's process running and stops attending to it, /// for a server that is going away and means to come back. See /// [`Driver::detach`]. pub fn detach(&self) { self.driver.detach(); } pub fn compact(&self) { self.driver.compact(); } pub fn subscribe(&self) -> broadcast::Receiver { self.events.subscribe() } pub fn transcript_path(&self) -> &Path { &self.transcript_path } /// The session's directory (attachments in, produced files out live in /// `attachments/` and `files/` under it). pub fn dir(&self) -> &Path { self.transcript_path .parent() .expect("transcript lives in the session dir") } /// Stores one uploaded attachment, returning the id `POST /message` /// references it by. Removed with the session directory on delete -- /// the same path out as everything else in it. pub fn save_attachment(&self, bytes: &[u8], content_type: &str) -> Result { // An unrecognized type is almost always a phone photo whose // content type the picker didn't set; jpg is the useful guess. let extension = crate::media::extension_for(content_type).unwrap_or("jpg"); let name = format!("{}.{extension}", random_hex()); let dir = self.dir().join("attachments"); wg_app_link::private::create_dir(&dir)?; std::fs::write(dir.join(&name), bytes) .with_context(|| format!("write attachment {name}"))?; Ok(name) } /// `setup_name` is passed in rather than stored: only the manager /// holds the config, and the label can change under a running session. fn info(&self, setup_name: &str, imported: bool) -> SessionInfo { SessionInfo { id: self.meta.id.clone(), provider: self.meta.provider.clone(), setup: self.meta.setup.clone(), setup_name: setup_name.to_string(), title: self.meta.title.clone(), model: self.shared.model.lock().unwrap().clone(), permission_mode: self.shared.permission_mode.lock().unwrap().clone(), imported, cwd: self.meta.cwd.clone(), status: *self.shared.status.lock().unwrap(), last_activity: *self.shared.last_activity.lock().unwrap(), created: self.meta.created, } } } struct Inner { config: Config, live: HashMap>, } pub struct SessionManager { config_path: PathBuf, /// Per-session directories (transcript, attachments, produced images) /// live under here, each named by session id. data_dir: PathBuf, /// Downloaded GGUF models, shared by every session that names one -- /// which is why they live beside the session directories rather than /// inside one. models_dir: PathBuf, inner: RwLock, } impl SessionManager { /// Loads the config and relaunches a driver for every persisted /// session -- for the real drivers that is the `--resume`/session-file /// crash-recovery story; the echo driver just starts fresh over the /// same transcript. Must be called inside a tokio runtime (each /// session spawns its event pump). pub fn new(config_path: PathBuf, data_dir: PathBuf, models_dir: PathBuf) -> Result { let config = Config::load(&config_path)?; wg_app_link::private::create_dir(&data_dir)?; let mut live = HashMap::new(); for meta in &config.sessions { // One unlaunchable session -- a corrupt transcript, an // unreachable ssh host, a provider that was edited away -- // shows as exited rather than taking the whole server down // with it, and can still be deleted from the phone. match resolve(&config, meta).and_then(|(setup, provider)| { launch( meta.clone(), &setup, &provider, &data_dir, &models_dir, None, ) }) { Ok(session) => { live.insert(meta.id.clone(), session); } Err(err) => { tracing::error!("couldn't relaunch session {}: {err:#}", meta.id); } } } let manager = Self { config_path, data_dir, models_dir, inner: RwLock::new(Inner { config, live }), }; Ok(manager) } /// Writes this machine into a config that has no setups, with the /// providers actually found on it. /// /// Discovered rather than assumed. Until 2026-08-28 this wrote a /// `claude-cli` provider unconditionally, so a fresh install on a /// machine without `claude` -- which is every machine but the dev VM /// -- offered a spawn option that could not work, and said so with the /// same confidence as a provider that had been checked for. Providers /// are discovered by asking the machine, and the local machine is not /// an exception to that. /// /// A discovery that fails seeds only `echo`, which is true wherever /// this server runs, and says so in the log. Seeding the hardcoded /// list on failure would be the original bug with an extra step, and /// seeding nothing would leave a fresh install with nothing to prove /// the pipe with. pub async fn seed_setup(&self) -> Result<()> { if !self.inner.read().unwrap().config.setups.is_empty() { return Ok(()); } let providers = match crate::setups::discover(&transport::Transport::Here).await { Ok(found) => found, Err(err) => { tracing::warn!( "couldn't ask this machine what it has ({err}); seeding {} only -- \ re-probe the setup from the app once that is fixed", crate::config::ECHO_PROVIDER ); vec![Config::echo_provider()] } }; let names: Vec<&str> = providers.iter().map(|p| p.name.as_str()).collect(); let mut inner = self.inner.write().unwrap(); if !inner.config.setups.is_empty() { return Ok(()); } let mut candidate = inner.config.clone(); candidate.setups.push(Config::seed(providers.clone())); candidate.save(&self.config_path)?; inner.config = candidate; tracing::info!( "no setups configured -- added \"{}\" with: {}", crate::config::LOCAL_SETUP, names.join(", ") ); Ok(()) } /// The one path by which the config changes. /// /// Clone, apply, save, and only then commit: a failed write leaves /// what was already there and reports why, so what this server /// believes and what is on disk cannot come apart. The ordering is /// the whole trick -- mutating in place and then saving would leave a /// server that had accepted a change nothing on disk records. fn update(&self, apply: impl FnOnce(&mut Config) -> Result) -> Result { let mut inner = self.inner.write().unwrap(); let mut candidate = inner.config.clone(); let outcome = apply(&mut candidate)?; candidate.save(&self.config_path)?; inner.config = candidate; Ok(outcome) } /// Adds a machine with the providers it was found to have. /// /// `providers` comes from probing rather than from the caller (see /// `crate::setups`), which is why this takes them as an argument: the /// probe is async and this is not, so the route does the asking and /// this does the writing. pub fn add_setup( &self, name: &str, ssh: Option, providers: Vec, ) -> Result { let name = name.trim().to_string(); if name.is_empty() { bail!("a setup needs a name"); } self.update(|config| { if config.setup_named(&name).is_some() { bail!("there is already a setup called \"{name}\""); } // Ids are derived once and then fixed, so a label can be // edited later without orphaning the sessions that named it. let mut id = crate::setups::id_from(&name); while config.setup(&id).is_some() { id = format!("{id}-{}", &random_hex()[..4]); } let setup = SetupConfig { id, name: name.clone(), ssh, providers, }; config.setups.push(setup.clone()); Ok(setup) }) } /// Renames a machine, or replaces what was discovered on it. pub fn update_setup( &self, id: &str, name: Option<&str>, providers: Option>, ) -> Result { self.update(|config| { if let Some(name) = name { let name = name.trim(); if name.is_empty() { bail!("a setup needs a name"); } if config.setups.iter().any(|s| s.name == name && s.id != id) { bail!("there is already a setup called \"{name}\""); } } let setup = config .setups .iter_mut() .find(|setup| setup.id == id) .with_context(|| format!("no setup with id \"{id}\""))?; if let Some(name) = name { setup.name = name.trim().to_string(); } if let Some(providers) = providers { setup.providers = providers; } Ok(setup.clone()) }) } /// Removes a machine, provided nothing is still running on it. /// /// Refused rather than cascaded: deleting a machine should not /// silently kill conversations, and the person asking is better placed /// to decide which of those sessions they still want. pub fn delete_setup(&self, id: &str) -> Result<()> { self.update(|config| { if config.setup(id).is_none() { bail!("no setup with id \"{id}\""); } let using: Vec<&str> = config .sessions .iter() .filter(|session| session.setup == id) .map(|session| session.title.as_str()) .collect(); if !using.is_empty() { bail!( "{} session(s) still run on it: {}. Delete them first.", using.len(), using.join(", "), ); } config.setups.retain(|setup| setup.id != id); Ok(()) }) } pub fn tokens(&self) -> Vec { self.inner.read().unwrap().config.tokens.clone() } /// Replaces the enrolled token list. With one device this is rotation: /// the old hash is invalidated the moment the new config is saved. pub fn set_tokens(&self, tokens: Vec) -> Result<()> { let mut inner = self.inner.write().unwrap(); let mut candidate = inner.config.clone(); candidate.tokens = tokens; candidate.save(&self.config_path)?; inner.config = candidate; Ok(()) } /// Lets go of every session's process, for a server that is going /// away and means to adopt them again when it comes back. /// /// Deliberately not a shutdown, and this is the load-bearing half of /// it: a backend restart -- a rebuild, a service restart, a crash -- /// must not end a turn somebody is waiting on. Each process keeps its /// record in the session directory, and `launch` finds it there rather /// than starting a second one against the same conversation. /// /// What this did before was ask them all to stop and then exit /// immediately, which stopped nothing reliably -- the grace timer died /// with the runtime -- and orphaned whatever survived with nothing /// written down to find it by. Processes leaked either way; what is /// different now is that they are left on purpose and can be picked /// back up. pub fn detach_all(&self) { let inner = self.inner.read().unwrap(); for session in inner.live.values() { session.detach(); } tracing::info!( "left {} session process(es) running to be reattached to", inner.live.len() ); } /// The session already continuing `source`, if there is one. /// /// Importing the same Claude Code session twice would leave two /// `--resume` processes appending to one transcript, each seeing the /// other's writes as work done elsewhere and replaying them. Nothing /// is corrupted, but both sessions show a conversation neither of them /// is having, which is worse than a refusal. pub fn session_importing(&self, source: &str) -> Option { let inner = self.inner.read().unwrap(); inner.config.sessions.iter().find_map(|meta| { let cursor = import::read_cursor(&self.data_dir.join(&meta.id))?; cursor .path .rsplit('/') .next()? .strip_suffix(".jsonl") .filter(|found| *found == source) .map(|_| meta.id.clone()) }) } /// Every session, in config order, with live status joined in. A /// session that failed to relaunch reports as exited. pub fn sessions(&self) -> Vec { let inner = self.inner.read().unwrap(); inner .config .sessions .iter() .map(|meta| match inner.live.get(&meta.id) { Some(session) => session.info( label_of(&inner.config, &meta.setup), import::read_cursor(&self.data_dir.join(&meta.id)).is_some(), ), None => SessionInfo { id: meta.id.clone(), setup: meta.setup.clone(), setup_name: label_of(&inner.config, &meta.setup).to_string(), provider: meta.provider.clone(), title: meta.title.clone(), model: meta.model.clone(), permission_mode: meta.permission_mode.clone(), imported: import::read_cursor(&self.data_dir.join(&meta.id)).is_some(), cwd: meta.cwd.clone(), status: SessionStatus::Exited, last_activity: meta.created, created: meta.created, }, }) .collect() } pub fn session(&self, id: &str) -> Option> { self.inner.read().unwrap().live.get(id).cloned() } /// Every provider this server offers, built-in echo included. /// Every machine this server can run something on, each with what it /// can run. One list rather than two, because the pair is the choice. pub fn setups(&self) -> Vec { self.inner.read().unwrap().config.setups.clone() } pub fn spawn_session(&self, spec: SpawnSpec) -> Result { self.spawn_seeded(spec, None) } /// Spawns a session that continues one the machine already had. /// /// The same path as any other spawn, with a [`Seed`] written into the /// session directory before the driver starts -- which is all an /// import is, because `claude.rs` already resumes when it finds a /// resume token. A separate spawn path would be a second way to start /// a session, and the driver would have to learn which one it was. pub fn spawn_imported(&self, spec: SpawnSpec, seed: Seed) -> Result { self.spawn_seeded(spec, Some(seed)) } fn spawn_seeded(&self, spec: SpawnSpec, seed: Option) -> Result { let mut inner = self.inner.write().unwrap(); let setup = inner .config .setup(&spec.setup) .with_context(|| { format!( "no setup with id \"{}\" -- configured: {}", spec.setup, // Ids, since that is what was looked up. Listing the // labels made the failure read as a contradiction: // "no setup named X -- configured: X", when X was a // label and the id was something else. names(inner.config.setups.iter().map(|s| s.id.as_str())), ) })? .clone(); let provider = setup .provider(&spec.provider) .with_context(|| { format!( "setup \"{}\" has no provider named \"{}\" -- it offers: {}", spec.setup, spec.provider, names(setup.providers.iter().map(|p| p.name.as_str())), ) })? .clone(); let id = unique_id(&inner.config); let title = spec .title .filter(|title| !title.trim().is_empty()) .unwrap_or_else(|| format!("{} session", provider.name)); let meta = SessionConfig { id: id.clone(), setup: setup.id.clone(), provider: provider.name.clone(), title, // No model unless one was chosen. This used to fall back to // the provider's first listed model, which sounds like a // default and is not one: that list is a shortcut for the // spawn screen, written in whatever order somebody typed it, // and its first entry happened to be `fable`. Every session // spawned without a model -- every import, since importing // asks for none -- silently became a fable session. Absent // means absent, and the CLI then uses whatever the person // configured for themselves. model: spec.model, cwd: spec.cwd, permission_mode: spec.permission_mode, params: spec.params, created: now(), }; let session = launch( meta.clone(), &setup, &provider, &self.data_dir, &self.models_dir, seed, )?; let mut candidate = inner.config.clone(); candidate.sessions.push(meta); if let Err(err) = candidate.save(&self.config_path) { // The path out of everything the launch created, taken in the // same change: drop the session and its directory so a failed // save leaves no orphan. drop(session); let _ = std::fs::remove_dir_all(self.data_dir.join(&id)); return Err(err); } inner.config = candidate; // Whether this one was seeded, which is the same question the // listing asks of the directory a moment later. let info = session.info( &setup.name, import::read_cursor(&self.data_dir.join(&id)).is_some(), ); inner.live.insert(id, session); Ok(info) } /// Changes a session's model: persisted (so a respawn keeps it and the /// list shows it) and handed to the driver, which switches in place /// where its dialect can. Through the manager, not the session, so the /// config and the live view can't disagree. /// Changes how much a session asks before acting, live and persisted. /// /// Alongside the model rather than folded into it: they are set at the /// same moment and by the same screen, but they answer different /// questions, and a caller changing one must not have to restate the /// other. pub fn set_session_permission_mode(&self, id: &str, mode: &str) -> Result<()> { let mut inner = self.inner.write().unwrap(); if !inner.config.sessions.iter().any(|meta| meta.id == id) { bail!("no session {id}"); } let mut candidate = inner.config.clone(); for meta in candidate.sessions.iter_mut().filter(|meta| meta.id == id) { meta.permission_mode = Some(mode.to_string()); } candidate.save(&self.config_path)?; inner.config = candidate; if let Some(session) = inner.live.get(id) { *session.shared.permission_mode.lock().unwrap() = Some(mode.to_string()); session.driver.set_permission_mode(mode); } Ok(()) } pub fn set_session_model(&self, id: &str, model: &str) -> Result<()> { let mut inner = self.inner.write().unwrap(); if !inner.config.sessions.iter().any(|meta| meta.id == id) { bail!("no session {id}"); } let mut candidate = inner.config.clone(); for meta in candidate.sessions.iter_mut().filter(|meta| meta.id == id) { meta.model = Some(model.to_string()); } candidate.save(&self.config_path)?; inner.config = candidate; if let Some(session) = inner.live.get(id) { *session.shared.model.lock().unwrap() = Some(model.to_string()); session.driver.set_model(model); } Ok(()) } /// Kills the process, releases everything the spawn created, and /// deletes the transcript and files -- the complete path out. pub fn delete_session(&self, id: &str) -> Result<()> { let mut inner = self.inner.write().unwrap(); if !inner.config.sessions.iter().any(|meta| meta.id == id) { bail!("no session {id}"); } let mut candidate = inner.config.clone(); candidate.sessions.retain(|meta| meta.id != id); candidate.save(&self.config_path)?; inner.config = candidate; if let Some(session) = inner.live.remove(id) { // Stopped, not detached: this is the one exit where the // process must not survive, because the conversation it // belongs to is being removed. See `Driver::stop`. session.driver.stop(); } let dir = self.data_dir.join(id); if dir.exists() { std::fs::remove_dir_all(&dir).with_context(|| format!("remove {}", dir.display()))?; } Ok(()) } } /// The provider and host a session's config names, or a message saying /// which one is missing. Both are looked up fresh at every launch, so /// editing either takes effect on the next respawn. fn resolve(config: &Config, meta: &SessionConfig) -> Result<(SetupConfig, ProviderConfig)> { let setup = config .setup(&meta.setup) .with_context(|| format!("no setup named \"{}\"", meta.setup))?; let provider = setup.provider(&meta.provider).with_context(|| { format!( "setup \"{}\" has no provider named \"{}\"", meta.setup, meta.provider ) })?; Ok((setup.clone(), provider.clone())) } /// A setup's current label, or its id when the setup has been deleted -- /// which is what a session left behind by a removed machine shows, and is /// better than an empty column or a guess at what it used to be called. fn label_of<'a>(config: &'a Config, id: &'a str) -> &'a str { config.setup(id).map_or(id, |setup| setup.name.as_str()) } /// Names for a failure message: what there is, so the reader can see what /// they meant instead of only that they were wrong. fn names<'a>(all: impl Iterator) -> String { let all: Vec<_> = all.collect(); if all.is_empty() { "none".to_string() } else { all.join(", ") } } /// 8 random bytes, hex -- short enough for a URL, unique enough forever at /// this scale. pub fn random_hex() -> String { use rand::Rng; let mut bytes = [0u8; 8]; rand::rng().fill_bytes(&mut bytes); bytes.iter().map(|b| format!("{b:02x}")).collect() } /// A [`random_hex`] id not already taken -- checked out of caution. fn unique_id(config: &Config) -> String { loop { let id = random_hex(); if !config.sessions.iter().any(|meta| meta.id == id) { return id; } } } /// Creates the session directory, opens its transcript (continuing the /// sequence numbering if one exists), starts the driver, and spawns the /// event pump connecting them. /// Keeps an imported session's transcript level with the file the CLI /// writes. /// /// Both this app and a terminal append to one file -- `--resume` continues /// the same transcript rather than forking, measured rather than assumed /// -- so the only hard question is which new lines are *ours*. They are /// already in the transcript, having arrived through the driver, and /// replaying them shows every message twice. /// /// Answered by counting what this session has recorded rather than by /// looking at its status. Status is the obvious signal and it is wrong: a /// turn that begins and ends between two polls reads as idle at both, and /// its output is then replayed on top of itself. The count cannot miss /// that, because the events went through the same pump either way. /// /// Its path out: the sink belongs to the session, so once that is dropped /// every send fails and this returns. Nothing else has to remember it. fn spawn_import_sync( transport: Transport, dir: PathBuf, mut cursor: import::Cursor, sink: EventSink, shared: Arc, ) { tokio::spawn(async move { // What the session had recorded when the cursor was last correct. let mut written_at_cursor = *shared.written.lock().unwrap(); loop { tokio::time::sleep(import::SYNC_INTERVAL).await; if sink.is_closed() { return; } let Ok(lines) = import::line_count(&transport, &cursor.path).await else { // A file that cannot be counted is not worth reporting: it // is usually a machine briefly away, and the next poll asks // again. continue; }; let written_now = *shared.written.lock().unwrap(); if written_now != written_at_cursor { // This session produced something since the cursor was set, // so the new lines are its own. Skip them and resynchronise. cursor.lines = lines; written_at_cursor = written_now; import::write_cursor(&dir, &cursor); continue; } if lines <= cursor.lines { continue; } match import::replay_after(&transport, &cursor.path, cursor.lines).await { Ok(events) => { tracing::info!( "{} grew by {} lines with nothing from here; replaying {} events", cursor.path, lines - cursor.lines, events.len(), ); let count = events.len() as u64; for event in events { if sink.send(event).is_err() { return; } } cursor.lines = lines; // The pump is about to record exactly these, so account // for them rather than reading a count that may not have // caught up yet. written_at_cursor = written_now + count; import::write_cursor(&dir, &cursor); } Err(err) => tracing::warn!("couldn't read new lines of {}: {err:#}", cursor.path), } } }); } /// What an imported session starts life with: the token that makes the CLI /// continue rather than begin, and the conversation so far. pub struct Seed { /// The CLI's own session id, written where `claude.rs` looks for it. pub resume: String, /// Where that session's file is and how much of it has been shown, so /// the session can keep itself up to date afterwards. pub cursor: import::Cursor, /// Replayed into the transcript so the phone shows the conversation it /// is joining. The CLI reads the real file itself, so this is what the /// reader sees rather than what the model is given. pub events: Vec, } fn launch( meta: SessionConfig, setup: &SetupConfig, provider: &ProviderConfig, data_dir: &Path, models_dir: &Path, seed: Option, ) -> Result> { let dir = data_dir.join(&meta.id); wg_app_link::private::create_dir(&dir)?; let transcript_path = dir.join("transcript.jsonl"); let mut transcript = Transcript::open(&transcript_path)?; // Before the driver starts, so the token is there when it looks and // the history is already in the transcript a phone will read. if let Some(seed) = seed { claude::write_resume_token(&dir, &seed.resume); import::write_cursor(&dir, &seed.cursor); let at = now(); for event in seed.events { transcript.append(event, at)?; } } let (sink, source) = mpsc::unbounded_channel(); let (events, _) = broadcast::channel(EVENT_BUFFER); let shared = Arc::new(Shared { // What it was last known to be doing, not an assumption. A driver // that has something to say corrects this within its first poll; // one adopting a process that has been quiet says nothing, and // this is then the only true answer available. status: Mutex::new( transcript::last_status(&transcript_path).unwrap_or(SessionStatus::Idle), ), last_activity: Mutex::new(now()), model: Mutex::new(meta.model.clone()), permission_mode: Mutex::new(meta.permission_mode.clone()), written: Mutex::new(0), }); // An imported session shares its transcript file with the CLI -- // `--resume` appends to the same one rather than forking, measured // rather than assumed -- so work done at a terminal belongs in this // session too, and arrives without anybody pressing anything. if let Some(cursor) = import::read_cursor(&dir) { spawn_import_sync( Transport::for_setup(setup), dir.clone(), cursor, sink.clone(), Arc::clone(&shared), ); } let driver: Box = match provider.kind { DriverKind::Echo => Box::new(EchoDriver::new(sink.clone())), DriverKind::LlamaCpp => Box::new(LlamaDriver::launch( &meta, provider, &Transport::for_setup(setup), models_dir, &transcript_path, &dir, sink.clone(), )?), DriverKind::ClaudeCli => Box::new(ClaudeDriver::launch( &meta, provider, &Transport::for_setup(setup), &dir, sink.clone(), )?), }; tokio::spawn(pump( transcript, source, Arc::clone(&shared), events.clone(), )); Ok(Arc::new(LiveSession { meta, driver, sink, events, transcript_path, shared, })) } /// The one writer of a session's transcript: assigns sequence numbers, /// appends, updates the shared status/activity view, fans out. Ends when /// every sender is dropped -- i.e. when the session is deleted and its /// last in-flight task finishes. /// /// The appends are synchronous file writes from an async task, /// deliberately: each is one small line on a local disk, and funneling /// them through one task is what makes the sequence numbering safe. async fn pump( mut transcript: Transcript, mut source: mpsc::UnboundedReceiver, shared: Arc, events: broadcast::Sender, ) { while let Some(event) = source.recv().await { let ts = now(); // Taking a message is how it enters the conversation, and the // conversation is what a phone renders -- so the event becomes the // message here rather than being carried alongside it. One rule // for where a user's message sits: where the session read it. let event = match event { Event::MessageTaken { text } => Event::UserMessage { text }, other => other, }; match transcript.append(event, ts) { Ok(entry) => { if let Event::Status { state } = &entry.event { *shared.status.lock().unwrap() = *state; } *shared.last_activity.lock().unwrap() = ts; *shared.written.lock().unwrap() += 1; // No subscribers is fine; the transcript already has it. let _ = events.send(entry); } Err(err) => tracing::error!("transcript append failed: {err:#}"), } } } #[cfg(test)] mod tests { use super::*; use std::time::Duration; fn echo_spec() -> SpawnSpec { SpawnSpec { params: Default::default(), setup: crate::config::LOCAL_SETUP_ID.to_string(), provider: crate::config::ECHO_PROVIDER.to_string(), title: None, model: None, cwd: None, permission_mode: None, } } /// Reads events from `rx` until `stop` matches one (returning all seen /// so far) or five seconds pass (panicking with what was seen). async fn collect_until( rx: &mut broadcast::Receiver, mut stop: impl FnMut(&Event) -> bool, ) -> Vec { let mut seen = Vec::new(); let deadline = tokio::time::Instant::now() + Duration::from_secs(5); loop { let entry = tokio::time::timeout_at(deadline, rx.recv()) .await .unwrap_or_else(|_| panic!("timed out; events so far: {seen:?}")) .expect("event stream closed"); let done = stop(&entry.event); seen.push(entry); if done { return seen; } } } fn is_idle(event: &Event) -> bool { matches!( event, Event::Status { state: SessionStatus::Idle } ) } /// Collects one full echo turn: everything up to the idle that follows /// the turn's `UsageDelta`. Stopping at the first idle would be racy -- /// the driver emits an idle at construction, and a subscriber attached /// just before the pump processes it would stop there, mid-spawn. async fn collect_turn(rx: &mut broadcast::Receiver) -> Vec { let mut saw_usage = false; collect_until(rx, |event| { saw_usage |= matches!(event, Event::UsageDelta { .. }); saw_usage && is_idle(event) }) .await } /// Writes this machine into `config_path` with echo and nothing else. /// /// Explicit rather than letting the manager seed itself: seeding now /// asks the machine what it has, so a test that relied on it would /// pass or fail depending on whether `claude` happens to be installed /// on whoever is running it. Echo is the only provider that is true /// everywhere, and the only one these tests need. fn seed_echo_only(config_path: &std::path::Path) { Config { setups: vec![Config::seed(vec![Config::echo_provider()])], ..Config::default() } .save(config_path) .expect("seed config"); } #[tokio::test] async fn spawn_message_and_delete_round_trip() { let dir = tempfile::tempdir().expect("tempdir"); let config_path = dir.path().join("config.ron"); let data_dir = dir.path().join("sessions"); seed_echo_only(&config_path); let manager = SessionManager::new( config_path.clone(), data_dir.clone(), data_dir.join("models"), ) .expect("manager"); let info = manager.spawn_session(echo_spec()).expect("spawn"); // Untitled sessions are named after the provider that runs them. assert_eq!(info.title, "echo session"); // Persisted: a fresh load of the config file knows the session. let persisted = Config::load(&config_path).expect("reload config"); assert_eq!(persisted.sessions.len(), 1); assert_eq!(persisted.sessions[0].id, info.id); let session = manager.session(&info.id).expect("live session"); let mut rx = session.subscribe(); session.send_message("hello there".to_string(), Vec::new()); let seen = collect_turn(&mut rx).await; // The user's message is in the stream, before the echoed reply. let user_at = seen .iter() .position(|entry| { matches!(&entry.event, Event::UserMessage { text } if text == "hello there") }) .expect("user message in the stream"); let echoed: String = seen[user_at..] .iter() .filter_map(|entry| match &entry.event { Event::AssistantText { delta } => Some(delta.as_str()), _ => None, }) .collect(); assert_eq!(echoed, "You said: hello there"); // The transcript replays the same events by cursor. let replay = transcript::read_after(session.transcript_path(), 0).expect("replay"); assert!(replay.len() >= seen.len()); let cursor = seen[user_at].seq; let after = transcript::read_after(session.transcript_path(), cursor).expect("replay"); assert_eq!(after.first().map(|entry| entry.seq), Some(cursor + 1)); // Delete is the complete path out: config, registry, and files. manager.delete_session(&info.id).expect("delete"); assert!(manager.sessions().is_empty()); assert!(manager.session(&info.id).is_none()); assert!(!data_dir.join(&info.id).exists()); assert!( Config::load(&config_path) .expect("reload") .sessions .is_empty() ); assert!(manager.delete_session(&info.id).is_err()); } #[tokio::test] async fn questions_round_trip_through_answer() { let dir = tempfile::tempdir().expect("tempdir"); seed_echo_only(&dir.path().join("config.ron")); let manager = SessionManager::new( dir.path().join("config.ron"), dir.path().join("sessions"), dir.path().join("models"), ) .expect("manager"); let info = manager.spawn_session(echo_spec()).expect("spawn"); let session = manager.session(&info.id).expect("live session"); let mut rx = session.subscribe(); session.send_message("/question deploy?".to_string(), Vec::new()); let seen = collect_until(&mut rx, |event| { matches!( event, Event::Status { state: SessionStatus::AwaitingInput } ) }) .await; let question_id = seen .iter() .find_map(|entry| match &entry.event { Event::Question { id, .. } => Some(id.clone()), _ => None, }) .expect("question event"); session.answer_question(&question_id, "Yes"); let seen = collect_until(&mut rx, is_idle).await; assert!(seen.iter().any(|entry| matches!( &entry.event, Event::Answered { id, answer } if *id == question_id && answer == "Yes" ))); } #[tokio::test] async fn a_restart_relaunches_sessions_and_continues_the_numbering() { let dir = tempfile::tempdir().expect("tempdir"); let config_path = dir.path().join("config.ron"); let data_dir = dir.path().join("sessions"); seed_echo_only(&config_path); let manager = SessionManager::new( config_path.clone(), data_dir.clone(), data_dir.join("models"), ) .expect("manager"); let info = manager.spawn_session(echo_spec()).expect("spawn"); let session = manager.session(&info.id).expect("live session"); let mut rx = session.subscribe(); session.send_message("first".to_string(), Vec::new()); let seen = collect_turn(&mut rx).await; let last_seq = seen.last().expect("events").seq; drop(rx); drop(session); drop(manager); // A new manager over the same state: the session is back, and new // events continue the sequence rather than restarting it -- which // is what makes a phone's cursor survive a backend restart. let manager = SessionManager::new(config_path, data_dir.clone(), data_dir.join("models")) .expect("manager restart"); let listed = manager.sessions(); assert_eq!(listed.len(), 1); assert_eq!(listed[0].id, info.id); let session = manager.session(&info.id).expect("relaunched session"); let mut rx = session.subscribe(); session.send_message("second".to_string(), Vec::new()); let seen = collect_turn(&mut rx).await; assert!(seen.first().expect("events").seq > last_seq); } }