use serde::{de::DeserializeOwned, Deserialize, Serialize}; use serde_json::{json, Value}; use std::env; use std::sync::{Arc, Mutex}; use tauri::{AppHandle, Emitter}; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::net::UnixStream; use tokio::time::{sleep, Duration}; #[derive(Serialize)] #[serde(rename_all = "camelCase")] struct BridgeRequest<'a> { version: &'static str, id: &'a str, op: &'a str, #[serde(skip_serializing_if = "Option::is_none")] agent_id: Option<&'a str>, #[serde(skip_serializing_if = "Option::is_none")] payload: Option, } pub fn default_socket_path() -> Result { env::var("PI_STATUS_BRIDGE_SOCKET") .or_else(|_| { env::var("XDG_RUNTIME_DIR") .map(|runtime| format!("{runtime}/pi-status-bridge/bridge.sock")) }) .map_err(|_| "Pi Status Bridge socket is unavailable".to_owned()) } async fn connect_and_send( socket_path: &str, operation: &str, agent_id: Option<&str>, payload: Option, ) -> Result, String> { let mut stream = UnixStream::connect(socket_path) .await .map_err(|error| format!("Could not connect to Pi Status Bridge: {error}"))?; let request = BridgeRequest { version: "v1", id: "tauri-ui", op: operation, agent_id, payload, }; let encoded = serde_json::to_string(&request).map_err(|error| error.to_string())?; stream .write_all(format!("{encoded}\n").as_bytes()) .await .map_err(|error| format!("Could not send bridge request: {error}"))?; Ok(BufReader::new(stream)) } fn result_from_response(value: Value) -> Result { if value.get("ok") != Some(&Value::Bool(true)) { return Err(value .get("error") .and_then(|error| error.get("message")) .and_then(Value::as_str) .unwrap_or("Pi Status Bridge rejected the request") .to_owned()); } value .get("result") .cloned() .ok_or_else(|| "Pi Status Bridge returned no result".to_owned()) } pub async fn request( socket_path: &str, operation: &str, agent_id: Option<&str>, payload: Option, ) -> Result { let mut reader = connect_and_send(socket_path, operation, agent_id, payload).await?; let mut response = String::new(); reader .read_line(&mut response) .await .map_err(|error| format!("Could not read bridge response: {error}"))?; let value: Value = serde_json::from_str(&response) .map_err(|error| format!("Bridge returned invalid JSON: {error}"))?; result_from_response(value) } async fn typed_request( socket_path: &str, operation: &str, payload: Option, ) -> Result { let value = request(socket_path, operation, None, payload).await?; serde_json::from_value(value) .map_err(|error| format!("Bridge returned an invalid {operation} result: {error}")) } #[derive(Clone, Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] pub struct RuntimeSummary { pub runtime_id: String, pub worktree_path: String, pub state: String, pub label: String, pub attention: bool, pub queue_count: u64, pub last_activity: String, pub opened_at: String, pub active_tool: Option, pub agent_id: Option, pub session_id: Option, pub session_path: Option, pub error: Option, } #[derive(Clone, Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] pub struct DirectoryWorkspace { pub worktree_path: String, pub is_home: bool, pub open_count: u64, pub working_count: u64, pub attention_count: u64, pub recovering_count: u64, pub error_count: u64, pub runtimes: Vec, } #[derive(Clone, Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] pub struct Workspace { pub bridge_instance_id: String, pub latest_seq: u64, pub directories: Vec, pub issue: Option, } #[derive(Clone, Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] pub struct WorkspaceSummary { pub bridge_instance_id: String, pub latest_seq: u64, pub open_count: u64, pub working_count: u64, pub attention_count: u64, pub recovering_count: u64, pub error_count: u64, pub directory_count: u64, pub resource_warning: bool, } #[derive(Clone, Debug, Deserialize, Serialize)] pub struct RuntimeResult { pub runtime: RuntimeSummary, } #[derive(Clone, Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] pub struct CloseRuntimeResult { pub runtime_id: String, pub session_path: Option, } #[derive(Clone, Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] pub struct DirectorySession { pub path: String, pub id: String, pub cwd: String, pub name: Option, pub parent_session_path: Option, pub created: Option, pub modified: String, pub message_count: u64, pub first_message: Option, pub is_current: bool, pub runtime_id: Option, } #[derive(Clone, Debug, Deserialize, Serialize)] pub struct DirectorySessionsResult { pub sessions: Vec, } #[derive(Clone, Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] pub struct RuntimeSnapshot { pub bridge_instance_id: String, pub latest_seq: u64, pub runtime: RuntimeSummary, pub state: Option, pub stats: Option, pub transcript: Option, pub commands: Option, pub models: Option, #[serde(default)] pub extensions: Vec, } pub async fn get_workspace(socket_path: &str) -> Result { typed_request(socket_path, "get_workspace", None).await } pub async fn get_workspace_summary(socket_path: &str) -> Result { typed_request(socket_path, "get_workspace_summary", None).await } pub async fn get_model_catalog(socket_path: &str) -> Result { request(socket_path, "get_model_catalog", None, None).await } pub async fn create_session_runtime( socket_path: &str, worktree_path: &str, ) -> Result { typed_request( socket_path, "create_session_runtime", Some(json!({ "worktreePath": worktree_path })), ) .await } pub async fn create_quick_runtime( socket_path: &str, worktree_path: &str, ) -> Result { typed_request( socket_path, "create_quick_runtime", Some(json!({ "worktreePath": worktree_path })), ) .await } pub async fn promote_quick_runtime( socket_path: &str, runtime_id: &str, ) -> Result { typed_request( socket_path, "promote_quick_runtime", Some(json!({ "runtimeId": runtime_id })), ) .await } pub async fn open_session_runtime( socket_path: &str, worktree_path: &str, session_path: &str, ) -> Result { typed_request( socket_path, "open_session_runtime", Some(json!({ "worktreePath": worktree_path, "sessionPath": session_path })), ) .await } pub async fn close_session_runtime( socket_path: &str, runtime_id: &str, ) -> Result { typed_request( socket_path, "close_session_runtime", Some(json!({ "runtimeId": runtime_id })), ) .await } pub async fn close_quick_runtime( socket_path: &str, runtime_id: &str, ) -> Result { typed_request( socket_path, "close_quick_runtime", Some(json!({ "runtimeId": runtime_id })), ) .await } pub async fn list_directory_sessions( socket_path: &str, worktree_path: &str, ) -> Result { typed_request( socket_path, "list_directory_sessions", Some(json!({ "worktreePath": worktree_path })), ) .await } pub async fn get_session_runtime_snapshot( socket_path: &str, runtime_id: &str, ) -> Result { typed_request( socket_path, "get_session_runtime_snapshot", Some(json!({ "runtimeId": runtime_id })), ) .await } pub async fn list_agents(socket_path: &str) -> Result { request(socket_path, "list_agents", None, None).await } pub async fn list_directories(socket_path: &str) -> Result { request(socket_path, "list_directories", None, None).await } pub async fn select_worktree(socket_path: &str, worktree_path: &str) -> Result { request( socket_path, "select_agent", None, Some(json!({ "worktreePath": worktree_path })), ) .await } pub async fn forget_directory(socket_path: &str, worktree_path: &str) -> Result { request( socket_path, "forget_directory", None, Some(json!({ "worktreePath": worktree_path })), ) .await } pub async fn list_sessions(socket_path: &str, agent_id: &str) -> Result { request(socket_path, "list_sessions", Some(agent_id), None).await } pub async fn switch_session( socket_path: &str, agent_id: &str, session_path: &str, ) -> Result { request( socket_path, "switch_session", Some(agent_id), Some(json!({ "sessionPath": session_path })), ) .await } pub async fn new_session(socket_path: &str, agent_id: &str) -> Result { request(socket_path, "new_session", Some(agent_id), None).await } pub async fn submit_prompt( socket_path: &str, agent_id: &str, message: &str, ) -> Result { request( socket_path, "submit_prompt", Some(agent_id), Some(json!({ "message": message })), ) .await } pub async fn abort(socket_path: &str, agent_id: &str) -> Result<(), String> { request(socket_path, "abort", Some(agent_id), None) .await .map(|_| ()) } pub async fn set_model( socket_path: &str, agent_id: &str, provider: &str, model_id: &str, ) -> Result<(), String> { request( socket_path, "set_model", Some(agent_id), Some(json!({ "provider": provider, "modelId": model_id })), ) .await .map(|_| ()) } pub async fn set_thinking_level( socket_path: &str, agent_id: &str, level: &str, ) -> Result<(), String> { request( socket_path, "set_thinking_level", Some(agent_id), Some(json!({ "level": level })), ) .await .map(|_| ()) } pub async fn set_session_name(socket_path: &str, agent_id: &str, name: &str) -> Result<(), String> { request( socket_path, "set_session_name", Some(agent_id), Some(json!({ "name": name })), ) .await .map(|_| ()) } pub async fn compact( socket_path: &str, agent_id: &str, custom_instructions: Option<&str>, ) -> Result<(), String> { request( socket_path, "compact", Some(agent_id), custom_instructions.map(|instructions| json!({ "customInstructions": instructions })), ) .await .map(|_| ()) } pub async fn command(socket_path: &str, agent_id: &str, operation: &str) -> Result<(), String> { request(socket_path, operation, Some(agent_id), None) .await .map(|_| ()) } pub async fn respond_to_extension( socket_path: &str, agent_id: &str, request_id: &str, response: Value, ) -> Result<(), String> { request( socket_path, "extension_response", Some(agent_id), Some(json!({ "requestId": request_id, "response": response })), ) .await .map(|_| ()) } #[derive(Clone, Debug, Deserialize, PartialEq, Serialize)] #[serde(rename_all = "camelCase")] pub struct WorkspaceEvent { pub bridge_instance_id: String, pub seq: u64, #[serde(flatten)] pub body: serde_json::Map, } #[derive(Clone, Debug, Deserialize)] #[serde(rename_all = "camelCase")] struct WorkspaceReplay { bridge_instance_id: String, first_available_seq: u64, latest_seq: u64, truncated: bool, events: Vec, } #[derive(Clone, Debug, Serialize)] #[serde( tag = "kind", rename_all = "camelCase", rename_all_fields = "camelCase" )] pub enum WorkspaceHostEvent { Connected { bridge_instance_id: String, first_available_seq: u64, latest_seq: u64, }, Event { event: WorkspaceEvent, }, ResetRequired { bridge_instance_id: String, latest_seq: u64, reason: String, }, Disconnected { message: String, retry_in_ms: u64, }, } #[derive(Clone, Debug, Default, PartialEq)] pub struct SubscriptionCursor { pub bridge_instance_id: Option, pub seq: u64, } fn accept_replay( cursor: &mut SubscriptionCursor, replay: &WorkspaceReplay, ) -> Vec { if let Some(expected) = &cursor.bridge_instance_id { if expected != &replay.bridge_instance_id { cursor.bridge_instance_id = Some(replay.bridge_instance_id.clone()); cursor.seq = replay.latest_seq; return vec![WorkspaceHostEvent::ResetRequired { bridge_instance_id: replay.bridge_instance_id.clone(), latest_seq: replay.latest_seq, reason: "epochChanged".into(), }]; } } if replay.truncated { cursor.bridge_instance_id = Some(replay.bridge_instance_id.clone()); cursor.seq = replay.latest_seq; return vec![WorkspaceHostEvent::ResetRequired { bridge_instance_id: replay.bridge_instance_id.clone(), latest_seq: replay.latest_seq, reason: "replayTruncated".into(), }]; } cursor.bridge_instance_id = Some(replay.bridge_instance_id.clone()); let mut output = vec![WorkspaceHostEvent::Connected { bridge_instance_id: replay.bridge_instance_id.clone(), first_available_seq: replay.first_available_seq, latest_seq: replay.latest_seq, }]; for event in &replay.events { if event.bridge_instance_id != replay.bridge_instance_id || event.seq <= cursor.seq { continue; } if event.seq != cursor.seq + 1 { cursor.seq = replay.latest_seq; output.push(WorkspaceHostEvent::ResetRequired { bridge_instance_id: replay.bridge_instance_id.clone(), latest_seq: replay.latest_seq, reason: "sequenceGap".into(), }); return output; } cursor.seq = event.seq; output.push(WorkspaceHostEvent::Event { event: event.clone(), }); } if cursor.seq < replay.latest_seq { cursor.seq = replay.latest_seq; output.push(WorkspaceHostEvent::ResetRequired { bridge_instance_id: replay.bridge_instance_id.clone(), latest_seq: replay.latest_seq, reason: "sequenceGap".into(), }); } output } fn accept_live( cursor: &mut SubscriptionCursor, event: WorkspaceEvent, ) -> Option { if cursor.bridge_instance_id.as_deref() != Some(&event.bridge_instance_id) { cursor.bridge_instance_id = Some(event.bridge_instance_id.clone()); cursor.seq = event.seq; return Some(WorkspaceHostEvent::ResetRequired { bridge_instance_id: event.bridge_instance_id, latest_seq: event.seq, reason: "epochChanged".into(), }); } if event.seq <= cursor.seq { return None; } if event.seq != cursor.seq + 1 { cursor.seq = event.seq; return Some(WorkspaceHostEvent::ResetRequired { bridge_instance_id: event.bridge_instance_id, latest_seq: event.seq, reason: "sequenceGap".into(), }); } cursor.seq = event.seq; Some(WorkspaceHostEvent::Event { event }) } #[derive(Debug, Default)] pub struct SubscriptionGeneration(Mutex); impl SubscriptionGeneration { pub fn replace(&self) -> Result { let mut current = self .0 .lock() .map_err(|_| "Could not replace workspace subscription".to_owned())?; *current = current.saturating_add(1); Ok(*current) } pub fn invalidate(&self) -> Result<(), String> { self.replace().map(|_| ()) } fn run_if_current( &self, expected_generation: u64, operation: impl FnOnce() -> Result, ) -> Result { let current = self .0 .lock() .map_err(|_| "Could not inspect workspace subscription".to_owned())?; if *current != expected_generation { return Err("Workspace subscription was replaced".to_owned()); } operation() } fn is_current(&self, expected_generation: u64) -> Result { self.0 .lock() .map(|current| *current == expected_generation) .map_err(|_| "Could not inspect workspace subscription".to_owned()) } } fn retry_delay(attempt: usize) -> u64 { match attempt { 0 => 250, 1 => 500, 2 => 1000, _ => 2000, } } fn emit_workspace( app: &AppHandle, generation: &SubscriptionGeneration, expected_generation: u64, event: WorkspaceHostEvent, ) -> Result<(), String> { generation.run_if_current(expected_generation, || { app.emit("workspace-bridge", event) .map_err(|error| format!("Could not publish workspace event: {error}")) }) } async fn subscribe_workspace_once( socket_path: &str, cursor: &mut SubscriptionCursor, app: &AppHandle, generation: &SubscriptionGeneration, expected_generation: u64, ) -> Result<(), String> { let mut reader = connect_and_send( socket_path, "subscribe_workspace", None, Some(json!({ "cursor": cursor.seq })), ) .await?; let mut handshake = String::new(); if reader .read_line(&mut handshake) .await .map_err(|error| format!("Bridge subscription failed: {error}"))? == 0 { return Err("Bridge subscription closed before handshake".into()); } let replay_value: Value = serde_json::from_str(&handshake) .map_err(|error| format!("Bridge emitted invalid JSON: {error}"))?; let replay: WorkspaceReplay = serde_json::from_value(result_from_response(replay_value)?) .map_err(|error| format!("Bridge returned an invalid workspace replay: {error}"))?; for event in accept_replay(cursor, &replay) { emit_workspace(app, generation, expected_generation, event)?; } loop { let mut line = String::new(); if reader .read_line(&mut line) .await .map_err(|error| format!("Bridge subscription failed: {error}"))? == 0 { return Err("Bridge subscription closed".into()); } let envelope: Value = serde_json::from_str(&line) .map_err(|error| format!("Bridge emitted invalid JSON: {error}"))?; let event: WorkspaceEvent = serde_json::from_value( envelope .get("event") .cloned() .ok_or_else(|| "Bridge event envelope is missing event".to_owned())?, ) .map_err(|error| format!("Bridge emitted an invalid workspace event: {error}"))?; if let Some(event) = accept_live(cursor, event) { emit_workspace(app, generation, expected_generation, event)?; } } } pub async fn subscribe_workspace( socket_path: String, mut cursor: SubscriptionCursor, app: AppHandle, generation: Arc, expected_generation: u64, ) -> Result<(), String> { let mut attempt = 0usize; loop { if !generation.is_current(expected_generation)? { return Ok(()); } match subscribe_workspace_once( &socket_path, &mut cursor, &app, &generation, expected_generation, ) .await { Ok(()) => attempt = 0, Err(error) => { let delay = retry_delay(attempt); if error == "Workspace subscription was replaced" { return Ok(()); } emit_workspace( &app, &generation, expected_generation, WorkspaceHostEvent::Disconnected { message: error, retry_in_ms: delay, }, )?; sleep(Duration::from_millis(delay)).await; attempt = (attempt + 1).min(3); } } } } pub async fn subscribe( socket_path: String, agent_id: String, cursor: u64, app: AppHandle, ) -> Result<(), String> { let mut reader = connect_and_send( &socket_path, "subscribe", Some(&agent_id), Some(json!({ "cursor": cursor })), ) .await?; loop { let mut line = String::new(); let bytes = reader .read_line(&mut line) .await .map_err(|error| format!("Bridge subscription failed: {error}"))?; if bytes == 0 { return Err("Bridge subscription closed".to_owned()); } let value: Value = serde_json::from_str(&line) .map_err(|error| format!("Bridge emitted invalid JSON: {error}"))?; if value.get("id") == Some(&Value::String("tauri-ui".to_owned())) { let result = result_from_response(value)?; for event in result .get("events") .and_then(Value::as_array) .into_iter() .flatten() { app.emit("bridge-event", event) .map_err(|error| format!("Could not publish bridge event: {error}"))?; } } else if value.get("type") == Some(&Value::String("event".to_owned())) { if let Some(event) = value.get("event") { app.emit("bridge-event", event) .map_err(|error| format!("Could not publish bridge event: {error}"))?; } } } } #[cfg(test)] mod tests { use super::*; use tokio::net::UnixListener; #[tokio::test] async fn sends_one_jsonl_request_and_unwraps_its_result() { let path = format!( "{}/pi-status-ui-{}.sock", env::temp_dir().display(), std::process::id() ); let _ = std::fs::remove_file(&path); let listener = UnixListener::bind(&path).expect("listener"); let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.expect("connection"); let mut line = String::new(); let mut reader = BufReader::new(stream); reader.read_line(&mut line).await.expect("request"); let request: Value = serde_json::from_str(&line).expect("JSON request"); assert_eq!(request["op"], "submit_prompt"); assert_eq!(request["agentId"], "agent-1"); assert_eq!(request["payload"]["message"], "Hello Pi"); let mut response = serde_json::to_vec(&json!({ "id": "tauri-ui", "ok": true, "result": { "accepted": true } })) .expect("response JSON"); response.push(b'\n'); reader .get_mut() .write_all(&response) .await .expect("response"); }); let result = request( &path, "submit_prompt", Some("agent-1"), Some(json!({ "message": "Hello Pi" })), ) .await .expect("request succeeds"); assert_eq!(result["accepted"], true); server.await.expect("server succeeds"); std::fs::remove_file(path).expect("socket cleanup"); } fn runtime(id: &str) -> Value { json!({ "runtimeId": id, "worktreePath": "/tmp/project", "state": "idle", "label": "Session", "attention": false, "queueCount": 0, "lastActivity": "2026-01-01T00:00:00Z", "openedAt": "2026-01-01T00:00:00Z" }) } async fn assert_typed_frame(operation: &'static str, payload: Value, call: F) where F: FnOnce(String) -> Fut, Fut: std::future::Future>, { let path = format!( "{}/pi-status-ui-frame-{operation}-{}-{}.sock", env::temp_dir().display(), std::process::id(), operation.len() ); let _ = std::fs::remove_file(&path); let listener = UnixListener::bind(&path).expect("listener"); let expected_payload = payload.clone(); let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.expect("connection"); let mut reader = BufReader::new(stream); let mut line = String::new(); reader.read_line(&mut line).await.unwrap(); let frame: Value = serde_json::from_str(&line).unwrap(); assert_eq!(frame["op"], operation); assert_eq!( frame.get("payload").cloned().unwrap_or(Value::Null), expected_payload ); let result = match operation { "close_session_runtime" => json!({"runtimeId":"r"}), "list_directory_sessions" => json!({"sessions":[]}), "get_session_runtime_snapshot" => { json!({"bridgeInstanceId":"b","latestSeq":0,"runtime":runtime("r"),"extensions":[]}) } _ => json!({"runtime":runtime("r")}), }; let response = format!("{}\n", json!({"id":"tauri-ui","ok":true,"result":result})); reader .get_mut() .write_all(response.as_bytes()) .await .unwrap(); }); call(path.clone()).await.unwrap(); server.await.unwrap(); std::fs::remove_file(path).unwrap(); } async fn assert_no_payload(operation: &'static str, result: Value, call: F) -> T where T: Send + 'static, F: FnOnce(String) -> Fut, Fut: std::future::Future>, { let path = format!( "{}/pi-status-ui-empty-{operation}-{}.sock", env::temp_dir().display(), std::process::id() ); let _ = std::fs::remove_file(&path); let listener = UnixListener::bind(&path).unwrap(); let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.unwrap(); let mut reader = BufReader::new(stream); let mut line = String::new(); reader.read_line(&mut line).await.unwrap(); let frame: Value = serde_json::from_str(&line).unwrap(); assert_eq!(frame["op"], operation); assert!(frame.get("payload").is_none()); let response = format!("{}\n", json!({"id":"tauri-ui","ok":true,"result":result})); reader .get_mut() .write_all(response.as_bytes()) .await .unwrap(); }); let value = call(path.clone()).await.unwrap(); server.await.unwrap(); std::fs::remove_file(path).unwrap(); value } #[tokio::test] async fn workspace_wrappers_emit_exact_operations_without_payloads() { let workspace: Workspace = assert_no_payload( "get_workspace", json!({"bridgeInstanceId":"b","latestSeq":0,"directories":[]}), |path| async move { get_workspace(&path).await }, ) .await; assert_eq!(workspace.bridge_instance_id, "b"); let summary: WorkspaceSummary = assert_no_payload( "get_workspace_summary", json!({"bridgeInstanceId":"b","latestSeq":0,"openCount":0,"workingCount":0,"attentionCount":0,"recoveringCount":0,"errorCount":0,"directoryCount":0,"resourceWarning":false}), |path| async move { get_workspace_summary(&path).await }, ) .await; assert_eq!(summary.open_count, 0); } #[tokio::test] async fn runtime_wrappers_emit_exact_operations_and_payloads() { assert_typed_frame( "create_session_runtime", json!({"worktreePath":"/tmp/project"}), |path| async move { create_session_runtime(&path, "/tmp/project") .await .map(|_| ()) }, ) .await; assert_typed_frame( "open_session_runtime", json!({"worktreePath":"/tmp/project","sessionPath":"/tmp/project/s.jsonl"}), |path| async move { open_session_runtime(&path, "/tmp/project", "/tmp/project/s.jsonl") .await .map(|_| ()) }, ) .await; assert_typed_frame( "close_session_runtime", json!({"runtimeId":"r"}), |path| async move { close_session_runtime(&path, "r").await.map(|_| ()) }, ) .await; assert_typed_frame( "list_directory_sessions", json!({"worktreePath":"/tmp/project"}), |path| async move { list_directory_sessions(&path, "/tmp/project") .await .map(|_| ()) }, ) .await; assert_typed_frame( "get_session_runtime_snapshot", json!({"runtimeId":"r"}), |path| async move { get_session_runtime_snapshot(&path, "r").await.map(|_| ()) }, ) .await; } #[test] fn decodes_multi_runtime_workspace_dormant_snapshot_and_directory_sessions() { let workspace: Workspace = serde_json::from_value(json!({ "bridgeInstanceId":"bridge", "latestSeq":4, "directories":[{"worktreePath":"/tmp/project","isHome":true,"openCount":2,"workingCount":1,"attentionCount":0,"recoveringCount":0,"errorCount":0,"runtimes":[runtime("one"),runtime("two")]}] })).unwrap(); assert_eq!(workspace.directories[0].runtimes.len(), 2); let snapshot: RuntimeSnapshot = serde_json::from_value(json!({"bridgeInstanceId":"bridge","latestSeq":4,"runtime":runtime("failed"),"extensions":[]})).unwrap(); assert!(snapshot.state.is_none()); let sessions: DirectorySessionsResult = serde_json::from_value(json!({ "sessions": [ {"path":"/tmp/project/one.jsonl","id":"one","cwd":"/tmp/project","modified":"2026-01-01T00:00:00Z","messageCount":2,"isCurrent":false,"runtimeId":"runtime-one"}, {"path":"/tmp/project/two.jsonl","id":"two","cwd":"/tmp/project","modified":"2026-01-02T00:00:00Z","messageCount":0,"isCurrent":false} ] })).unwrap(); assert_eq!( sessions.sessions[0].runtime_id.as_deref(), Some("runtime-one") ); assert!(sessions.sessions[1].runtime_id.is_none()); } #[tokio::test] async fn workspace_subscription_frame_uses_global_cursor_without_agent() { let path = format!( "{}/pi-status-ui-workspace-subscribe-{}.sock", env::temp_dir().display(), std::process::id() ); let _ = std::fs::remove_file(&path); let listener = UnixListener::bind(&path).unwrap(); let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.unwrap(); let mut reader = BufReader::new(stream); let mut line = String::new(); reader.read_line(&mut line).await.unwrap(); let frame: Value = serde_json::from_str(&line).unwrap(); assert_eq!(frame["op"], "subscribe_workspace"); assert_eq!(frame["payload"], json!({"cursor":17})); assert!(frame.get("agentId").is_none()); }); let _reader = connect_and_send( &path, "subscribe_workspace", None, Some(json!({"cursor":17})), ) .await .unwrap(); server.await.unwrap(); std::fs::remove_file(path).unwrap(); } fn event(epoch: &str, seq: u64) -> WorkspaceEvent { serde_json::from_value( json!({"bridgeInstanceId":epoch,"seq":seq,"type":"runtime_event","data":{}}), ) .unwrap() } #[test] fn workspace_host_events_use_exact_camel_case_json() { let event = event("bridge-1", 7); let cases = [ ( WorkspaceHostEvent::Connected { bridge_instance_id: "bridge-1".into(), first_available_seq: 3, latest_seq: 7, }, json!({ "kind": "connected", "bridgeInstanceId": "bridge-1", "firstAvailableSeq": 3, "latestSeq": 7 }), ), ( WorkspaceHostEvent::Event { event: event.clone(), }, json!({ "kind": "event", "event": event }), ), ( WorkspaceHostEvent::ResetRequired { bridge_instance_id: "bridge-1".into(), latest_seq: 7, reason: "sequenceGap".into(), }, json!({ "kind": "resetRequired", "bridgeInstanceId": "bridge-1", "latestSeq": 7, "reason": "sequenceGap" }), ), ( WorkspaceHostEvent::Disconnected { message: "offline".into(), retry_in_ms: 500, }, json!({ "kind": "disconnected", "message": "offline", "retryInMs": 500 }), ), ]; for (host_event, expected) in cases { assert_eq!(serde_json::to_value(host_event).unwrap(), expected); } } #[test] fn replacement_waits_for_an_in_flight_emission_and_rejects_stale_work() { use std::sync::mpsc; use std::thread; use std::time::Duration as StdDuration; let generation = Arc::new(SubscriptionGeneration::default()); let first = generation.replace().unwrap(); let (entered_tx, entered_rx) = mpsc::channel(); let (release_tx, release_rx) = mpsc::channel(); let emitter_generation = Arc::clone(&generation); let emitter = thread::spawn(move || { emitter_generation.run_if_current(first, || { entered_tx.send(()).unwrap(); release_rx.recv().unwrap(); Ok(()) }) }); entered_rx.recv().unwrap(); let replacement_generation = Arc::clone(&generation); let (replaced_tx, replaced_rx) = mpsc::channel(); let replacement = thread::spawn(move || { let next = replacement_generation.replace().unwrap(); replaced_tx.send(next).unwrap(); }); assert!(replaced_rx .recv_timeout(StdDuration::from_millis(50)) .is_err()); release_tx.send(()).unwrap(); assert!(emitter.join().unwrap().is_ok()); assert_eq!(replaced_rx.recv().unwrap(), first + 1); replacement.join().unwrap(); assert_eq!( generation.run_if_current(first, || Ok(())), Err("Workspace subscription was replaced".to_owned()) ); assert_eq!( (0..6).map(retry_delay).collect::>(), vec![250, 500, 1000, 2000, 2000, 2000] ); } #[test] fn replay_and_live_decisions_handle_epochs_truncation_and_gaps() { let mut cursor = SubscriptionCursor { bridge_instance_id: Some("a".into()), seq: 1, }; let replay = WorkspaceReplay { bridge_instance_id: "a".into(), first_available_seq: 1, latest_seq: 3, truncated: false, events: vec![event("a", 2), event("a", 3)], }; assert_eq!(accept_replay(&mut cursor, &replay).len(), 3); assert_eq!(cursor.seq, 3); assert!(matches!( accept_live(&mut cursor, event("a", 4)), Some(WorkspaceHostEvent::Event { .. }) )); assert!(accept_live(&mut cursor, event("a", 4)).is_none()); assert!( matches!(accept_live(&mut cursor, event("a",6)), Some(WorkspaceHostEvent::ResetRequired { reason, .. }) if reason == "sequenceGap") ); let changed = WorkspaceReplay { bridge_instance_id: "b".into(), first_available_seq: 1, latest_seq: 9, truncated: false, events: vec![], }; assert!( matches!(&accept_replay(&mut cursor, &changed)[0], WorkspaceHostEvent::ResetRequired { reason, .. } if reason == "epochChanged") ); let truncated = WorkspaceReplay { bridge_instance_id: "b".into(), first_available_seq: 10, latest_seq: 12, truncated: true, events: vec![], }; assert!( matches!(&accept_replay(&mut cursor, &truncated)[0], WorkspaceHostEvent::ResetRequired { reason, .. } if reason == "replayTruncated") ); } }