diff --git a/Cargo.lock b/Cargo.lock index 41e8bef..d49ee7e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -883,6 +883,7 @@ dependencies = [ "tauri-plugin-dialog", "tauri-plugin-opener", "tauri-plugin-store", + "tokio", ] [[package]] @@ -1937,6 +1938,7 @@ dependencies = [ "schemars 1.2.1", "serde", "serde_json", + "tokio", "ts-rs", ] @@ -1945,12 +1947,14 @@ name = "inkfinite-core" version = "0.0.0" dependencies = [ "automerge", + "getrandom 0.4.3", "perfect_freehand", "proptest", "schemars 1.2.1", "serde", "serde_json", "thiserror 2.0.18", + "tokio", "ts-rs", ] diff --git a/Cargo.toml b/Cargo.toml index 66b5ca9..17e51a9 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -21,6 +21,8 @@ perfect_freehand = "=0.1.1" schemars = "1.2.1" serde = { version = "1.0", features = ["derive"] } serde_json = "1.0" +getrandom = "0.4.3" +tokio = { version = "1.48", features = ["io-util", "macros", "net", "rt-multi-thread", "sync"] } ts-rs = { version = "12.0.1", features = ["no-serde-warnings", "serde-json-impl"] } thiserror = "2.0" diff --git a/README.md b/README.md index 1fb3cc5..54e9f23 100644 --- a/README.md +++ b/README.md @@ -42,6 +42,7 @@ pnpm tauri dev ## CLI The `inkfinite` CLI works on `.inkfinite` files while the desktop app is closed. +It can also inspect a running desktop app through authenticated local IPC. Run it through Cargo during development: @@ -112,6 +113,22 @@ cargo run -p inkfinite-cli --bin inkfinite -- capabilities --json Run `inkfinite --help` for examples and the complete command reference. +### Live desktop control + +With the desktop app running, use `app status`, `app inspect`, and `app query` to +read its open sessions. `app focus` asks the frontend to bring its window forward: + +```sh +inkfinite app status --json +inkfinite app inspect --json +inkfinite app query --role architecture.service --json +inkfinite app focus +``` + +The desktop publishes a per-user Unix-domain socket on Unix-like systems or a +per-user named pipe on Windows. A protected discovery file carries a random +process token, and requests use versioned length-prefixed frames. + ## Inkfinite files Canonical `.inkfinite` files contain the document and its change history in a diff --git a/ROADMAP.md b/ROADMAP.md index 5908e3a..95d84ef 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -1,6 +1,6 @@ # Inkfinite vNext / Version 2 -Status: implementation in progress; V2-01 through V2-15 are complete +Status: implementation in progress; V2-01 through V2-17 are complete This is the product and architecture contract for vNext. [TODO.md](TODO.md) is the implementation queue. @@ -60,8 +60,11 @@ the CLI and a bundled `SKILL.md`; MCP and UI automation are not part of vNext. the locked atomic file boundary. The CLI also exposes generated schemas, global `--json` and `--non-interactive` options, task-oriented help, and a machine-readable capability contract. -- [TODO.md](TODO.md) starts the remaining work at V2-17: authenticated local - control, followed by sync, agent packaging, release verification, and v1 +- The desktop owns an authenticated, versioned local IPC server. The CLI can + inspect open sessions, query the same materialized records as file mode, and + request frontend focus without a TCP listener or background daemon. +- [TODO.md](TODO.md) starts the remaining work at V2-18: reviewable live + proposals, followed by sync, agent packaging, release verification, and v1 compatibility removal. ## Architecture @@ -284,6 +287,7 @@ inkfinite query architecture.inkfinite --role architecture.service --json inkfinite validate architecture.inkfinite inkfinite schema document inkfinite capabilities --json +inkfinite app status --json ``` V2-16 added mutating and rendering commands: @@ -324,8 +328,8 @@ never edit document bytes manually. ## Live control and collaboration -The CLI may connect to the running desktop app for `app status`, `inspect`, -`query`, `propose`, `apply`, and `focus`. +The CLI connects to the running desktop app for `app status`, `app inspect`, +`app query`, and `app focus`. V2-18 will add proposals and explicit apply. `app propose` is the default agent path. Rust validates it and the UI shows a ghost preview plus created, changed, and deleted IDs. Rejection changes nothing. @@ -333,10 +337,11 @@ Partial acceptance creates a new transaction from the selected operations and revalidates it against current heads. `app apply` requires explicit user authorization. -Use a per-user Unix-domain socket or Windows named pipe, a per-install or +V2-17 uses a per-user Unix-domain socket or Windows named pipe, a protected per-session token, protocol versions, length-prefixed messages, and strict size -limits. Do not expose public TCP or local HTTP. The Tauri process hosts the -server; vNext has no background daemon. +limits. It does not expose public TCP or local HTTP. The Tauri process owns the +server and removes its discovery record when it exits; vNext has no background +daemon. Automerge sync between trusted peers is a vNext deliverable. The sync layer must be transport-independent and prove offline concurrent edits between two app @@ -410,8 +415,9 @@ deterministic SVG, and the mutating file-mode CLI all passed. - **Milestone 5, layers and styles:** complete. - **Milestone 6, CLI and SVG:** complete. The file-mode CLI exposes the SVG renderer and validated, atomic mutation commands. -- **Milestone 7, live control and sync:** add authenticated local IPC, reviewable - proposals, explicit apply, and two-replica offline convergence. +- **Milestone 7, live control and sync:** authenticated local IPC is complete; + reviewable proposals, explicit apply, and two-replica offline convergence + remain. - **Milestone 8, agent and release readiness:** bundle the agent skill, record release evidence, replace useful v1 coverage with v2-native fixtures, remove the unreleased v1 compatibility surface, and rerun the release matrix. diff --git a/TODO.md b/TODO.md index 9c19055..48d83de 100644 --- a/TODO.md +++ b/TODO.md @@ -90,54 +90,11 @@ opacity without regressing existing stencils. Shipped ordered layers with visibility, locking, active-layer state, opacity, and a complete Svelte panel. -Blocked by: V2-07, V2-09, V2-11 - -Acceptance criteria: - -- [x] New and imported pages always have a default layer; migration preserves - exact shape order and is idempotent. -- [x] Rendering follows page layer order and child order, skips hidden layers, - and composites layer opacity without leaking canvas state. -- [x] Hit testing, marquee, selection UI, editing, and agent transactions ignore - hidden shapes and reject locked-layer changes. -- [x] New shapes use the active layer. Moving and reordering layers or shapes is - one undoable transaction and converges under concurrent edits. -- [x] The panel lists, selects, creates, renames, reorders, hides, locks, deletes, - and changes opacity with accessible controls. -- [x] Deleting a non-empty layer requires an explicit move destination or an - explicit content deletion; the last layer cannot disappear. - -Verification: - -- Run model/engine, renderer, web component, migration, undo, and two-replica - layer tests. - ### V2-13: Add shape opacity and finish active-layer stencils Exposed fill and stroke opacity, completed the curated built-in stencil set, and made stencil insertion obey active-layer rules. -Blocked by: V2-12 - -Acceptance criteria: - -- [x] Applicable shapes have validated fill and stroke opacity in `0..=1`, with - accessible inspector controls and deterministic Canvas output. -- [x] Existing files default to their current opaque appearance; stroke opacity - already present on freehand shapes migrates without drift. -- [x] The built-in library covers the intended flowchart, UI, and developer - diagram set without adding sharing or community-library infrastructure. -- [x] Palette click and drag insertion place every stencil shape in the active - layer, preserve grouping, snap when enabled, select the result, and create - one undoable transaction. -- [x] Built-in stencil fixtures pass in visible, hidden, locked, and translucent - layers. - -Verification: - -- Run focused model, Canvas renderer, stencil, inspector, migration, and undo - tests. V2-14 adds the matching headless SVG coverage. - ## Milestone 6: CLI and headless rendering Exit when scripts and agents can inspect, safely change, validate, and render a @@ -145,26 +102,8 @@ closed document without desktop code. ### V2-14: Render deterministic SVG -Added headless SVG rendering for all built-in shapes, layers, -bindings, transforms, opacity, text, and Markdown. - -Blocked by: V2-04, V2-06, V2-13 - -Acceptance criteria: - -- [x] Output is deterministic for a snapshot and supports page, layer, selection, - and region filtering. -- [x] Hidden and locked semantics, ordering, opacity, arrow routing/labels, - Markdown, text wrapping, and freehand strokes match Canvas fixtures. -- [x] Missing fonts or assets produce explicit warnings and deterministic - fallbacks. -- [x] Snapshot tests cover every built-in shape and a representative full board. - -Verification: - -```sh -cargo test -p inkfinite-core render:: -``` +Added headless SVG rendering for all built-in shapes, layers, bindings, transforms, +opacity, text, and Markdown. ### V2-15: Ship read-only and schema CLI commands @@ -172,64 +111,11 @@ Implemented `new`, `inspect`, `query`, `validate`, `schema`, and `capabilities` in file mode with stable human and machine output, global output controls, and task-oriented help. -Blocked by: V2-06, V2-07 - -Acceptance criteria: - -- [x] Commands work while the desktop app is closed and contain no business - logic outside shared crates. -- [x] `--json` never prompts, writes only machine data to stdout, sends - diagnostics to stderr, and returns documented stable exit codes. It works - before or after a subcommand. -- [x] Inspect/query report document heads and support semantic, hierarchy, layer, - kind, and bounds filters. -- [x] Schema and capability output matches generated artifacts and is snapshot - tested on Unix and Windows path conventions. -- [x] Top-level help includes common examples, version discovery, documentation - and issue links, and typo suggestions. Each subcommand has realistic - examples and clear value names. Running without a subcommand prints - concise help and returns the usage exit code. -- [x] `capabilities --json` reports the global `--json` and - `--non-interactive` options alongside commands, schemas, filters, format, - protocol, and exit codes. - -Verification: - -```sh -cargo test -p inkfinite-cli -``` - ### V2-16: Ship mutating CLI commands and SVG output -What to build: Add generic apply plus structured shape, connection, layout, and +Added generic apply plus structured shape, connection, layout, and render commands, all through the transaction engine. -Blocked by: V2-05, V2-14, V2-15 - -Acceptance criteria: - -- [x] `apply` accepts a transaction from a file or stdin and supports dry-run, - inspected heads, record preconditions, and deterministic JSON results. -- [x] `shape create/patch/delete`, `connect`, and `layout` build ordinary - transactions and honor layers, locks, permissions, and semantic selectors. -- [x] Failed validation, stale preconditions, file locks, or write errors leave - the original byte-for-byte unchanged. -- [x] Results report previous/current heads, transaction ID, created/updated/ - deleted IDs, repairs, and warnings. -- [x] `render` writes deterministic SVG without opening the desktop app. -- [x] New subcommands preserve the global `--json` and `--non-interactive` - options, stdout/stderr separation, stable exit codes, unambiguous names, - descriptive long flags, built-in examples, and project support links. -- [x] Help, `capabilities --json`, README.md, TODO.md, ROADMAP.md, and CLI - integration tests describe the same shipped interface. - -Verification: - -- Run CLI workflow tests for inspect → dry-run → apply → validate → reopen → - render, including every failure path. Cover top-level and subcommand help, - both placements of global options, version output, and the capability - contract. - ## Milestone 7: Live control and CRDT sync Exit when a running desktop app can be inspected and changed safely and two @@ -237,19 +123,19 @@ offline app replicas converge after reconnecting. ### V2-17: Add authenticated local IPC -What to build: Host a versioned local-socket server in Tauri and connect the CLI +Implemented: Host a versioned local-socket server in Tauri and connect the CLI for status, inspect, query, and focus. Blocked by: V2-08, V2-15 Acceptance criteria: -- [ ] Unix-domain sockets and Windows named pipes use per-user names, a protected +- [x] Unix-domain sockets and Windows named pipes use per-user names, a protected per-install/session token, length-prefix framing, versions, and size limits. -- [ ] The server exposes no TCP/HTTP listener and stops with the Tauri process. -- [ ] Read-only app commands return the same protocol records and query results as +- [x] The server exposes no TCP/HTTP listener and stops with the Tauri process. +- [x] Read-only app commands return the same protocol records and query results as file mode; focus emits a small frontend notification. -- [ ] Tests reject wrong tokens, oversized/truncated frames, unsupported versions, +- [x] Tests reject wrong tokens, oversized/truncated frames, unsupported versions, malformed JSON, replayed request IDs, and unavailable app sessions. Verification: @@ -431,5 +317,4 @@ Verification: ## Frontier -V2-17 is the current frontier. Add authenticated local IPC for read-only live -status, inspection, queries, and focus. +V2-18 is the current frontier. Add reviewable live proposals and explicit apply. diff --git a/apps/desktop/src-tauri/Cargo.toml b/apps/desktop/src-tauri/Cargo.toml index 1b63a50..ab2a538 100644 --- a/apps/desktop/src-tauri/Cargo.toml +++ b/apps/desktop/src-tauri/Cargo.toml @@ -25,3 +25,4 @@ tauri-plugin-store = "2" inkfinite-core.workspace = true serde = { version = "1", features = ["derive"] } serde_json = "1" +tokio.workspace = true diff --git a/apps/desktop/src-tauri/src/ipc.rs b/apps/desktop/src-tauri/src/ipc.rs new file mode 100644 index 0000000..ceaec41 --- /dev/null +++ b/apps/desktop/src-tauri/src/ipc.rs @@ -0,0 +1,259 @@ +//! Local IPC server for read-only live CLI control. + +use std::sync::{Arc, Mutex}; + +use inkfinite_core::ipc::{ + dispatch, endpoint_name, ensure_ipc_directory, random_secret, read_frame, remove_discovery, write_discovery, + write_frame, AppRequest, AppResponse, DiscoveryRecord, IpcError, RequestEnvelope, RequestGuard, ResponseEnvelope, + FOCUS_EVENT, +}; +use inkfinite_core::proto::{ProtocolError, PROTOCOL_ID, PROTOCOL_VERSION}; +use inkfinite_core::session::SessionService; +use serde_json::json; +use tauri::{AppHandle, Emitter}; +use tokio::sync::oneshot; + +/// Handle used by the Tauri lifecycle to stop the local server and clean up discovery. +pub struct IpcServerHandle { + shutdown: Mutex>>, + discovery: DiscoveryRecord, + endpoint: String, +} + +impl IpcServerHandle { + /// Requests shutdown and removes only this server's published endpoint metadata. + pub fn stop(&self) { + if let Ok(mut shutdown) = self.shutdown.lock() { + if let Some(sender) = shutdown.take() { + let _ = sender.send(()); + } + } + cleanup_owned_endpoint(&self.discovery, &self.endpoint); + } +} + +impl Drop for IpcServerHandle { + fn drop(&mut self) { + self.stop(); + } +} + +/// Starts the per-user local server for one desktop process. +/// +/// # Errors +/// +/// Returns [`IpcError`] when the endpoint or protected discovery record cannot +/// be created. The server is not published until both are ready. +pub async fn start(app: AppHandle, service: Arc>) -> Result { + ensure_ipc_directory()?; + let endpoint = endpoint_name(); + #[cfg(unix)] + let listener = bind_unix_endpoint(&endpoint)?; + #[cfg(windows)] + let listener = create_named_pipe(&endpoint, true)?; + + let token = random_secret().map_err(|error| IpcError::Unavailable(error.to_string()))?; + let discovery = DiscoveryRecord { + protocol_id: PROTOCOL_ID.into(), + version: PROTOCOL_VERSION, + endpoint: endpoint.clone(), + token: token.clone(), + }; + if let Err(error) = write_discovery(&inkfinite_core::ipc::discovery_path(), &discovery) { + remove_endpoint(&endpoint); + return Err(error); + } + + let (shutdown_sender, shutdown_receiver) = oneshot::channel(); + let guard = Arc::new(Mutex::new(RequestGuard::new(token))); + let discovery_for_task = discovery.clone(); + #[cfg(unix)] + tauri::async_runtime::spawn(run_unix( + listener, + app, + service, + guard, + shutdown_receiver, + discovery_for_task, + )); + #[cfg(windows)] + tauri::async_runtime::spawn(run_windows( + listener, + endpoint.clone(), + app, + service, + guard, + shutdown_receiver, + discovery_for_task, + )); + + Ok(IpcServerHandle { shutdown: Mutex::new(Some(shutdown_sender)), discovery, endpoint }) +} + +#[cfg(unix)] +fn bind_unix_endpoint(endpoint: &str) -> Result { + use std::os::unix::fs::{FileTypeExt, PermissionsExt}; + + let path = std::path::Path::new(endpoint); + let listener = match std::os::unix::net::UnixListener::bind(path) { + Ok(listener) => listener, + Err(error) if error.kind() == std::io::ErrorKind::AddrInUse => { + let metadata = std::fs::symlink_metadata(path)?; + if !metadata.file_type().is_socket() { + return Err(IpcError::Unavailable(format!( + "IPC endpoint is not a Unix socket: {}", + path.display() + ))); + } + if std::os::unix::net::UnixStream::connect(path).is_ok() { + return Err(IpcError::Unavailable( + "another Inkfinite desktop process owns the IPC endpoint".into(), + )); + } + std::fs::remove_file(path)?; + std::os::unix::net::UnixListener::bind(path)? + } + Err(error) => return Err(IpcError::Io(error)), + }; + listener.set_nonblocking(true)?; + let mut permissions = std::fs::metadata(path)?.permissions(); + permissions.set_mode(0o600); + std::fs::set_permissions(path, permissions)?; + tokio::net::UnixListener::from_std(listener).map_err(IpcError::from) +} + +#[cfg(windows)] +fn create_named_pipe( + endpoint: &str, first: bool, +) -> Result { + use tokio::net::windows::named_pipe::ServerOptions; + + let mut options = ServerOptions::new(); + options.reject_remote_clients(true).first_pipe_instance(first); + options.create(endpoint).map_err(IpcError::from) +} + +#[cfg(unix)] +async fn run_unix( + listener: tokio::net::UnixListener, app: AppHandle, service: Arc>, + guard: Arc>, mut shutdown: oneshot::Receiver<()>, discovery: DiscoveryRecord, +) { + let mut connections = tokio::task::JoinSet::new(); + loop { + tokio::select! { + _ = &mut shutdown => break, + accepted = listener.accept() => match accepted { + Ok((stream, _)) => { + let app = app.clone(); + let service = Arc::clone(&service); + let guard = Arc::clone(&guard); + connections.spawn(handle_connection(stream, app, service, guard)); + } + Err(_) => break, + }, + } + } + connections.abort_all(); + while connections.join_next().await.is_some() {} + cleanup_owned_endpoint(&discovery, &discovery.endpoint); +} + +#[cfg(windows)] +async fn run_windows( + mut listener: tokio::net::windows::named_pipe::NamedPipeServer, endpoint: String, app: AppHandle, + service: Arc>, guard: Arc>, mut shutdown: oneshot::Receiver<()>, + discovery: DiscoveryRecord, +) { + loop { + let connected = tokio::select! { + _ = &mut shutdown => break, + result = listener.connect() => result, + }; + if connected.is_err() { + break; + } + let app = app.clone(); + let service = Arc::clone(&service); + let guard = Arc::clone(&guard); + tokio::select! { + _ = &mut shutdown => break, + _ = handle_connection(listener, app, service, guard) => {}, + } + listener = match create_named_pipe(&endpoint, false) { + Ok(listener) => listener, + Err(_) => break, + }; + } + cleanup_owned_endpoint(&discovery, &discovery.endpoint); +} + +async fn handle_connection( + mut stream: S, app: AppHandle, service: Arc>, guard: Arc>, +) where + S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin, +{ + let request = match read_frame::(&mut stream).await { + Ok(request) => request, + Err(_) => return, + }; + let request_id = request.request_id.clone(); + let result = match guard.lock() { + Ok(mut guard) => match guard.validate(&request) { + Ok(()) => dispatch_request(&app, &service, request.request), + Err(error) => Err(error), + }, + Err(_) => Err(protocol_error( + "ipc_guard_unavailable", + "the IPC authentication guard is unavailable", + )), + }; + let response = ResponseEnvelope { request_id, result }; + let _ = write_frame(&mut stream, &response).await; +} + +fn dispatch_request( + app: &AppHandle, service: &Arc>, request: AppRequest, +) -> Result { + let mut service = service.lock().map_err(|_| { + protocol_error( + "session_service_unavailable", + "the desktop session service lock is poisoned", + ) + })?; + let response = dispatch(&mut service, request)?; + if matches!(response, AppResponse::Focused) { + app.emit(FOCUS_EVENT, json!({ "source": "cli" })).map_err(|error| { + protocol_error( + "focus_notification_failed", + format!("could not notify the desktop frontend: {error}"), + ) + })?; + } + Ok(response) +} + +fn protocol_error(code: &str, message: impl Into) -> ProtocolError { + ProtocolError { code: code.into(), message: message.into(), details: None } +} + +fn cleanup_owned_endpoint(discovery: &DiscoveryRecord, endpoint: &str) { + let path = inkfinite_core::ipc::discovery_path(); + let owns_discovery = inkfinite_core::ipc::read_discovery(&path).is_ok_and(|current| current == *discovery); + if owns_discovery { + let _ = remove_discovery(&path, discovery); + remove_endpoint(endpoint); + } +} + +#[cfg(unix)] +fn remove_endpoint(endpoint: &str) { + use std::os::unix::fs::FileTypeExt; + + let path = std::path::Path::new(endpoint); + if std::fs::symlink_metadata(path).is_ok_and(|metadata| metadata.file_type().is_socket()) { + let _ = std::fs::remove_file(path); + } +} + +#[cfg(windows)] +fn remove_endpoint(_endpoint: &str) {} diff --git a/apps/desktop/src-tauri/src/lib.rs b/apps/desktop/src-tauri/src/lib.rs index 1d6b2b0..8464771 100644 --- a/apps/desktop/src-tauri/src/lib.rs +++ b/apps/desktop/src-tauri/src/lib.rs @@ -1,13 +1,22 @@ mod files; +mod ipc; mod session; +use tauri::Manager; + #[cfg_attr(mobile, tauri::mobile_entry_point)] pub fn run() { - tauri::Builder::default() + let builder = tauri::Builder::default() .manage(session::DesktopState::default()) .plugin(tauri_plugin_opener::init()) .plugin(tauri_plugin_dialog::init()) .plugin(tauri_plugin_store::Builder::default().build()) + .setup(|app| { + let service = app.state::().service_handle(); + let server = tauri::async_runtime::block_on(ipc::start(app.handle().clone(), service))?; + app.manage(server); + Ok(()) + }) .invoke_handler(tauri::generate_handler![ session::create_document, session::open_document, @@ -24,7 +33,16 @@ pub fn run() { files::rename_file, files::delete_file, files::pick_workspace_directory - ]) - .run(tauri::generate_context!()) - .expect("error while running tauri application"); + ]); + let app = builder + .build(tauri::generate_context!()) + .expect("error while building tauri application"); + + app.run(|app_handle, event| { + if matches!(event, tauri::RunEvent::Exit) { + if let Some(server) = app_handle.try_state::() { + server.stop(); + } + } + }); } diff --git a/apps/desktop/src-tauri/src/session.rs b/apps/desktop/src-tauri/src/session.rs index 7345aa2..38f2f29 100644 --- a/apps/desktop/src-tauri/src/session.rs +++ b/apps/desktop/src-tauri/src/session.rs @@ -1,6 +1,6 @@ -//! Typed Tauri commands for Rust-owned document sessions. +//! Tauri commands for document sessions. -use std::sync::{Mutex, MutexGuard}; +use std::sync::{Arc, Mutex, MutexGuard}; use inkfinite_core::proto::{DocumentPath, ProtocolError, Query, QueryResult, SessionId, TransactionDraft}; use inkfinite_core::session::{ @@ -9,18 +9,26 @@ use inkfinite_core::session::{ use inkfinite_core::{ActorId, ChangeHash, DocumentId}; use tauri::State; -/// Tauri-managed owner of every open document session in this app process. +type Result = std::result::Result; + +/// Owner of every open document session in this app process. pub struct DesktopState { - service: Mutex, + service: Arc>, } impl Default for DesktopState { fn default() -> Self { - Self { service: Mutex::new(SessionService::new()) } + Self { service: Arc::new(Mutex::new(SessionService::new())) } + } +} + +impl DesktopState { + pub fn service_handle(&self) -> Arc> { + Arc::clone(&self.service) } } -fn lock_service(state: &DesktopState) -> Result, ProtocolError> { +fn lock_service(state: &DesktopState) -> Result> { state.service.lock().map_err(|_| ProtocolError { code: "session_service_unavailable".into(), message: "the desktop session service lock is poisoned".into(), @@ -29,35 +37,14 @@ fn lock_service(state: &DesktopState) -> Result, } fn to_protocol_error(error: SessionError) -> ProtocolError { - let code = match &error { - SessionError::NotFound(_) => "session_not_found", - SessionError::ActorMismatch { .. } => "actor_mismatch", - SessionError::StaleHeads => "stale_heads", - SessionError::AlreadyOpen { .. } => "document_already_open", - SessionError::File(file_error) => match file_error { - inkfinite_core::file::FileError::Locked { .. } => "document_locked", - inkfinite_core::file::FileError::AlreadyExists { .. } => "document_already_exists", - inkfinite_core::file::FileError::InvalidV1(_) - | inkfinite_core::file::FileError::Json(_) - | inkfinite_core::file::FileError::UnsupportedFormat { .. } - | inkfinite_core::file::FileError::UnsupportedShapeKind { .. } - | inkfinite_core::file::FileError::SamePath { .. } - | inkfinite_core::file::FileError::RecoveryNotFound { .. } - | inkfinite_core::file::FileError::InvalidRecovery(_) - | inkfinite_core::file::FileError::RecoveryAhead { .. } - | inkfinite_core::file::FileError::Engine(_) - | inkfinite_core::file::FileError::Io { .. } => "document_file_error", - }, - SessionError::Engine(_) => "document_engine_error", - }; - ProtocolError { code: code.into(), message: error.to_string(), details: None } + inkfinite_core::ipc::session_protocol_error(&error) } /// Creates a new canonical `.inkfinite` file and opens its session. #[tauri::command] pub fn create_document( state: State<'_, DesktopState>, path: String, document_id: String, actor_id: String, page_name: Option, -) -> Result { +) -> Result { lock_service(&state)? .create( path, @@ -70,9 +57,7 @@ pub fn create_document( /// Opens a canonical document or imports a selected frozen v1 JSON file. #[tauri::command] -pub fn open_document( - state: State<'_, DesktopState>, path: String, actor_id: String, -) -> Result { +pub fn open_document(state: State<'_, DesktopState>, path: String, actor_id: String) -> Result { lock_service(&state)? .open(path, ActorId::new(actor_id)) .map_err(to_protocol_error) @@ -80,7 +65,7 @@ pub fn open_document( /// Returns the current snapshot and session state. #[tauri::command] -pub fn snapshot(state: State<'_, DesktopState>, session_id: String) -> Result { +pub fn snapshot(state: State<'_, DesktopState>, session_id: String) -> Result { lock_service(&state)? .status(&SessionId(session_id)) .map_err(to_protocol_error) @@ -90,7 +75,7 @@ pub fn snapshot(state: State<'_, DesktopState>, session_id: String) -> Result, session_id: String, transaction: TransactionDraft, -) -> Result { +) -> Result { lock_service(&state)? .commit(&SessionId(session_id), transaction) .map_err(to_protocol_error) @@ -98,9 +83,7 @@ pub fn commit( /// Compensates the latest transaction for the session actor. #[tauri::command] -pub fn undo( - state: State<'_, DesktopState>, session_id: String, actor_id: String, -) -> Result { +pub fn undo(state: State<'_, DesktopState>, session_id: String, actor_id: String) -> Result { let session_id = SessionId(session_id); let actor_id = ActorId::new(actor_id); lock_service(&state)? @@ -110,9 +93,7 @@ pub fn undo( /// Reapplies the latest compensated transaction for the session actor. #[tauri::command] -pub fn redo( - state: State<'_, DesktopState>, session_id: String, actor_id: String, -) -> Result { +pub fn redo(state: State<'_, DesktopState>, session_id: String, actor_id: String) -> Result { let session_id = SessionId(session_id); let actor_id = ActorId::new(actor_id); lock_service(&state)? @@ -124,7 +105,7 @@ pub fn redo( #[tauri::command] pub fn save( state: State<'_, DesktopState>, session_id: String, expected_heads: Vec, -) -> Result { +) -> Result { lock_service(&state)? .save(&SessionId(session_id), &expected_heads) .map_err(to_protocol_error) @@ -134,7 +115,7 @@ pub fn save( #[tauri::command] pub fn save_as( state: State<'_, DesktopState>, session_id: String, path: DocumentPath, expected_heads: Vec, -) -> Result { +) -> Result { lock_service(&state)? .save_as(&SessionId(session_id), path.0, &expected_heads) .map_err(to_protocol_error) @@ -142,7 +123,7 @@ pub fn save_as( /// Queries records through the shared deterministic query implementation. #[tauri::command] -pub fn query(state: State<'_, DesktopState>, session_id: String, query: Query) -> Result { +pub fn query(state: State<'_, DesktopState>, session_id: String, query: Query) -> Result { lock_service(&state)? .query(&SessionId(session_id), &query) .map_err(to_protocol_error) @@ -150,7 +131,7 @@ pub fn query(state: State<'_, DesktopState>, session_id: String, query: Query) - /// Validates the current session without changing it. #[tauri::command] -pub fn validate(state: State<'_, DesktopState>, session_id: String) -> Result { +pub fn validate(state: State<'_, DesktopState>, session_id: String) -> Result { lock_service(&state)? .validate(&SessionId(session_id)) .map_err(to_protocol_error) @@ -158,7 +139,7 @@ pub fn validate(state: State<'_, DesktopState>, session_id: String) -> Result, session_id: String) -> Result<(), ProtocolError> { +pub fn close(state: State<'_, DesktopState>, session_id: String) -> Result<()> { lock_service(&state)? .close(&SessionId(session_id)) .map_err(to_protocol_error) diff --git a/apps/desktop/src/routes/+layout.svelte b/apps/desktop/src/routes/+layout.svelte index 1d6fb31..73802b9 100644 --- a/apps/desktop/src/routes/+layout.svelte +++ b/apps/desktop/src/routes/+layout.svelte @@ -1,11 +1,31 @@ diff --git a/crates/inkfinite-cli/Cargo.toml b/crates/inkfinite-cli/Cargo.toml index 1423949..6cb2708 100644 --- a/crates/inkfinite-cli/Cargo.toml +++ b/crates/inkfinite-cli/Cargo.toml @@ -21,6 +21,7 @@ inkfinite-core.workspace = true schemars.workspace = true serde.workspace = true serde_json.workspace = true +tokio.workspace = true ts-rs.workspace = true [lints] diff --git a/crates/inkfinite-cli/src/cli/app.rs b/crates/inkfinite-cli/src/cli/app.rs new file mode 100644 index 0000000..bd65691 --- /dev/null +++ b/crates/inkfinite-cli/src/cli/app.rs @@ -0,0 +1,142 @@ +//! Commands for inspecting and focusing a running desktop app. + +use inkfinite_core::ipc::{self, AppRequest, AppResponse, IpcError}; +use inkfinite_core::proto::{Query, RecordId}; +use inkfinite_core::{LayerId, PageId}; + +use super::args::{AppCommand, AppInspectArgs, AppQueryArgs}; +use super::support::{map_output_error, write_heads, write_json}; +use super::{CliError, EXIT_INPUT, EXIT_INVALID, Result, Write, anyhow, json}; + +/// Runs one authenticated read-only desktop command. +pub fn run_app_command(command: AppCommand, json_output: bool, stdout: &mut dyn Write) -> Result<()> { + match command { + AppCommand::Status => status(json_output, stdout), + AppCommand::Inspect(args) => inspect(args, json_output, stdout), + AppCommand::Query(args) => query(args, json_output, stdout), + AppCommand::Focus => focus(json_output, stdout), + } +} + +fn status(json_output: bool, stdout: &mut dyn Write) -> Result<()> { + let response = send(AppRequest::Status)?; + let AppResponse::Status(statuses) = response else { + return unexpected_response("status"); + }; + if json_output { + return write_json(stdout, &statuses); + } + + writeln!(stdout, "Open sessions: {}", statuses.len()).map_err(map_output_error)?; + for status in statuses { + writeln!( + stdout, + "session\t{}\t{}\t{}", + status.session_id.0, + status.path.0, + if status.dirty { "dirty" } else { "saved" } + ) + .map_err(map_output_error)?; + write_heads(stdout, &status.snapshot.heads)?; + } + Ok(()) +} + +fn inspect(args: AppInspectArgs, json_output: bool, stdout: &mut dyn Write) -> Result<()> { + let response = send(AppRequest::Inspect { session_id: args.session_id.map(inkfinite_core::proto::SessionId) })?; + let AppResponse::Snapshot(snapshot) = response else { + return unexpected_response("inspect"); + }; + if json_output { + return write_json(stdout, &snapshot); + } + writeln!(stdout, "Document: {}", snapshot.document_id).map_err(map_output_error)?; + writeln!(stdout, "Format: {} {}", snapshot.format, snapshot.format_version).map_err(map_output_error)?; + write_heads(stdout, &snapshot.heads)?; + writeln!(stdout, "Pages: {}", snapshot.document.pages.len()).map_err(map_output_error)?; + writeln!(stdout, "Layers: {}", snapshot.document.layers.len()).map_err(map_output_error)?; + writeln!(stdout, "Shapes: {}", snapshot.document.shapes.len()).map_err(map_output_error)?; + writeln!(stdout, "Bindings: {}", snapshot.document.bindings.len()).map_err(map_output_error)?; + writeln!(stdout, "Assets: {}", snapshot.document.assets.len()).map_err(map_output_error) +} + +fn query(args: AppQueryArgs, json_output: bool, stdout: &mut dyn Write) -> Result<()> { + let query = Query { + id: args.id, + name: args.name, + role: args.role, + tag: args.tag, + shape_kind: args.shape_kind, + page_id: args.page.map(PageId::from), + layer_id: args.layer.map(LayerId::from), + parent_id: args.parent, + bounds: args.bounds, + }; + let response = + send(AppRequest::Query { session_id: args.session_id.map(inkfinite_core::proto::SessionId), query })?; + let AppResponse::QueryResult(result) = response else { + return unexpected_response("query"); + }; + if json_output { + return write_json(stdout, &result); + } + + write_heads(stdout, &result.heads)?; + writeln!(stdout, "Matches: {}", result.records.len()).map_err(map_output_error)?; + for record in result.records { + match record { + RecordId::Page(id) => writeln!(stdout, "page\t{id}"), + RecordId::Layer(id) => writeln!(stdout, "layer\t{id}"), + RecordId::Shape(id) => match result.bounds.get(&id) { + Some(bounds) => writeln!( + stdout, + "shape\t{id}\t{},{},{},{}", + bounds.x, bounds.y, bounds.width, bounds.height + ), + None => writeln!(stdout, "shape\t{id}"), + }, + RecordId::Binding(id) => writeln!(stdout, "binding\t{id}"), + RecordId::Asset(id) => writeln!(stdout, "asset\t{id}"), + } + .map_err(map_output_error)?; + } + Ok(()) +} + +fn focus(json_output: bool, stdout: &mut dyn Write) -> Result<()> { + let response = send(AppRequest::Focus)?; + if !matches!(response, AppResponse::Focused) { + return unexpected_response("focus"); + } + if json_output { + write_json(stdout, &json!({ "focused": true })) + } else { + writeln!(stdout, "Desktop focus requested").map_err(map_output_error) + } +} + +fn send(request: AppRequest) -> Result { + let runtime = tokio::runtime::Runtime::new() + .map_err(|error| CliError::new(EXIT_INPUT, anyhow!(error)).context("could not start the IPC runtime"))?; + let response = runtime.block_on(ipc::send(request)).map_err(map_ipc_error)?; + response + .result + .map_err(|error| CliError::new(EXIT_INVALID, anyhow!("[{}] {}", error.code, error.message))) +} + +fn map_ipc_error(error: IpcError) -> CliError { + let exit_code = match &error { + IpcError::FrameTooLarge { .. } | IpcError::DiscoveryTooLarge { .. } | IpcError::MalformedJson(_) => { + EXIT_INVALID + } + IpcError::TruncatedFrame | IpcError::Unavailable(_) | IpcError::Io(_) => EXIT_INPUT, + }; + CliError::new(exit_code, anyhow!(error)).context("could not contact the running desktop app") +} + +fn unexpected_response(command: &str) -> Result<()> { + Err(CliError::new( + EXIT_INVALID, + anyhow!("desktop app returned an unexpected response for {command}"), + )) +} diff --git a/crates/inkfinite-cli/src/cli/apply.rs b/crates/inkfinite-cli/src/cli/apply.rs index 9f80eb8..9f9c0e7 100644 --- a/crates/inkfinite-cli/src/cli/apply.rs +++ b/crates/inkfinite-cli/src/cli/apply.rs @@ -1,8 +1,8 @@ use super::mutation::commit_mutation; use super::support::{open_document, portable_path}; -use super::{ApplyArgs, CliError, EXIT_INPUT, EXIT_INVALID, Path, Read, TransactionDraft, Write, fs, io}; +use super::{ApplyArgs, CliError, EXIT_INPUT, EXIT_INVALID, Path, Read, Result, TransactionDraft, Write, fs, io}; -pub fn apply_transaction(args: &ApplyArgs, json_output: bool, stdout: &mut dyn Write) -> Result<(), CliError> { +pub fn apply_transaction(args: &ApplyArgs, json_output: bool, stdout: &mut dyn Write) -> Result<()> { let transaction_json = if args.transaction == Path::new("-") { let mut input = String::new(); io::stdin() diff --git a/crates/inkfinite-cli/src/cli/args.rs b/crates/inkfinite-cli/src/cli/args.rs index 51ae628..10f5794 100644 --- a/crates/inkfinite-cli/src/cli/args.rs +++ b/crates/inkfinite-cli/src/cli/args.rs @@ -5,13 +5,14 @@ use super::{ArgGroup, Args, Bounds, Parser, PathBuf, Subcommand, ValueEnum, pars name = "inkfinite", version, about = "Work with Inkfinite documents from the command line", - long_about = "Create, inspect, query, edit, validate, and render canonical Inkfinite documents while the desktop app is closed." + long_about = "Create, inspect, query, edit, validate, and render canonical Inkfinite documents, or inspect a running desktop app." )] #[command(after_help = "Examples: inkfinite new architecture.inkfinite inkfinite inspect architecture.inkfinite --json inkfinite apply architecture.inkfinite --transaction transaction.json --dry-run inkfinite render architecture.inkfinite --output architecture.svg + inkfinite app status --json Documentation: https://github.com/stormlightlabs/inkfinite#file-mode-cli Report issues: https://github.com/stormlightlabs/inkfinite/issues @@ -60,6 +61,9 @@ pub enum Command { inkfinite query architecture.inkfinite --kind rect --bounds 0,0,1920,1080 ")] Query(QueryArgs), + /// Inspect or focus a running desktop app over authenticated local IPC. + #[command(subcommand)] + App(AppCommand), /// Load and validate a canonical document. #[command(after_help = "Examples: @@ -433,6 +437,78 @@ pub struct QueryArgs { pub bounds: Option, } +#[derive(Debug, Subcommand)] +#[allow(clippy::large_enum_variant)] +pub enum AppCommand { + /// List the sessions currently open in the desktop app. + #[command(after_help = "Example: + + inkfinite app status --json +")] + Status, + /// Print the current snapshot from the desktop app. + #[command(after_help = "Examples: + + inkfinite app inspect --json + inkfinite app inspect --session-id session:1 --json +")] + Inspect(AppInspectArgs), + /// Query the current desktop session using shared semantic filters. + #[command(after_help = "Examples: + + inkfinite app query --role architecture.service --json + inkfinite app query --session-id session:1 --kind rect +")] + Query(AppQueryArgs), + /// Ask the desktop frontend to focus its main window. + #[command(after_help = "Example: + + inkfinite app focus +")] + Focus, +} + +#[derive(Debug, Args)] +pub struct AppInspectArgs { + /// Inspect this session, or the only open session when omitted. + #[arg(long, value_name = "SESSION_ID")] + pub session_id: Option, +} + +#[derive(Debug, Args)] +pub struct AppQueryArgs { + /// Query this session, or the only open session when omitted. + #[arg(long, value_name = "SESSION_ID")] + pub session_id: Option, + /// Match an exact record ID. + #[arg(long)] + pub id: Option, + /// Match an exact display name. + #[arg(long)] + pub name: Option, + /// Match an exact semantic role. + #[arg(long)] + pub role: Option, + /// Match one exact semantic tag. + #[arg(long)] + pub tag: Option, + /// Match an exact shape registry key. + #[arg(long = "kind")] + pub shape_kind: Option, + /// Restrict results to a page. + #[arg(long)] + pub page: Option, + /// Restrict results to a layer. + #[arg(long)] + pub layer: Option, + /// Restrict shapes to one direct parent. + #[arg(long)] + pub parent: Option, + /// Restrict shapes to bounds formatted as x,y,width,height. + #[arg(long, value_name = "X,Y,WIDTH,HEIGHT", value_parser = parse_bounds)] + pub bounds: Option, +} + #[derive(Clone, Copy, Debug, ValueEnum)] pub enum SchemaKind { Document, diff --git a/crates/inkfinite-cli/src/cli/contract.rs b/crates/inkfinite-cli/src/cli/contract.rs index a5af48f..039aa48 100644 --- a/crates/inkfinite-cli/src/cli/contract.rs +++ b/crates/inkfinite-cli/src/cli/contract.rs @@ -1,9 +1,11 @@ use super::support::{map_json_error, map_output_error, write_json}; +use super::{CliError, SchemaKind, Value, Write}; use super::{ - CliError, DOCUMENT_SCHEMA, EXIT_CONFLICT, EXIT_INPUT, EXIT_INVALID, INKFINITE_FORMAT_ID, INKFINITE_FORMAT_VERSION, + DOCUMENT_SCHEMA, EXIT_CONFLICT, EXIT_INPUT, EXIT_INVALID, INKFINITE_FORMAT_ID, INKFINITE_FORMAT_VERSION, PROTOCOL_ERROR_SCHEMA, PROTOCOL_ID, PROTOCOL_REQUEST_SCHEMA, PROTOCOL_RESPONSE_SCHEMA, PROTOCOL_VERSION, - SchemaKind, TRANSACTION_SCHEMA, Value, Write, builtin_shape_kinds, json, + TRANSACTION_SCHEMA, }; +use super::{builtin_shape_kinds, json}; pub fn print_schema(kind: SchemaKind, stdout: &mut dyn Write) -> Result<(), CliError> { match kind { @@ -34,7 +36,7 @@ pub fn print_schema(kind: SchemaKind, stdout: &mut dyn Write) -> Result<(), CliE pub fn print_capabilities(json_output: bool, stdout: &mut dyn Write) -> Result<(), CliError> { let capabilities = json!({ - "commands": ["new", "inspect", "query", "validate", "apply", "shape", "connect", "layout", "render", "schema", "capabilities"], + "commands": ["new", "inspect", "query", "app", "validate", "apply", "shape", "connect", "layout", "render", "schema", "capabilities"], "exit_codes": { "conflict": EXIT_CONFLICT, "input": EXIT_INPUT, @@ -46,6 +48,11 @@ pub fn print_capabilities(json_output: bool, stdout: &mut dyn Write) -> Result<( "format": { "id": INKFINITE_FORMAT_ID, "version": INKFINITE_FORMAT_VERSION }, "global_options": ["--json", "--non-interactive"], "json_stdout_is_machine_only": true, + "live_mode": { + "commands": ["status", "inspect", "query", "focus"], + "transport": "authenticated_local_socket", + "tcp_or_http": false + }, "mutation_commands": { "apply": ["--transaction", "--dry-run"], "connect": ["--binding-id", "--source", "--source-role", "--target", "--target-role", "--dry-run"], @@ -71,9 +78,10 @@ pub fn print_capabilities(json_output: bool, stdout: &mut dyn Write) -> Result<( .map_err(map_output_error)?; writeln!( stdout, - "Commands: new, inspect, query, validate, apply, shape, connect, layout, render, schema, capabilities" + "Commands: new, inspect, query, app, validate, apply, shape, connect, layout, render, schema, capabilities" ) .map_err(map_output_error)?; + writeln!(stdout, "Live mode: app status, app inspect, app query, app focus").map_err(map_output_error)?; writeln!(stdout, "Global options: --json, --non-interactive").map_err(map_output_error)?; writeln!(stdout, "Schemas: document, transaction, protocol").map_err(map_output_error)?; writeln!(stdout, "JSON: machine data on stdout; diagnostics on stderr").map_err(map_output_error) diff --git a/crates/inkfinite-cli/src/cli/document.rs b/crates/inkfinite-cli/src/cli/document.rs index d59545a..e281a29 100644 --- a/crates/inkfinite-cli/src/cli/document.rs +++ b/crates/inkfinite-cli/src/cli/document.rs @@ -3,10 +3,10 @@ use super::support::{ }; use super::{ ACTOR_ID, ActorId, CliError, DocumentFile, DocumentId, EXIT_INVALID, FileOutputArgs, LayerId, NewArgs, PageId, - Query, QueryArgs, RecordId, Write, anyhow, blank_document, json, validate_document, + Query, QueryArgs, RecordId, Result, Write, anyhow, blank_document, json, validate_document, }; -pub fn create_document(args: NewArgs, json_output: bool, stdout: &mut dyn Write) -> Result<(), CliError> { +pub fn create_document(args: NewArgs, json_output: bool, stdout: &mut dyn Write) -> Result<()> { let document_id = DocumentId::from(args.document_id.unwrap_or_else(|| default_document_id(&args.path))); if document_id.as_str().trim().is_empty() { return Err(CliError::new(EXIT_INVALID, anyhow!("document ID must not be empty"))); @@ -33,7 +33,7 @@ pub fn create_document(args: NewArgs, json_output: bool, stdout: &mut dyn Write) } } -pub fn inspect_document(args: &FileOutputArgs, json_output: bool, stdout: &mut dyn Write) -> Result<(), CliError> { +pub fn inspect_document(args: &FileOutputArgs, json_output: bool, stdout: &mut dyn Write) -> Result<()> { let mut file = open_document(&args.path)?; let snapshot = file.snapshot().map_err(map_file_error)?; if json_output { @@ -50,7 +50,7 @@ pub fn inspect_document(args: &FileOutputArgs, json_output: bool, stdout: &mut d writeln!(stdout, "Assets: {}", snapshot.document.assets.len()).map_err(map_output_error) } -pub fn query_document(args: QueryArgs, json_output: bool, stdout: &mut dyn Write) -> Result<(), CliError> { +pub fn query_document(args: QueryArgs, json_output: bool, stdout: &mut dyn Write) -> Result<()> { let mut file = open_document(&args.path)?; let query = Query { id: args.id, @@ -97,7 +97,7 @@ pub fn query_document(args: QueryArgs, json_output: bool, stdout: &mut dyn Write Ok(()) } -pub fn validate_file(args: &FileOutputArgs, json_output: bool, stdout: &mut dyn Write) -> Result<(), CliError> { +pub fn validate_file(args: &FileOutputArgs, json_output: bool, stdout: &mut dyn Write) -> Result<()> { let mut file = open_document(&args.path)?; let snapshot = file.snapshot().map_err(map_file_error)?; validate_document(&snapshot.document) diff --git a/crates/inkfinite-cli/src/cli/layout.rs b/crates/inkfinite-cli/src/cli/layout.rs index 0ac3168..8907331 100644 --- a/crates/inkfinite-cli/src/cli/layout.rs +++ b/crates/inkfinite-cli/src/cli/layout.rs @@ -2,9 +2,7 @@ use super::mutation::{commit_mutation, select_layout_shapes, structured_transact use super::support::open_document; use super::{AlignmentArg, AxisArg, BTreeMap, CliError, LayoutAxis, LayoutCommand, Operation, ShapeAlignment, Write}; -pub fn run_layout_command( - command: LayoutCommand, json_output: bool, stdout: &mut dyn Write, -) -> Result<(), CliError> { +pub fn run_layout_command(command: LayoutCommand, json_output: bool, stdout: &mut dyn Write) -> Result<(), CliError> { match command { LayoutCommand::Align(args) => { let mut file = open_document(&args.path)?; diff --git a/crates/inkfinite-cli/src/cli/mod.rs b/crates/inkfinite-cli/src/cli/mod.rs index 4f84c3e..0c1b691 100644 --- a/crates/inkfinite-cli/src/cli/mod.rs +++ b/crates/inkfinite-cli/src/cli/mod.rs @@ -34,6 +34,8 @@ const PROTOCOL_REQUEST_SCHEMA: &str = include_str!("../../../../schemas/protocol const PROTOCOL_RESPONSE_SCHEMA: &str = include_str!("../../../../schemas/protocol-response.schema.json"); const PROTOCOL_ERROR_SCHEMA: &str = include_str!("../../../../schemas/protocol-error.schema.json"); +pub type Result = std::result::Result; + #[derive(Debug)] pub struct CliError { exit_code: i32, @@ -54,6 +56,7 @@ impl CliError { } } +mod app; mod apply; mod args; mod connect; @@ -73,11 +76,12 @@ use support::parse_bounds; pub use args::{Cli, Command}; -pub fn run(command: Command, json_output: bool, stdout: &mut dyn Write) -> Result<(), CliError> { +pub fn run(command: Command, json_output: bool, stdout: &mut dyn Write) -> Result<()> { match command { Command::New(args) => document::create_document(args, json_output, stdout), Command::Inspect(args) => document::inspect_document(&args, json_output, stdout), Command::Query(args) => document::query_document(args, json_output, stdout), + Command::App(command) => app::run_app_command(command, json_output, stdout), Command::Validate(args) => document::validate_file(&args, json_output, stdout), Command::Apply(args) => apply::apply_transaction(&args, json_output, stdout), Command::Shape(command) => shape::run_shape_command(command, json_output, stdout), diff --git a/crates/inkfinite-cli/src/cli/mutation.rs b/crates/inkfinite-cli/src/cli/mutation.rs index 1f4a50b..18d686d 100644 --- a/crates/inkfinite-cli/src/cli/mutation.rs +++ b/crates/inkfinite-cli/src/cli/mutation.rs @@ -1,7 +1,7 @@ use super::support::{map_file_error, map_output_error, portable_path, write_heads, write_json}; use super::{ ACTOR_ID, ActorId, BTreeSet, CliError, DocumentFile, EXIT_INPUT, EXIT_INVALID, LayoutSelectionArgs, Operation, - Origin, Path, RecordId, Serialize, ShapeId, Timestamp, TransactionDraft, TransactionId, Write, anyhow, fs, + Origin, Path, RecordId, Result, Serialize, ShapeId, Timestamp, TransactionDraft, TransactionId, Write, anyhow, fs, }; #[derive(Serialize)] @@ -19,8 +19,8 @@ struct MutationResult { pub fn structured_transaction( file: &mut DocumentFile, transaction_id: Option, mut default_id: String, description: String, - operations: Vec, -) -> Result { + ops: Vec, +) -> Result { let base_heads = file.heads().map_err(map_file_error)?; let transaction_id = transaction_id.unwrap_or_else(|| { let head_suffix = base_heads @@ -38,14 +38,14 @@ pub fn structured_transaction( origin: Origin::Agent, base_heads, description, - operations, + operations: ops, timestamp: Timestamp(0), }) } pub fn commit_mutation( file: &mut DocumentFile, transaction: TransactionDraft, dry_run: bool, json_output: bool, stdout: &mut dyn Write, -) -> Result<(), CliError> { +) -> Result<()> { let previous_heads = file.heads().map_err(map_file_error)?; let commit = file.commit(transaction).map_err(map_file_error)?; if !dry_run { @@ -76,7 +76,7 @@ pub fn commit_mutation( pub fn select_unique_shape( file: &mut DocumentFile, shape_id: Option<&str>, name: Option<&str>, role: Option<&str>, -) -> Result { +) -> Result { let snapshot = file.snapshot().map_err(map_file_error)?; let matches: Vec = snapshot .document @@ -102,9 +102,7 @@ pub fn select_unique_shape( } } -pub fn select_layout_shapes( - file: &mut DocumentFile, selection: LayoutSelectionArgs, -) -> Result, CliError> { +pub fn select_layout_shapes(file: &mut DocumentFile, selection: LayoutSelectionArgs) -> Result> { let snapshot = file.snapshot().map_err(map_file_error)?; let mut selected: BTreeSet = selection.shape_ids.into_iter().map(ShapeId::from).collect(); if let Some(role) = selection.role { @@ -131,7 +129,7 @@ pub fn select_layout_shapes( Ok(selected.into_iter().collect()) } -pub fn read_json_argument(argument: &str, description: &str) -> Result { +pub fn read_json_argument(argument: &str, description: &str) -> Result { let Some(path) = argument.strip_prefix('@') else { return Ok(argument.to_owned()); }; diff --git a/crates/inkfinite-cli/src/cli/shape.rs b/crates/inkfinite-cli/src/cli/shape.rs index 8bb2399..197d13c 100644 --- a/crates/inkfinite-cli/src/cli/shape.rs +++ b/crates/inkfinite-cli/src/cli/shape.rs @@ -2,12 +2,12 @@ use super::mutation::{commit_mutation, read_json_argument, select_unique_shape, use super::support::open_document; use super::{ ACTOR_ID, ActorId, BTreeMap, CliError, EXIT_INVALID, LayerId, Opacity, Operation, Origin, Provenance, - RecordVersion, SemanticMetadata, ShapeCommand, ShapeCreateArgs, ShapeDeleteArgs, ShapeId, ShapeKind, ShapeParent, - ShapePatch, ShapePatchArgs, ShapeRecord, ShapeStyle, SiblingAnchor, Timestamp, Transform, Value, Vec2, Write, - anyhow, builtin_shape_kinds, + RecordVersion, Result, SemanticMetadata, ShapeCommand, ShapeCreateArgs, ShapeDeleteArgs, ShapeId, ShapeKind, + ShapeParent, ShapePatch, ShapePatchArgs, ShapeRecord, ShapeStyle, SiblingAnchor, Timestamp, Transform, Value, Vec2, + Write, anyhow, builtin_shape_kinds, }; -pub fn run_shape_command(command: ShapeCommand, json_output: bool, stdout: &mut dyn Write) -> Result<(), CliError> { +pub fn run_shape_command(command: ShapeCommand, json_output: bool, stdout: &mut dyn Write) -> Result<()> { match command { ShapeCommand::Create(args) => create_shape(args, json_output, stdout), ShapeCommand::Patch(args) => patch_shape(args, json_output, stdout), @@ -15,7 +15,7 @@ pub fn run_shape_command(command: ShapeCommand, json_output: bool, stdout: &mut } } -fn create_shape(args: ShapeCreateArgs, json_output: bool, stdout: &mut dyn Write) -> Result<(), CliError> { +fn create_shape(args: ShapeCreateArgs, json_output: bool, stdout: &mut dyn Write) -> Result<()> { if !args.x.is_finite() || !args.y.is_finite() || !args.rotation.is_finite() { return Err(CliError::new( EXIT_INVALID, @@ -81,7 +81,7 @@ fn create_shape(args: ShapeCreateArgs, json_output: bool, stdout: &mut dyn Write commit_mutation(&mut file, transaction, args.mutation.dry_run, json_output, stdout) } -fn patch_shape(args: ShapePatchArgs, json_output: bool, stdout: &mut dyn Write) -> Result<(), CliError> { +fn patch_shape(args: ShapePatchArgs, json_output: bool, stdout: &mut dyn Write) -> Result<()> { let patch_json = read_json_argument(&args.patch, "shape patch")?; let patch: ShapePatch = serde_json::from_str(&patch_json) .map_err(|error| CliError::new(EXIT_INVALID, error).context("could not parse ShapePatch JSON"))?; @@ -102,7 +102,7 @@ fn patch_shape(args: ShapePatchArgs, json_output: bool, stdout: &mut dyn Write) commit_mutation(&mut file, transaction, args.mutation.dry_run, json_output, stdout) } -fn delete_shape(args: ShapeDeleteArgs, json_output: bool, stdout: &mut dyn Write) -> Result<(), CliError> { +fn delete_shape(args: ShapeDeleteArgs, json_output: bool, stdout: &mut dyn Write) -> Result<()> { let mut file = open_document(&args.path)?; let shape_id = select_unique_shape( &mut file, diff --git a/crates/inkfinite-cli/src/cli/support.rs b/crates/inkfinite-cli/src/cli/support.rs index 50d6445..e492e9b 100644 --- a/crates/inkfinite-cli/src/cli/support.rs +++ b/crates/inkfinite-cli/src/cli/support.rs @@ -1,9 +1,9 @@ use super::{ ACTOR_ID, ActorId, Bounds, CliError, DocumentFile, EXIT_CONFLICT, EXIT_INPUT, EXIT_INVALID, EngineError, FileError, - Path, Write, io, + Path, Result, Write, io, }; -pub fn open_document(path: &Path) -> Result { +pub fn open_document(path: &Path) -> Result { DocumentFile::open(path, ActorId::from(ACTOR_ID)) .map_err(map_file_error) .map_err(|error| error.context(format!("could not open {}", portable_path(path)))) @@ -31,12 +31,12 @@ pub fn portable_path(path: &Path) -> String { path.to_string_lossy().replace('\\', "/") } -pub fn parse_bounds(value: &str) -> Result { +pub fn parse_bounds(value: &str) -> std::result::Result { let values = value .split(',') .map(str::trim) .map(str::parse::) - .collect::, _>>() + .collect::, _>>() .map_err(|_| "bounds must contain four numbers: x,y,width,height".to_owned())?; let [x, y, width, height] = values.as_slice() else { return Err("bounds must contain four numbers: x,y,width,height".into()); @@ -47,7 +47,7 @@ pub fn parse_bounds(value: &str) -> Result { Ok(Bounds { x: *x, y: *y, width: *width, height: *height }) } -pub fn write_heads(stdout: &mut dyn Write, heads: &[inkfinite_core::ChangeHash]) -> Result<(), CliError> { +pub fn write_heads(stdout: &mut dyn Write, heads: &[inkfinite_core::ChangeHash]) -> Result<()> { writeln!( stdout, "Heads: {}", @@ -60,7 +60,7 @@ pub fn write_heads(stdout: &mut dyn Write, heads: &[inkfinite_core::ChangeHash]) .map_err(map_output_error) } -pub fn write_json(stdout: &mut dyn Write, value: &impl serde::Serialize) -> Result<(), CliError> { +pub fn write_json(stdout: &mut dyn Write, value: &impl serde::Serialize) -> Result<()> { serde_json::to_writer_pretty(&mut *stdout, value).map_err(map_json_error)?; writeln!(stdout).map_err(map_output_error) } diff --git a/crates/inkfinite-cli/tests/cli.rs b/crates/inkfinite-cli/tests/cli.rs index 371fd38..09d6b56 100644 --- a/crates/inkfinite-cli/tests/cli.rs +++ b/crates/inkfinite-cli/tests/cli.rs @@ -3,9 +3,15 @@ use std::fs; use std::io::Write as _; use std::path::{Path, PathBuf}; use std::process::{Command, Output, Stdio}; +#[cfg(unix)] +use std::sync::Mutex; use std::sync::atomic::{AtomicU64, Ordering}; +#[cfg(unix)] +use std::thread; use inkfinite_core::file::DocumentFile; +#[cfg(unix)] +use inkfinite_core::ipc::{self, AppRequest, AppResponse, DiscoveryRecord, RequestEnvelope, ResponseEnvelope}; use inkfinite_core::proto::{Operation, TransactionDraft, TransactionId}; use inkfinite_core::{ ActorId, DocumentId, Opacity, Origin, Provenance, RecordVersion, SemanticMetadata, ShapeId, ShapeKind, ShapeParent, @@ -14,6 +20,8 @@ use inkfinite_core::{ use serde_json::Value; static TEMP_COUNTER: AtomicU64 = AtomicU64::new(0); +#[cfg(unix)] +static IPC_TEST_LOCK: Mutex<()> = Mutex::new(()); #[test] fn help_makes_common_tasks_and_support_paths_discoverable() { @@ -31,6 +39,12 @@ fn help_makes_common_tasks_and_support_paths_discoverable() { assert!(query_help.contains("--bounds ")); assert!(query_help.contains("inkfinite query architecture.inkfinite --role architecture.service --json")); + let app_query_help = run(["help", "app", "query"]); + assert_success(&app_query_help); + let app_query_help = String::from_utf8(app_query_help.stdout).unwrap(); + assert!(app_query_help.contains("inkfinite app query --role architecture.service --json")); + assert!(app_query_help.contains("--session-id ")); + let version = run(["--version"]); assert_success(&version); assert_eq!( @@ -200,6 +214,17 @@ fn schemas_and_capabilities_match_checked_in_contracts() { assert_eq!(capabilities_json["path_format"], "forward_slashes"); assert_eq!(capabilities_json["format"]["id"], "inkfinite.document"); assert_eq!(capabilities_json["protocol"]["id"], "inkfinite.protocol"); + assert!( + capabilities_json["commands"] + .as_array() + .unwrap() + .contains(&Value::String("app".into())) + ); + assert_eq!( + capabilities_json["live_mode"]["transport"], + "authenticated_local_socket" + ); + assert_eq!(capabilities_json["live_mode"]["tcp_or_http"], false); assert_eq!(capabilities_json["mutation_commands"]["shape"][0], "create"); assert!( capabilities_json["commands"] @@ -537,6 +562,107 @@ fn json_failures_keep_stdout_clean_and_use_stable_exit_codes() { assert!(usage.stdout.is_empty()); } +#[cfg(unix)] +#[test] +fn live_commands_read_shared_records_from_the_authenticated_local_server() { + use std::os::unix::fs::FileTypeExt; + use std::os::unix::net::{UnixListener, UnixStream}; + + use inkfinite_core::proto::QueryResult; + + let _guard = IPC_TEST_LOCK.lock().unwrap(); + ipc::ensure_ipc_directory().unwrap(); + let endpoint = ipc::endpoint_name(); + let endpoint_path = Path::new(&endpoint); + if UnixStream::connect(endpoint_path).is_ok() { + return; + } + if let Ok(metadata) = fs::symlink_metadata(endpoint_path) { + if !metadata.file_type().is_socket() { + return; + } + fs::remove_file(endpoint_path).unwrap(); + } + + let temporary = TestDirectory::new("live-ipc"); + let document_path = temporary.path.join("live.inkfinite"); + let document_id = DocumentId::from("document:live-ipc"); + let document = blank_document(&document_id, Some("Live")); + let mut file = DocumentFile::create(&document_path, document_id, ActorId::from("actor:test"), document).unwrap(); + let snapshot = file.snapshot().unwrap(); + drop(file); + + let previous_discovery = ipc::read_discovery(&ipc::discovery_path()).ok(); + let discovery = DiscoveryRecord { + protocol_id: inkfinite_core::proto::PROTOCOL_ID.into(), + version: inkfinite_core::proto::PROTOCOL_VERSION, + endpoint: endpoint.clone(), + token: "cli-ipc-test-token".into(), + }; + let listener = UnixListener::bind(endpoint_path).unwrap(); + listener.set_nonblocking(true).unwrap(); + ipc::write_discovery(&ipc::discovery_path(), &discovery).unwrap(); + + let server_snapshot = snapshot.clone(); + let server_token = discovery.token.clone(); + let server = thread::spawn(move || { + let runtime = tokio::runtime::Runtime::new().unwrap(); + runtime.block_on(async move { + let listener = tokio::net::UnixListener::from_std(listener).unwrap(); + for _ in 0..4 { + let (mut stream, _) = listener.accept().await.unwrap(); + let request = ipc::read_frame::(&mut stream).await.unwrap(); + assert_eq!(request.token, server_token); + let request_id = request.request_id; + let response = match request.request { + AppRequest::Status => AppResponse::Status(Vec::new()), + AppRequest::Inspect { .. } => AppResponse::Snapshot(server_snapshot.clone()), + AppRequest::Query { .. } => AppResponse::QueryResult(QueryResult { + heads: server_snapshot.heads.clone(), + records: Vec::new(), + bounds: BTreeMap::new(), + }), + AppRequest::Focus => AppResponse::Focused, + }; + ipc::write_frame(&mut stream, &ResponseEnvelope { request_id, result: Ok(response) }) + .await + .unwrap(); + } + }); + }); + + let status = run(["app", "status", "--json"]); + let inspected = run(["app", "inspect", "--json"]); + let queried = run(["app", "query", "--json"]); + let focused = run(["app", "focus"]); + server.join().unwrap(); + + let _ = ipc::remove_discovery(&ipc::discovery_path(), &discovery); + let _ = fs::remove_file(endpoint_path); + if let Some(previous) = previous_discovery { + ipc::write_discovery(&ipc::discovery_path(), &previous).unwrap(); + } + + assert_success(&status); + assert_eq!(parse_stdout(&status), Value::Array(Vec::new())); + assert_success(&inspected); + assert_eq!( + parse_stdout(&inspected), + serde_json::to_value(snapshot.clone()).unwrap() + ); + assert_success(&queried); + assert_eq!( + parse_stdout(&queried)["heads"], + serde_json::to_value(snapshot.heads).unwrap() + ); + assert_success(&focused); + assert!( + String::from_utf8(focused.stdout) + .unwrap() + .contains("Desktop focus requested") + ); +} + fn run(args: [&str; N]) -> Output { Command::new(env!("CARGO_BIN_EXE_inkfinite")) .args(args) diff --git a/crates/inkfinite-core/Cargo.toml b/crates/inkfinite-core/Cargo.toml index 049a86a..b86535c 100644 --- a/crates/inkfinite-core/Cargo.toml +++ b/crates/inkfinite-core/Cargo.toml @@ -12,6 +12,8 @@ perfect_freehand.workspace = true schemars.workspace = true serde.workspace = true serde_json.workspace = true +getrandom.workspace = true +tokio.workspace = true ts-rs.workspace = true thiserror.workspace = true diff --git a/crates/inkfinite-core/src/ipc/mod.rs b/crates/inkfinite-core/src/ipc/mod.rs index bcbadda..7d53b48 100644 --- a/crates/inkfinite-core/src/ipc/mod.rs +++ b/crates/inkfinite-core/src/ipc/mod.rs @@ -1,4 +1,717 @@ -//! Authenticated local IPC boundary. +//! Authenticated local IPC contracts, framing, discovery, and request dispatch. +//! +//! The transport stays deliberately small: one authenticated JSON request and +//! one response per local connection. Desktop owns the listener lifecycle and +//! frontend integration; this module owns the wire contract and shared session +//! behavior used by desktop and the CLI. + +use std::collections::{BTreeSet, VecDeque}; +use std::fmt::Write as _; +use std::fs::{self, OpenOptions}; +use std::io::{self, Write as _}; +use std::path::{Path, PathBuf}; + +use serde::de::DeserializeOwned; +use serde::{Deserialize, Serialize}; +use thiserror::Error; +use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt}; + +use crate::DocumentSnapshot; +use crate::file::FileError; +use crate::proto::{PROTOCOL_ID, PROTOCOL_VERSION, Query, QueryResult, SessionId}; +use crate::session::{SessionError, SessionService, SessionStatus}; pub use crate::engine::{CommitResult, TransactionDraft}; pub use crate::proto::{ProtocolError, Request, Response}; + +/// Largest accepted IPC payload, excluding the four-byte length prefix. +pub const MAX_FRAME_SIZE: usize = 1024 * 1024; + +/// Largest accepted discovery record. +pub const MAX_DISCOVERY_SIZE: usize = 16 * 1024; + +/// Tauri event emitted when a live client asks the desktop frontend to focus. +pub const FOCUS_EVENT: &str = "inkfinite-focus"; + +const REPLAY_WINDOW: usize = 4096; +const DISCOVERY_FILE: &str = "ipc.json"; +const IPC_DIRECTORY_PREFIX: &str = "inkfinite-"; + +/// Local endpoint and authentication material published by the desktop process. +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct DiscoveryRecord { + /// Stable protocol identifier. + pub protocol_id: String, + /// Wire protocol version accepted by the server. + pub version: u32, + /// Unix-domain socket path or Windows named-pipe name. + pub endpoint: String, + /// Random token scoped to this desktop process. + pub token: String, +} + +/// Read-only operation accepted by the V2-17 desktop server. +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case", tag = "type")] +#[allow(clippy::large_enum_variant)] +pub enum AppRequest { + /// List open desktop document sessions. + Status, + /// Return the current snapshot for an open session. + Inspect { + /// Session to inspect, or the only open session when omitted. + session_id: Option, + }, + /// Run the shared query implementation against an open session. + Query { + /// Session to query, or the only open session when omitted. + session_id: Option, + /// Shared semantic and hierarchy filters. + query: Query, + }, + /// Ask the desktop frontend to bring its main window forward. + Focus, +} + +/// Authenticated, replay-protected request frame. +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct RequestEnvelope { + /// Stable protocol identifier. + pub protocol_id: String, + /// Wire protocol version requested by the client. + pub version: u32, + /// Unique identifier used to reject replayed requests. + pub request_id: String, + /// Token read from the protected discovery record. + pub token: String, + /// Requested read-only operation. + pub request: AppRequest, +} + +/// Successful result from a V2-17 read-only app command. +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case", tag = "type", content = "value")] +pub enum AppResponse { + /// Current state of every open document session. + Status(Vec), + /// Shared materialized document snapshot. + Snapshot(DocumentSnapshot), + /// Shared deterministic query result. + QueryResult(QueryResult), + /// The focus notification was emitted. + Focused, +} + +/// Response frame correlated with one request. +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct ResponseEnvelope { + /// Request identifier supplied by the client. + pub request_id: String, + /// Successful response or stable protocol error. + pub result: Result, +} + +/// Local IPC framing, discovery, or connection failure. +#[derive(Debug, Error)] +pub enum IpcError { + /// A frame exceeded [`MAX_FRAME_SIZE`]. + #[error("IPC frame size {size} exceeds the {max} byte limit")] + FrameTooLarge { size: usize, max: usize }, + /// A discovery record exceeded [`MAX_DISCOVERY_SIZE`]. + #[error("IPC discovery record exceeds the {max} byte limit")] + DiscoveryTooLarge { max: usize }, + /// The peer closed the stream before sending the declared frame length. + #[error("IPC frame was truncated")] + TruncatedFrame, + /// A frame did not contain a valid JSON contract. + #[error("IPC frame contains malformed JSON: {0}")] + MalformedJson(#[from] serde_json::Error), + /// Discovery metadata was unavailable or invalid. + #[error("desktop app session is unavailable: {0}")] + Unavailable(String), + /// Local transport I/O failed. + #[error(transparent)] + Io(#[from] io::Error), +} + +/// Authentication and replay state retained for the lifetime of one server. +#[derive(Debug)] +pub struct RequestGuard { + token: String, + seen: BTreeSet, + order: VecDeque, +} + +impl RequestGuard { + /// Creates a guard for one process-scoped authentication token. + #[must_use] + pub fn new(token: String) -> Self { + Self { token, seen: BTreeSet::new(), order: VecDeque::new() } + } + + /// Checks the protocol, version, token, and replay identifier. + /// + /// # Errors + /// + /// Returns a stable [`ProtocolError`] without dispatching invalid requests. + pub fn validate(&mut self, request: &RequestEnvelope) -> Result<(), ProtocolError> { + if request.protocol_id != PROTOCOL_ID { + return Err(protocol_error( + "unsupported_protocol", + "unsupported IPC protocol identifier", + )); + } + if request.version != PROTOCOL_VERSION { + return Err(protocol_error( + "unsupported_version", + format!("unsupported IPC protocol version {}", request.version), + )); + } + if request.token != self.token { + return Err(protocol_error("unauthorized", "invalid IPC authentication token")); + } + if request.request_id.trim().is_empty() { + return Err(protocol_error("invalid_request_id", "IPC request ID must not be empty")); + } + if !self.seen.insert(request.request_id.clone()) { + return Err(protocol_error( + "replayed_request", + "IPC request ID has already been used", + )); + } + self.order.push_back(request.request_id.clone()); + if self.order.len() > REPLAY_WINDOW { + let Some(expired) = self.order.pop_front() else { return Ok(()) }; + self.seen.remove(&expired); + } + Ok(()) + } +} + +/// Reads one length-prefixed JSON value with strict truncation and size checks. +/// +/// # Errors +/// +/// Returns [`IpcError`] for I/O, size, truncation, or JSON failures. +pub async fn read_frame(stream: &mut (impl AsyncRead + Unpin)) -> Result { + let mut prefix = [0_u8; 4]; + read_exact_frame(stream, &mut prefix).await?; + let size = u32::from_be_bytes(prefix) as usize; + if size > MAX_FRAME_SIZE { + return Err(IpcError::FrameTooLarge { size, max: MAX_FRAME_SIZE }); + } + let mut payload = vec![0; size]; + read_exact_frame(stream, &mut payload).await?; + Ok(serde_json::from_slice(&payload)?) +} + +/// Writes one length-prefixed JSON value after enforcing the frame limit. +/// +/// # Errors +/// +/// Returns [`IpcError`] for serialization, size, or I/O failures. +pub async fn write_frame(stream: &mut (impl AsyncWrite + Unpin), value: &T) -> Result<(), IpcError> { + let payload = serde_json::to_vec(value)?; + if payload.len() > MAX_FRAME_SIZE { + return Err(IpcError::FrameTooLarge { size: payload.len(), max: MAX_FRAME_SIZE }); + } + let length = u32::try_from(payload.len()) + .map_err(|_| IpcError::FrameTooLarge { size: payload.len(), max: MAX_FRAME_SIZE })?; + stream.write_all(&length.to_be_bytes()).await?; + stream.write_all(&payload).await?; + stream.flush().await?; + Ok(()) +} + +/// Returns the per-user directory shared by the desktop and CLI. +#[must_use] +pub fn ipc_directory() -> PathBuf { + runtime_directory().join(format!("{IPC_DIRECTORY_PREFIX}{}", user_component())) +} + +/// Returns the protected discovery-file path shared by the desktop and CLI. +#[must_use] +pub fn discovery_path() -> PathBuf { + ipc_directory().join(DISCOVERY_FILE) +} + +/// Returns a platform-local endpoint name for this user. +#[must_use] +pub fn endpoint_name() -> String { + #[cfg(unix)] + { + ipc_directory().join("control.sock").to_string_lossy().into_owned() + } + #[cfg(windows)] + { + format!(r"\\.\pipe\inkfinite-{}", user_component()) + } +} + +/// Creates and protects the per-user IPC directory. +/// +/// # Errors +/// +/// Returns [`IpcError::Unavailable`] when the directory cannot be made private. +pub fn ensure_ipc_directory() -> Result { + let path = ipc_directory(); + ensure_private_directory(&path)?; + Ok(path) +} + +/// Generates a cryptographically random process token or request identifier. +/// +/// # Errors +/// +/// Returns the operating-system random source failure. +pub fn random_secret() -> Result { + let mut bytes = [0_u8; 32]; + getrandom::fill(&mut bytes)?; + let mut encoded = String::with_capacity(bytes.len() * 2); + for byte in bytes { + write!(&mut encoded, "{byte:02x}").expect("writing to a String cannot fail"); + } + Ok(encoded) +} + +/// Publishes a protected discovery record atomically. +/// +/// # Errors +/// +/// Returns [`IpcError`] when the record is invalid or cannot be written. +pub fn write_discovery(path: &Path, record: &DiscoveryRecord) -> Result<(), IpcError> { + validate_discovery_record(record)?; + let parent = path + .parent() + .ok_or_else(|| IpcError::Unavailable("IPC discovery path has no parent directory".into()))?; + ensure_private_directory(parent)?; + let payload = serde_json::to_vec(record)?; + if payload.len() > MAX_DISCOVERY_SIZE { + return Err(IpcError::DiscoveryTooLarge { max: MAX_DISCOVERY_SIZE }); + } + + let filename = path + .file_name() + .and_then(|name| name.to_str()) + .ok_or_else(|| IpcError::Unavailable("IPC discovery path has no valid filename".into()))?; + let temporary = parent.join(format!(".{filename}.{}.tmp", std::process::id())); + let write_result = write_private_file(&temporary, &payload).and_then(|()| { + #[cfg(windows)] + if path.exists() { + fs::remove_file(path)?; + } + fs::rename(&temporary, path).map_err(IpcError::from) + }); + if write_result.is_err() { + let _ = fs::remove_file(&temporary); + } + write_result +} + +/// Reads and validates a protected discovery record. +/// +/// # Errors +/// +/// Returns [`IpcError::Unavailable`] when no usable desktop server is published. +pub fn read_discovery(path: &Path) -> Result { + let metadata = fs::symlink_metadata(path).map_err(|error| IpcError::Unavailable(error.to_string()))?; + if !metadata.file_type().is_file() { + return Err(IpcError::Unavailable("IPC discovery path is not a regular file".into())); + } + if metadata.len() > MAX_DISCOVERY_SIZE as u64 { + return Err(IpcError::DiscoveryTooLarge { max: MAX_DISCOVERY_SIZE }); + } + let bytes = fs::read(path).map_err(|error| IpcError::Unavailable(error.to_string()))?; + if bytes.len() > MAX_DISCOVERY_SIZE { + return Err(IpcError::DiscoveryTooLarge { max: MAX_DISCOVERY_SIZE }); + } + let record: DiscoveryRecord = + serde_json::from_slice(&bytes).map_err(|error| IpcError::Unavailable(error.to_string()))?; + validate_discovery_record(&record)?; + Ok(record) +} + +/// Removes a discovery record only when it still belongs to this desktop session. +/// +/// # Errors +/// +/// Returns [`IpcError`] when the record cannot be read or removed. +pub fn remove_discovery(path: &Path, expected: &DiscoveryRecord) -> Result<(), IpcError> { + let current = match read_discovery(path) { + Ok(record) => record, + Err(IpcError::Unavailable(_)) if !path.exists() => return Ok(()), + Err(error) => return Err(error), + }; + if current == *expected { + fs::remove_file(path)?; + } + Ok(()) +} + +/// Dispatches one authenticated read-only app request through the shared session service. +/// +/// # Errors +/// +/// Returns a stable [`ProtocolError`] for unavailable sessions or document +/// failures. No request in this module mutates a document. +pub fn dispatch(service: &mut SessionService, request: AppRequest) -> Result { + match request { + AppRequest::Status => service + .statuses() + .map(AppResponse::Status) + .map_err(|error| session_protocol_error(&error)), + AppRequest::Inspect { session_id } => { + let session_id = service + .resolve_session_id(session_id.as_ref()) + .map_err(|error| session_protocol_error(&error))?; + service + .status(&session_id) + .map(|status| AppResponse::Snapshot(status.snapshot)) + .map_err(|error| session_protocol_error(&error)) + } + AppRequest::Query { session_id, query } => { + let session_id = service + .resolve_session_id(session_id.as_ref()) + .map_err(|error| session_protocol_error(&error))?; + service + .query(&session_id, &query) + .map(AppResponse::QueryResult) + .map_err(|error| session_protocol_error(&error)) + } + AppRequest::Focus => Ok(AppResponse::Focused), + } +} + +/// Converts a session failure into the shared protocol error contract. +#[must_use] +pub fn session_protocol_error(error: &SessionError) -> ProtocolError { + let code = match &error { + SessionError::NotFound(_) => "session_not_found", + SessionError::SessionSelectionRequired { open_sessions: 0 } => "app_session_unavailable", + SessionError::SessionSelectionRequired { .. } => "session_selection_required", + SessionError::ActorMismatch { .. } => "actor_mismatch", + SessionError::StaleHeads => "stale_heads", + SessionError::AlreadyOpen { .. } => "document_already_open", + SessionError::File(file_error) => match file_error { + FileError::Locked { .. } => "document_locked", + FileError::AlreadyExists { .. } => "document_already_exists", + FileError::InvalidV1(_) + | FileError::Json(_) + | FileError::UnsupportedFormat { .. } + | FileError::UnsupportedShapeKind { .. } + | FileError::SamePath { .. } + | FileError::RecoveryNotFound { .. } + | FileError::InvalidRecovery(_) + | FileError::RecoveryAhead { .. } + | FileError::Engine(_) + | FileError::Io { .. } => "document_file_error", + }, + SessionError::Engine(_) => "document_engine_error", + }; + protocol_error(code, error.to_string()) +} + +/// Sends one request to the currently published desktop process. +/// +/// # Errors +/// +/// Returns discovery, connection, framing, or response-correlation failures. +pub async fn send(request: AppRequest) -> Result { + let discovery = read_discovery(&discovery_path())?; + if discovery.endpoint != endpoint_name() { + return Err(IpcError::Unavailable( + "IPC discovery endpoint is not the per-user Inkfinite endpoint".into(), + )); + } + let request_id = random_secret().map_err(|error| IpcError::Unavailable(error.to_string()))?; + let envelope = RequestEnvelope { + protocol_id: PROTOCOL_ID.into(), + version: PROTOCOL_VERSION, + request_id: request_id.clone(), + token: discovery.token, + request, + }; + let response = send_to_endpoint(&discovery.endpoint, &envelope).await?; + if response.request_id != request_id { + return Err(IpcError::Unavailable( + "desktop response used the wrong request ID".into(), + )); + } + Ok(response) +} + +#[cfg(unix)] +async fn send_to_endpoint(endpoint: &str, request: &RequestEnvelope) -> Result { + let mut stream = tokio::net::UnixStream::connect(endpoint).await?; + write_frame(&mut stream, request).await?; + read_frame(&mut stream).await +} + +#[cfg(windows)] +async fn send_to_endpoint(endpoint: &str, request: &RequestEnvelope) -> Result { + use tokio::net::windows::named_pipe::ClientOptions; + + let mut stream = ClientOptions::new().open(endpoint)?; + write_frame(&mut stream, request).await?; + read_frame(&mut stream).await +} + +fn runtime_directory() -> PathBuf { + #[cfg(unix)] + if let Some(path) = std::env::var_os("XDG_RUNTIME_DIR").filter(|value| !value.is_empty()) { + return PathBuf::from(path); + } + std::env::temp_dir() +} + +fn user_component() -> String { + let user = if cfg!(windows) { + std::env::var("USERNAME").ok() + } else { + std::env::var("USER").ok().or_else(|| std::env::var("USERNAME").ok()) + }; + let sanitized: String = user + .unwrap_or_else(|| "user".into()) + .chars() + .filter(|character| character.is_ascii_alphanumeric() || matches!(character, '-' | '_')) + .take(32) + .collect(); + if sanitized.is_empty() { "user".into() } else { sanitized } +} + +fn validate_discovery_record(record: &DiscoveryRecord) -> Result<(), IpcError> { + if record.protocol_id != PROTOCOL_ID { + return Err(IpcError::Unavailable( + "discovery record has an unsupported protocol identifier".into(), + )); + } + if record.version != PROTOCOL_VERSION { + return Err(IpcError::Unavailable( + "discovery record has an unsupported protocol version".into(), + )); + } + if record.endpoint.trim().is_empty() { + return Err(IpcError::Unavailable("discovery record has an empty endpoint".into())); + } + if record.token.trim().is_empty() { + return Err(IpcError::Unavailable( + "discovery record has an empty authentication token".into(), + )); + } + Ok(()) +} + +fn ensure_private_directory(path: &Path) -> Result<(), IpcError> { + fs::create_dir_all(path)?; + let metadata = fs::symlink_metadata(path)?; + if !metadata.file_type().is_dir() { + return Err(IpcError::Unavailable(format!( + "IPC path is not a directory: {}", + path.display() + ))); + } + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + + let mut permissions = metadata.permissions(); + permissions.set_mode(0o700); + fs::set_permissions(path, permissions)?; + } + Ok(()) +} + +fn write_private_file(path: &Path, payload: &[u8]) -> Result<(), IpcError> { + let mut options = OpenOptions::new(); + options.write(true).create_new(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + + options.mode(0o600); + } + let mut file = options.open(path)?; + file.write_all(payload)?; + file.sync_all()?; + Ok(()) +} + +async fn read_exact_frame(stream: &mut (impl AsyncRead + Unpin), bytes: &mut [u8]) -> Result<(), IpcError> { + match stream.read_exact(bytes).await { + Ok(_) => Ok(()), + Err(error) if error.kind() == io::ErrorKind::UnexpectedEof => Err(IpcError::TruncatedFrame), + Err(error) => Err(IpcError::Io(error)), + } +} + +fn protocol_error(code: &str, message: impl Into) -> ProtocolError { + ProtocolError { code: code.into(), message: message.into(), details: None } +} + +#[cfg(test)] +mod tests { + use std::path::PathBuf; + use std::sync::atomic::{AtomicU64, Ordering}; + + use tokio::io::AsyncWriteExt as _; + + use super::*; + use crate::proto::Query; + + static TEST_COUNTER: AtomicU64 = AtomicU64::new(0); + + fn request(token: &str, request_id: &str) -> RequestEnvelope { + RequestEnvelope { + protocol_id: PROTOCOL_ID.into(), + version: PROTOCOL_VERSION, + request_id: request_id.into(), + token: token.into(), + request: AppRequest::Status, + } + } + + #[test] + fn guard_rejects_wrong_tokens_versions_and_replays() { + let mut guard = RequestGuard::new("secret".into()); + assert_eq!( + guard.validate(&request("wrong", "one")).unwrap_err().code, + "unauthorized" + ); + + let mut unsupported_protocol = request("secret", "protocol"); + unsupported_protocol.protocol_id = "other.protocol".into(); + assert_eq!( + guard.validate(&unsupported_protocol).unwrap_err().code, + "unsupported_protocol" + ); + + let mut unsupported = request("secret", "two"); + unsupported.version += 1; + assert_eq!(guard.validate(&unsupported).unwrap_err().code, "unsupported_version"); + + let valid = request("secret", "three"); + guard.validate(&valid).unwrap(); + assert_eq!(guard.validate(&valid).unwrap_err().code, "replayed_request"); + } + + #[tokio::test] + async fn framing_rejects_oversized_truncated_and_malformed_frames() { + let (mut writer, mut reader) = tokio::io::duplex(64); + writer + .write_all( + &u32::try_from(MAX_FRAME_SIZE + 1) + .expect("test frame size fits") + .to_be_bytes(), + ) + .await + .unwrap(); + assert!(matches!( + read_frame::(&mut reader).await, + Err(IpcError::FrameTooLarge { .. }) + )); + + let (mut writer, mut reader) = tokio::io::duplex(64); + writer.write_all(&8_u32.to_be_bytes()).await.unwrap(); + writer.write_all(b"{}").await.unwrap(); + drop(writer); + assert!(matches!( + read_frame::(&mut reader).await, + Err(IpcError::TruncatedFrame) + )); + + let (mut writer, mut reader) = tokio::io::duplex(64); + writer.write_all(&2_u32.to_be_bytes()).await.unwrap(); + writer.write_all(b"{]").await.unwrap(); + assert!(matches!( + read_frame::(&mut reader).await, + Err(IpcError::MalformedJson(_)) + )); + } + + #[test] + fn discovery_is_atomic_and_rejects_invalid_records() { + let root = test_directory(); + let path = root.join("ipc.json"); + let record = DiscoveryRecord { + protocol_id: PROTOCOL_ID.into(), + version: PROTOCOL_VERSION, + endpoint: endpoint_name(), + token: "secret".into(), + }; + write_discovery(&path, &record).unwrap(); + assert_eq!(read_discovery(&path).unwrap(), record); + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + + assert_eq!(std::fs::metadata(&root).unwrap().permissions().mode() & 0o777, 0o700); + assert_eq!(std::fs::metadata(&path).unwrap().permissions().mode() & 0o777, 0o600); + } + + let malformed = root.join("malformed.json"); + std::fs::write(&malformed, b"[]").unwrap(); + assert!(matches!(read_discovery(&malformed), Err(IpcError::Unavailable(_)))); + + remove_discovery(&path, &record).unwrap(); + assert!(!path.exists()); + remove_test_directory(root); + } + + #[test] + fn dispatch_returns_shared_records_and_reports_unavailable_sessions() { + let root = test_directory(); + let path = root.join("board.inkfinite"); + let mut service = SessionService::new(); + let opened = service + .create( + &path, + crate::DocumentId::from("document:ipc"), + crate::ActorId::from("actor:test"), + None, + ) + .unwrap(); + + let status = dispatch(&mut service, AppRequest::Status).unwrap(); + assert!(matches!(status, AppResponse::Status(statuses) if statuses.len() == 1)); + let snapshot = dispatch(&mut service, AppRequest::Inspect { session_id: None }).unwrap(); + assert!(matches!(snapshot, AppResponse::Snapshot(snapshot) if snapshot == opened.status.snapshot)); + let query = dispatch( + &mut service, + AppRequest::Query { session_id: None, query: Query::default() }, + ) + .unwrap(); + assert!(matches!(query, AppResponse::QueryResult(result) if result.heads == opened.status.snapshot.heads)); + + service.close(&opened.session_id).unwrap(); + let error = dispatch(&mut service, AppRequest::Inspect { session_id: None }).unwrap_err(); + assert_eq!(error.code, "app_session_unavailable"); + remove_test_directory(root); + } + + #[test] + fn endpoint_is_scoped_to_a_user_directory() { + let directory = ipc_directory(); + assert!( + directory + .file_name() + .and_then(|name| name.to_str()) + .is_some_and(|name| name.starts_with(IPC_DIRECTORY_PREFIX)) + ); + assert!(endpoint_name().contains(directory.to_string_lossy().as_ref())); + } + + fn test_directory() -> PathBuf { + let id = TEST_COUNTER.fetch_add(1, Ordering::Relaxed); + let path = std::env::temp_dir().join(format!("inkfinite-ipc-test-{id}")); + let _ = std::fs::remove_dir_all(&path); + std::fs::create_dir_all(&path).unwrap(); + path + } + + fn remove_test_directory(path: PathBuf) { + let _ = std::fs::remove_dir_all(path); + } +} diff --git a/crates/inkfinite-core/src/session.rs b/crates/inkfinite-core/src/session.rs index 2821345..be2a7f1 100644 --- a/crates/inkfinite-core/src/session.rs +++ b/crates/inkfinite-core/src/session.rs @@ -1,8 +1,10 @@ -//! Rust-owned document sessions used by desktop and local adapters. +//! Document sessions used by desktop and local adapters. //! -//! A session keeps the durable file boundary, the materialized snapshot, and -//! actor-scoped history together. Adapters can expose this service over Tauri, -//! local IPC, or a CLI without moving document bytes into the frontend. +//! A session keeps the file, the materialized snapshot, and actor-scoped +//! history together. +//! +//! Adapters can expose this service over Tauri, local IPC, or a CLI without +//! moving document bytes into the frontend. use std::collections::BTreeMap; use std::path::{Path, PathBuf}; @@ -82,6 +84,12 @@ pub enum SessionError { /// No open session has the requested identifier. #[error("session not found: {0:?}")] NotFound(SessionId), + /// A live command omitted its session while zero or multiple sessions exist. + #[error("session selection required; {open_sessions} sessions are open")] + SessionSelectionRequired { + /// Number of sessions currently available for selection. + open_sessions: usize, + }, /// The actor does not own the session's local mutation stream. #[error("actor {actual} does not own session actor {expected}")] ActorMismatch { @@ -179,6 +187,42 @@ impl SessionService { session.status(session_id) } + /// Returns current state for every open session in stable identifier order. + /// + /// # Errors + /// + /// Returns a snapshot or recovery error from any open session. + pub fn statuses(&mut self) -> Result, SessionError> { + self.sessions + .iter_mut() + .map(|(session_id, session)| session.status(session_id)) + .collect() + } + + /// Resolves an explicit session, or the only open session when omitted. + /// + /// # Errors + /// + /// Returns a typed error when no session is open or more than one session + /// requires the caller to choose explicitly. + pub fn resolve_session_id(&self, requested: Option<&SessionId>) -> Result { + if let Some(session_id) = requested { + if self.sessions.contains_key(session_id) { + return Ok(session_id.clone()); + } + return Err(SessionError::NotFound(session_id.clone())); + } + match self.sessions.len() { + 1 => self + .sessions + .keys() + .next() + .cloned() + .ok_or(SessionError::SessionSelectionRequired { open_sessions: 0 }), + _ => Err(SessionError::SessionSelectionRequired { open_sessions: self.sessions.len() }), + } + } + /// Commits one actor-owned transaction and returns its materialized patch. /// /// # Errors