diff --git a/cmux-tui/bindings/conformance/runner.py b/cmux-tui/bindings/conformance/runner.py index b1642fc88fc1..7f1832f20c52 100644 --- a/cmux-tui/bindings/conformance/runner.py +++ b/cmux-tui/bindings/conformance/runner.py @@ -29,7 +29,7 @@ BUILD = HERE / ".build" / "resource-v2" LANGUAGES = ("python", "typescript", "rust", "go", "java", "cpp", "zig") PROTOCOL = "cmux.protocol/2" -TRANSPORTED_OPERATION_COUNT = 125 +TRANSPORTED_OPERATION_COUNT = 126 MAX_REQUEST_BYTES = 4 * 1024 * 1024 MAX_STREAM_MESSAGES = 256 MAX_STREAM_BYTES = 16 * 1024 * 1024 diff --git a/cmux-tui/bindings/cpp/.cmux-resource-api.json b/cmux-tui/bindings/cpp/.cmux-resource-api.json index 101c7568d031..190e3a4058ec 100644 --- a/cmux-tui/bindings/cpp/.cmux-resource-api.json +++ b/cmux-tui/bindings/cpp/.cmux-resource-api.json @@ -1,5 +1,5 @@ { - "catalog_sha256": "2479d038387ac83af736875b3314a4b33d867243191a50bab477d0a723bc7ab0", + "catalog_sha256": "051b96446a4622e99f979f49ff1a17db214feaff318e40a315e34ef3533c10df", "operations": { "agent.list": { "class": "read" @@ -85,6 +85,9 @@ "machine.list": { "class": "read" }, + "notification.ack": { + "class": "mutation" + }, "notification.create": { "class": "mutation" }, diff --git a/cmux-tui/bindings/go/.cmux-resource-api.json b/cmux-tui/bindings/go/.cmux-resource-api.json index 101c7568d031..190e3a4058ec 100644 --- a/cmux-tui/bindings/go/.cmux-resource-api.json +++ b/cmux-tui/bindings/go/.cmux-resource-api.json @@ -1,5 +1,5 @@ { - "catalog_sha256": "2479d038387ac83af736875b3314a4b33d867243191a50bab477d0a723bc7ab0", + "catalog_sha256": "051b96446a4622e99f979f49ff1a17db214feaff318e40a315e34ef3533c10df", "operations": { "agent.list": { "class": "read" @@ -85,6 +85,9 @@ "machine.list": { "class": "read" }, + "notification.ack": { + "class": "mutation" + }, "notification.create": { "class": "mutation" }, diff --git a/cmux-tui/bindings/java/.cmux-resource-api.json b/cmux-tui/bindings/java/.cmux-resource-api.json index 101c7568d031..190e3a4058ec 100644 --- a/cmux-tui/bindings/java/.cmux-resource-api.json +++ b/cmux-tui/bindings/java/.cmux-resource-api.json @@ -1,5 +1,5 @@ { - "catalog_sha256": "2479d038387ac83af736875b3314a4b33d867243191a50bab477d0a723bc7ab0", + "catalog_sha256": "051b96446a4622e99f979f49ff1a17db214feaff318e40a315e34ef3533c10df", "operations": { "agent.list": { "class": "read" @@ -85,6 +85,9 @@ "machine.list": { "class": "read" }, + "notification.ack": { + "class": "mutation" + }, "notification.create": { "class": "mutation" }, diff --git a/cmux-tui/bindings/python/.cmux-resource-api.json b/cmux-tui/bindings/python/.cmux-resource-api.json index 101c7568d031..190e3a4058ec 100644 --- a/cmux-tui/bindings/python/.cmux-resource-api.json +++ b/cmux-tui/bindings/python/.cmux-resource-api.json @@ -1,5 +1,5 @@ { - "catalog_sha256": "2479d038387ac83af736875b3314a4b33d867243191a50bab477d0a723bc7ab0", + "catalog_sha256": "051b96446a4622e99f979f49ff1a17db214feaff318e40a315e34ef3533c10df", "operations": { "agent.list": { "class": "read" @@ -85,6 +85,9 @@ "machine.list": { "class": "read" }, + "notification.ack": { + "class": "mutation" + }, "notification.create": { "class": "mutation" }, diff --git a/cmux-tui/bindings/rust/.cmux-resource-api.json b/cmux-tui/bindings/rust/.cmux-resource-api.json index 101c7568d031..190e3a4058ec 100644 --- a/cmux-tui/bindings/rust/.cmux-resource-api.json +++ b/cmux-tui/bindings/rust/.cmux-resource-api.json @@ -1,5 +1,5 @@ { - "catalog_sha256": "2479d038387ac83af736875b3314a4b33d867243191a50bab477d0a723bc7ab0", + "catalog_sha256": "051b96446a4622e99f979f49ff1a17db214feaff318e40a315e34ef3533c10df", "operations": { "agent.list": { "class": "read" @@ -85,6 +85,9 @@ "machine.list": { "class": "read" }, + "notification.ack": { + "class": "mutation" + }, "notification.create": { "class": "mutation" }, diff --git a/cmux-tui/bindings/typescript/.cmux-resource-api.json b/cmux-tui/bindings/typescript/.cmux-resource-api.json index 101c7568d031..190e3a4058ec 100644 --- a/cmux-tui/bindings/typescript/.cmux-resource-api.json +++ b/cmux-tui/bindings/typescript/.cmux-resource-api.json @@ -1,5 +1,5 @@ { - "catalog_sha256": "2479d038387ac83af736875b3314a4b33d867243191a50bab477d0a723bc7ab0", + "catalog_sha256": "051b96446a4622e99f979f49ff1a17db214feaff318e40a315e34ef3533c10df", "operations": { "agent.list": { "class": "read" @@ -85,6 +85,9 @@ "machine.list": { "class": "read" }, + "notification.ack": { + "class": "mutation" + }, "notification.create": { "class": "mutation" }, diff --git a/cmux-tui/bindings/zig/.cmux-resource-api.json b/cmux-tui/bindings/zig/.cmux-resource-api.json index 101c7568d031..190e3a4058ec 100644 --- a/cmux-tui/bindings/zig/.cmux-resource-api.json +++ b/cmux-tui/bindings/zig/.cmux-resource-api.json @@ -1,5 +1,5 @@ { - "catalog_sha256": "2479d038387ac83af736875b3314a4b33d867243191a50bab477d0a723bc7ab0", + "catalog_sha256": "051b96446a4622e99f979f49ff1a17db214feaff318e40a315e34ef3533c10df", "operations": { "agent.list": { "class": "read" @@ -85,6 +85,9 @@ "machine.list": { "class": "read" }, + "notification.ack": { + "class": "mutation" + }, "notification.create": { "class": "mutation" }, diff --git a/cmux-tui/crates/cmux-tui-core/src/mux.rs b/cmux-tui/crates/cmux-tui-core/src/mux.rs index 25c6c5390d66..a46d5bbe22dd 100644 --- a/cmux-tui/crates/cmux-tui-core/src/mux.rs +++ b/cmux-tui/crates/cmux-tui-core/src/mux.rs @@ -8,7 +8,7 @@ mod resource_topology; pub(crate) use resource_content::ResourceEffectProjection; use public_projections::{RestoredPublicProjections, restore_public_projections}; -use std::collections::{HashMap, HashSet, VecDeque}; +use std::collections::{BTreeSet, HashMap, HashSet, VecDeque}; use std::fmt; use std::ops::{Deref, DerefMut}; use std::path::Path; @@ -1041,6 +1041,18 @@ pub struct TreeDelta { pub workspace_revision: Option, } +/// A durable client install identity: non-empty, at most 128 ASCII graphic +/// bytes. Shared by per-client focus memory and per-client notification reads. +pub(crate) fn validate_client_id(client_id: &str) -> anyhow::Result<()> { + if client_id.is_empty() + || client_id.len() > 128 + || !client_id.bytes().all(|byte| byte.is_ascii_graphic()) + { + anyhow::bail!("bad request: invalid client_id"); + } + Ok(()) +} + #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum NotificationLevel { Info, @@ -1173,6 +1185,49 @@ fn agent_state_for_hook_kind(kind: &str) -> Option { }) } +/// Title, body, and level for the notification an agent hook event earns, +/// or `None` for transitions that need no attention (session start, turn +/// start, child lifecycle, session end). +fn agent_hook_notification( + ingress: &crate::JournalIngress, +) -> Option<(String, String, NotificationLevel)> { + const BODY_MAX_CHARS: usize = 512; + let (verb, level) = match ingress.kind.as_str() { + "agent.turn.completed" => ("finished", NotificationLevel::Info), + "agent.approval.requested" => ("needs approval", NotificationLevel::Warning), + "agent.question.requested" => ("asked a question", NotificationLevel::Warning), + "agent.plan_review.requested" => ("requested plan review", NotificationLevel::Warning), + "agent.error.reported" => ("reported an error", NotificationLevel::Error), + _ => return None, + }; + let adapter = ingress + .payload + .get("adapter") + .and_then(|adapter| adapter.get("id")) + .and_then(Value::as_str) + .filter(|id| !id.is_empty()) + .unwrap_or("Agent"); + let mut agent = String::with_capacity(adapter.len()); + let mut chars = adapter.chars(); + if let Some(first) = chars.next() { + agent.extend(first.to_uppercase()); + agent.push_str(chars.as_str()); + } + // Prompt and message text is redacted before the journal accepts it, so + // the body is the one structural field a viewer can act on: the tool an + // approval is waiting on. Everything else stays empty rather than leaking + // a redaction marker into the notification feed. + let normalized = ingress.payload.get("normalized"); + let body = ["tool_name"] + .into_iter() + .filter_map(|field| normalized.and_then(|value| value.get(field)).and_then(Value::as_str)) + .map(str::trim) + .find(|text| !text.is_empty() && *text != "[redacted]") + .map(|text| text.chars().take(BODY_MAX_CHARS).collect::()) + .unwrap_or_default(); + Some((format!("{agent} {verb}"), body, level)) +} + #[derive(Debug, Clone)] pub struct AgentRecord { pub surface: SurfaceId, @@ -2251,6 +2306,16 @@ pub struct Mux { placement_notifications: Mutex>, terminal_notifications: Mutex>, notification_ledger: Mutex>, + /// Per-client read marks. The shared unread marker above answers "does + /// this terminal need attention on the shared console"; this map answers + /// "has this client install seen this notification", so several remote + /// clients of one session keep independent unread state. + notification_reads: Mutex>>, + /// Notification ids the in-memory ledger evicted whose durable read marks + /// are still to be pruned. Pruning happens only after a create commits, + /// and only for ids the committed receipts no longer retain, so a failed + /// create cannot orphan marks the next restart would rebuild. + notification_read_prunes: Mutex>, resource_machine_service: OnceLock>, journal_kernel: Arc, journal_ingress: crate::journal_ingress::JournalIngressSender, @@ -2521,6 +2586,7 @@ impl Mux { agent_hook_fences, terminal_notifications, notification_ledger, + notification_reads, } = restore_public_projections(&state, registry.public_projections()?)?; let journal_producers = registry.journal_producer_manifests()?; let session_public_id = registry.session_id().clone(); @@ -2639,6 +2705,8 @@ impl Mux { placement_notifications: Mutex::new(HashMap::new()), terminal_notifications: Mutex::new(terminal_notifications), notification_ledger: Mutex::new(notification_ledger), + notification_reads: Mutex::new(notification_reads), + notification_read_prunes: Mutex::new(Vec::new()), resource_machine_service: OnceLock::new(), journal_kernel, journal_ingress, @@ -5739,6 +5807,21 @@ impl Mux { ) else { return Ok(()); }; + // Attention-worthy transitions become durable notifications before + // the agent report commits. The notification key is derived from the + // journal sequence, so a retry after a crash between the two commits + // replays the notification instead of posting it twice, and the fence + // (stored by the report) still advances exactly once. + if let Some((title, body, level)) = agent_hook_notification(ingress) { + self.create_durable_notification( + &format!("agent-hook-notification-{sequence}"), + title, + body, + level, + Some(surface), + ) + .with_context(|| format!("agent hook notification for sequence {sequence}"))?; + } // The record's session field is a human-facing label; native agent // session ids are opaque, so views fall back to their own context. let marker = if state == AgentState::Done { @@ -9048,6 +9131,9 @@ impl Mux { } } + /// Post a notification from the legacy `notify` verb. This is the same + /// durable path as `notification.create`, under a fresh key, so remote + /// subscribers of the resource feed and a restarted daemon see it too. pub fn post_notification( &self, title: String, @@ -9055,24 +9141,9 @@ impl Mux { level: NotificationLevel, surface: Option, ) -> anyhow::Result { - let public_id = NotificationPublicId::random()?; - let terminal_id = surface.and_then(|surface| { - let state = self.state.lock().unwrap(); - state - .surfaces - .get(&surface) - .or_else(|| state.terminal_runtime_by_id(surface)) - .and_then(|surface| surface.terminal_public_id().cloned()) - }); - Ok(self.post_resource_notification( - public_id, - title, - body, - level, - surface, - terminal_id, - now_ms(), - )) + let key = format!("notify-{}", crate::workspace_registry::new_uuid_v4()); + self.create_durable_notification(&key, title, body, level, surface)? + .context("fresh notify key unexpectedly replayed") } #[allow(clippy::too_many_arguments)] @@ -9099,8 +9170,18 @@ impl Mux { created_at_ms, surface, }); + let mut evicted = Vec::new(); while ledger.len() > NOTIFICATION_LEDGER_CAPACITY { - ledger.pop_front(); + if let Some(old) = ledger.pop_front() { + evicted.push(old.id); + } + } + if !evicted.is_empty() { + let mut reads = self.notification_reads.lock().unwrap(); + for id in &evicted { + reads.remove(id); + } + self.notification_read_prunes.lock().unwrap().extend(evicted); } } // Shared topology focus is only a default projection. A frontend must @@ -9162,6 +9243,279 @@ impl Mux { notifications } + /// Client ids that acknowledged `notification`, sorted and unique. + pub fn notification_read_by(&self, notification: &NotificationPublicId) -> Vec { + self.notification_reads + .lock() + .unwrap() + .get(notification) + .map(|clients| clients.iter().cloned().collect()) + .unwrap_or_default() + } + + /// The public snapshot row for one ledger entry. Every producer of a + /// `notification` resource value goes through here so create results, + /// acknowledgement deltas, and session snapshots cannot drift. + pub(crate) fn notification_snapshot_value( + &self, + notification: &ResourceNotification, + session_id: &SessionPublicId, + read_by: &[String], + ) -> Value { + let unread = notification + .terminal_id + .as_ref() + .and_then(|terminal_id| self.terminal_notification(terminal_id)) + .is_some_and(|marker| marker.unread); + let mut value = serde_json::json!({ + "id": notification.id, + "session_id": session_id, + "title": notification.title, + "body": notification.body, + "level": notification.level.as_str(), + "created_at_ms": notification.created_at_ms.to_string(), + "unread": unread, + "read_by": read_by, + }); + if let Some(terminal_id) = ¬ification.terminal_id { + value["terminal_id"] = serde_json::json!(terminal_id); + } + value + } + + /// Post a notification through the durable `notification.create` effect + /// path under a caller-owned idempotency key. Replaying the same key + /// returns `None` without posting again, so journal-driven producers + /// (agent hooks) and the legacy `notify` verb share one durable ledger + /// with the resource API and survive a daemon restart. A fresh post + /// returns the session-local legacy notification id. + pub(crate) fn create_durable_notification( + &self, + idempotency_key: &str, + title: String, + body: String, + level: NotificationLevel, + surface: Option, + ) -> anyhow::Result> { + const OPERATION: &str = "notification.create"; + let terminal_id = surface.and_then(|surface| { + let state = self.state.lock().unwrap(); + state + .surfaces + .get(&surface) + .or_else(|| state.terminal_runtime_by_id(surface)) + .and_then(|surface| surface.terminal_public_id().cloned()) + }); + let fingerprint = serde_json::json!({ + "operation": OPERATION, + "origin": "durable-notification", + "title": title, + "body": body, + "level": level.as_str(), + "terminal_id": terminal_id, + }); + let committed = |outcome: ResourceEffectOutcome| match outcome { + ResourceEffectOutcome::Success(_) => Ok(None), + ResourceEffectOutcome::Failure(error) => Err(anyhow::Error::new(error)), + }; + let preparation = match self.lookup_resource_effect( + idempotency_key, + OPERATION, + &fingerprint, + )? { + Some(preparation) => preparation, + None => { + let intent = serde_json::json!({ + "notification_id": NotificationPublicId::random().map_err(anyhow::Error::new)?, + "title": title, + "body": body, + "level": level.as_str(), + "terminal_id": terminal_id, + "created_at_ms": now_ms(), + }); + self.prepare_resource_effect( + idempotency_key, + OPERATION, + &fingerprint, + &intent, + None, + None, + )? + } + }; + let intent = match preparation { + ResourceEffectPreparation::Committed { outcome, .. } => { + return committed(outcome); + } + ResourceEffectPreparation::Indeterminate => { + // The post may or may not have happened before a crash. A + // notification is advisory, so the caller proceeds without it + // rather than retrying the same key forever; the agent-hook + // fold in particular must still commit its report and fence. + self.report_internal_diagnostic("notification effect indeterminate, skipped"); + return Ok(None); + } + ResourceEffectPreparation::Execute { .. } => { + self.mark_resource_effect_executing(idempotency_key, OPERATION, &fingerprint)? + } + }; + let notification_id: NotificationPublicId = + serde_json::from_value(intent["notification_id"].clone()) + .context("stored notification intent has an invalid identity")?; + let created_at_ms = intent + .get("created_at_ms") + .and_then(Value::as_u64) + .context("stored notification intent has an invalid timestamp")?; + let numeric_id = self.post_resource_notification( + notification_id.clone(), + title.clone(), + body.clone(), + level, + surface, + terminal_id.clone(), + created_at_ms, + ); + let session_id = self.workspace_registry.lock().unwrap().session_id().clone(); + let value = self.notification_snapshot_value( + &ResourceNotification { + id: notification_id.clone(), + title, + body, + level, + terminal_id, + created_at_ms, + surface, + }, + &session_id, + &[], + ); + let outcome = ResourceEffectOutcome::Success(value.clone()); + let deltas = serde_json::json!([{ + "kind":"upsert", + "sequence":0, + "resource":"notification", + "id":notification_id, + "value":value, + }]); + if let Err(error) = self.commit_resource_effect( + idempotency_key, + OPERATION, + &fingerprint, + &outcome, + Some(&deltas), + ) { + let _ = self.mark_resource_effect_indeterminate(idempotency_key); + return Err(error.context("notification effect commit failed")); + } + self.prune_evicted_notification_reads(); + Ok(Some(numeric_id)) + } + + /// Delete durable read marks for evicted notifications that the committed + /// receipts no longer retain. Called after a notification create commits; + /// ids still retained durably stay queued for a later create. + pub(crate) fn prune_evicted_notification_reads(&self) { + let candidates = std::mem::take(&mut *self.notification_read_prunes.lock().unwrap()); + if candidates.is_empty() { + return; + } + let remaining = + match self.workspace_registry.lock().unwrap().prune_notification_reads(&candidates) { + Ok(remaining) => remaining, + Err(_) => { + self.report_internal_diagnostic("notification read-mark prune deferred"); + candidates + } + }; + if !remaining.is_empty() { + self.notification_read_prunes.lock().unwrap().extend(remaining); + } + } + + /// Record that `client_id` read `notifications`. Unknown ids are reported, + /// not rejected: a bounded ledger may have evicted them, and an + /// acknowledgement of something already gone is complete by definition. + /// The refreshed rows are published as one resource revision so every + /// subscribed client converges on the same `read_by` sets. + pub(crate) fn ack_notifications( + &self, + mutation: &WorkspaceMutation, + expected_revision: Option, + client_id: &str, + notifications: &[NotificationPublicId], + ) -> anyhow::Result { + const OPERATION: &str = "notification.ack"; + validate_client_id(client_id)?; + let fingerprint = serde_json::json!({ + "operation": OPERATION, + "client_id": client_id, + "notifications": notifications, + }); + let mut registry = self.workspace_registry.lock().unwrap(); + if let Some(replay) = registry.replay_resource_patch(mutation, OPERATION, &fingerprint)? { + return Ok(replay); + } + let session_id = registry.session_id().clone(); + let mut acknowledged: Vec = Vec::new(); + let mut unknown: Vec = Vec::new(); + let mut deltas = Vec::new(); + { + let ledger = self.notification_ledger.lock().unwrap(); + let reads = self.notification_reads.lock().unwrap(); + for id in notifications { + if acknowledged.contains(id) || unknown.contains(id) { + continue; + } + match ledger.iter().find(|entry| &entry.id == id) { + Some(entry) => { + let mut read_by = reads.get(id).cloned().unwrap_or_default(); + read_by.insert(client_id.to_string()); + let read_by = read_by.into_iter().collect::>(); + let value = self.notification_snapshot_value(entry, &session_id, &read_by); + deltas.push(serde_json::json!({ + "kind":"upsert", + "sequence":0, + "resource":"notification", + "id":id, + "value":value, + })); + acknowledged.push(id.clone()); + } + None => unknown.push(id.clone()), + } + } + } + let result = serde_json::json!({ + "client_id": client_id, + "acknowledged": acknowledged, + "unknown": unknown, + }); + let commit = registry.commit_notification_ack( + mutation, + &fingerprint, + expected_revision, + client_id, + &acknowledged, + now_ms(), + &result, + &Value::Array(deltas), + )?; + if !commit.replayed { + { + let mut reads = self.notification_reads.lock().unwrap(); + for id in &acknowledged { + reads.entry(id.clone()).or_default().insert(client_id.to_string()); + } + } + self.state.lock().unwrap().resource_revision = commit.revision; + } + drop(registry); + if !commit.replayed { + self.publish_resource_event(); + } + Ok(commit) + } + pub fn report_agent( &self, surface: SurfaceId, @@ -22627,6 +22981,331 @@ mod tests { ); } + #[test] + fn agent_hook_transitions_post_durable_notifications_once() { + let mux = test_mux(); + let surface = mux.new_workspace(None, None).unwrap(); + let surface_id = surface.id; + let terminal_id = surface.terminal_public_id().cloned().unwrap(); + let ingress = |event: &str, native: Value| { + crate::agent_hooks::agent_hook_journal_ingress( + "claude", + event, + Some(&terminal_id.to_string()), + native, + ) + .unwrap() + }; + mux.apply_agent_hook_record(&ingress("SessionStart", serde_json::json!({})), 1).unwrap(); + mux.apply_agent_hook_record(&ingress("UserPromptSubmit", serde_json::json!({})), 2) + .unwrap(); + assert!(mux.resource_notifications(16).is_empty(), "start and prompt need no attention"); + + mux.apply_agent_hook_record( + &ingress( + "PermissionRequest", + serde_json::json!({"tool_name":"Bash","message":"Allow Bash(rm -rf build)?"}), + ), + 3, + ) + .unwrap(); + let posted = mux.resource_notifications(16); + assert_eq!(posted.len(), 1); + assert_eq!(posted[0].title, "Claude needs approval"); + assert_eq!(posted[0].body, "Bash", "redacted prompt text must not leak; tool name may"); + assert_eq!(posted[0].level, NotificationLevel::Warning); + assert_eq!(posted[0].terminal_id.as_ref(), Some(&terminal_id)); + assert!(mux.surface_notification(surface_id).is_some_and(|marker| marker.unread)); + + // A replayed sequence (restart repair) is a fence no-op and must not + // post a second notification. + mux.apply_agent_hook_record( + &ingress("PermissionRequest", serde_json::json!({"message":"again"})), + 3, + ) + .unwrap(); + assert_eq!(mux.resource_notifications(16).len(), 1); + + mux.apply_agent_hook_record(&ingress("Stop", serde_json::json!({})), 4).unwrap(); + let posted = mux.resource_notifications(16); + assert_eq!(posted.len(), 2); + assert_eq!(posted[0].title, "Claude finished"); + assert_eq!(posted[0].level, NotificationLevel::Info); + + // The rows are durable resource effects: the public snapshot carries + // them with an empty per-client read set. + let snapshot = crate::resource_api::public_session_snapshot(&mux).unwrap(); + let rows = snapshot["notifications"].as_array().unwrap(); + assert_eq!(rows.len(), 2); + for row in rows { + assert_eq!(row["read_by"], serde_json::json!([])); + assert_eq!(row["terminal_id"], serde_json::json!(terminal_id)); + } + } + + #[test] + fn notification_ack_is_per_client_and_replay_safe() { + let mux = test_mux(); + let surface = mux.new_workspace(None, None).unwrap(); + let surface_id = surface.id; + let first = mux + .post_notification("one".into(), "".into(), NotificationLevel::Info, Some(surface_id)) + .unwrap(); + let second = mux + .post_notification("two".into(), "".into(), NotificationLevel::Error, Some(surface_id)) + .unwrap(); + assert_ne!(first, second); + let ledger = mux.resource_notifications(16); + let (newest, oldest) = (ledger[0].id.clone(), ledger[1].id.clone()); + let revision_before = mux.with_state(|state| state.resource_revision); + let epoch_before = mux.resource_event_epoch(); + + let mutation = WorkspaceMutation::new("ack-a-1", "test").unwrap(); + let ack = + mux.ack_notifications(&mutation, None, "mac-a", std::slice::from_ref(&oldest)).unwrap(); + assert!(!ack.replayed); + assert_eq!(ack.result["acknowledged"], serde_json::json!([oldest])); + assert_eq!(ack.result["unknown"], serde_json::json!([])); + assert_eq!(ack.revision, revision_before + 1); + assert_eq!(mux.with_state(|state| state.resource_revision), revision_before + 1); + assert!(mux.resource_event_epoch() > epoch_before, "subscribers must wake"); + assert_eq!(mux.notification_read_by(&oldest), vec!["mac-a".to_string()]); + assert!(mux.notification_read_by(&newest).is_empty()); + // The shared console marker is not a per-client read. + assert!(mux.surface_notification(surface_id).is_some_and(|marker| marker.unread)); + + // Same key, same input: replay without a new revision or a new row. + let replay = + mux.ack_notifications(&mutation, None, "mac-a", std::slice::from_ref(&oldest)).unwrap(); + assert!(replay.replayed); + assert_eq!(replay.revision, revision_before + 1); + assert_eq!(mux.with_state(|state| state.resource_revision), revision_before + 1); + + // A second client keeps its own state and sees the merged set. + let mutation_b = WorkspaceMutation::new("ack-b-1", "test").unwrap(); + let unknown = NotificationPublicId::random().unwrap(); + let ack_b = mux + .ack_notifications( + &mutation_b, + None, + "mac-b", + &[newest.clone(), oldest.clone(), unknown.clone(), oldest.clone()], + ) + .unwrap(); + assert_eq!(ack_b.result["acknowledged"], serde_json::json!([newest, oldest])); + assert_eq!(ack_b.result["unknown"], serde_json::json!([unknown])); + assert_eq!( + mux.notification_read_by(&oldest), + vec!["mac-a".to_string(), "mac-b".to_string()] + ); + assert_eq!(mux.notification_read_by(&newest), vec!["mac-b".to_string()]); + + let snapshot = crate::resource_api::public_session_snapshot(&mux).unwrap(); + let rows = snapshot["notifications"].as_array().unwrap(); + let row = |id: &NotificationPublicId| { + rows.iter().find(|row| row["id"] == serde_json::json!(id)).cloned().unwrap() + }; + assert_eq!(row(&oldest)["read_by"], serde_json::json!(["mac-a", "mac-b"])); + assert_eq!(row(&newest)["read_by"], serde_json::json!(["mac-b"])); + + // The ack publishes one upsert delta per acknowledged row so remote + // feeds converge without a snapshot. + let page = mux.resource_events_after(revision_before).unwrap(); + let ack_batch = + page.batches.iter().find(|batch| batch.revision == revision_before + 2).unwrap(); + let changes = ack_batch.changes.as_array().unwrap(); + assert_eq!(changes.len(), 2); + assert!(changes.iter().all(|change| { + change["resource"] == "notification" + && change["kind"] == "upsert" + && change["value"]["read_by"] + .as_array() + .unwrap() + .contains(&serde_json::json!("mac-b")) + })); + + let bad = WorkspaceMutation::new("ack-bad", "test").unwrap(); + assert!(mux.ack_notifications(&bad, None, "has space", &[oldest]).is_err()); + } + + #[test] + fn notification_reads_survive_restart_and_eviction_prunes_them() { + let root = std::env::temp_dir() + .join(format!("cmux-notification-reads-{}", WorkspacePublicId::random().unwrap())); + let session = "notification-reads"; + let open = || { + let registry = WorkspaceRegistry::open(&root, session).unwrap(); + Mux::from_workspace_registry( + session.into(), + SurfaceOptions::default(), + registry, + ProviderWorkspaceState::default(), + true, + ) + .unwrap() + }; + let mux = open(); + let surface = mux.new_workspace(None, None).unwrap(); + let surface_id = surface.id; + mux.post_notification("kept".into(), "".into(), NotificationLevel::Info, Some(surface_id)) + .unwrap(); + let kept = mux.resource_notifications(1)[0].id.clone(); + let mutation = WorkspaceMutation::new("ack-restart", "test").unwrap(); + mux.ack_notifications(&mutation, None, "mac-a", std::slice::from_ref(&kept)).unwrap(); + drop(mux); + + let mux = open(); + assert_eq!(mux.notification_read_by(&kept), vec!["mac-a".to_string()]); + let snapshot = crate::resource_api::public_session_snapshot(&mux).unwrap(); + let row = snapshot["notifications"] + .as_array() + .unwrap() + .iter() + .find(|row| row["id"] == serde_json::json!(kept)) + .cloned() + .unwrap(); + assert_eq!(row["read_by"], serde_json::json!(["mac-a"])); + // Replay across restart still returns the original commit. + let replay = + mux.ack_notifications(&mutation, None, "mac-a", std::slice::from_ref(&kept)).unwrap(); + assert!(replay.replayed); + + // Evict `kept` from the bounded ledger, then acknowledge something + // else: the stale read row must be gone from memory and from disk. + let surface_id = mux.resource_surface_for_terminal( + mux.resource_notifications(1)[0].terminal_id.as_ref().unwrap(), + ); + for index in 0..256 { + mux.post_notification( + format!("fill-{index}"), + "".into(), + NotificationLevel::Info, + surface_id, + ) + .unwrap(); + } + assert!(mux.resource_notifications(256).iter().all(|entry| entry.id != kept)); + assert!(mux.notification_read_by(&kept).is_empty()); + // The prune rides the committed create that evicted `kept`, not an + // acknowledgement, and only once the receipts no longer retain it. + let newest = mux.resource_notifications(1)[0].id.clone(); + let prune = WorkspaceMutation::new("ack-newest", "test").unwrap(); + mux.ack_notifications(&prune, None, "mac-a", std::slice::from_ref(&newest)).unwrap(); + let stale_rows = mux + .workspace_registry + .lock() + .unwrap() + .durable_notification_read_clients(kept.as_str()) + .unwrap(); + assert!(stale_rows.is_empty(), "evicted notification kept read rows: {stale_rows:?}"); + drop(mux); + let mux = open(); + assert!(mux.notification_read_by(&kept).is_empty()); + assert_eq!(mux.notification_read_by(&newest), vec!["mac-a".to_string()]); + let _ = std::fs::remove_dir_all(&root); + } + + /// Random creates, acks from several clients, replays, and restarts must + /// keep one invariant: a retained notification's `read_by` is exactly the + /// set of clients that acknowledged it while it was retained. + #[test] + fn notification_read_state_converges_under_random_operations() { + let root = std::env::temp_dir() + .join(format!("cmux-notification-fuzz-{}", WorkspacePublicId::random().unwrap())); + let session = "notification-fuzz"; + let open = || { + let registry = WorkspaceRegistry::open(&root, session).unwrap(); + Mux::from_workspace_registry( + session.into(), + SurfaceOptions::default(), + registry, + ProviderWorkspaceState::default(), + true, + ) + .unwrap() + }; + let mut mux = open(); + let surface = mux.new_workspace(None, None).unwrap(); + let terminal_id = surface.terminal_public_id().cloned().unwrap(); + let clients = ["mac-a", "mac-b", "phone-c"]; + let mut expected: HashMap> = HashMap::new(); + let mut retained: VecDeque = VecDeque::new(); + let mut seed: u64 = 0x9e37_79b9_7f4a_7c15; + let mut next = || { + seed ^= seed << 13; + seed ^= seed >> 7; + seed ^= seed << 17; + seed + }; + let mut ack_counter = 0u64; + let mut last_ack: Option<(WorkspaceMutation, String, Vec)> = None; + for step in 0..400u64 { + match next() % 10 { + 0..=3 => { + let surface_id = mux.resource_surface_for_terminal(&terminal_id); + mux.post_notification( + format!("n-{step}"), + "".into(), + NotificationLevel::Info, + surface_id, + ) + .unwrap(); + let id = mux.resource_notifications(1)[0].id.clone(); + retained.push_back(id.clone()); + expected.insert(id, BTreeSet::new()); + while retained.len() > 256 { + let evicted = retained.pop_front().unwrap(); + expected.remove(&evicted); + } + } + 4..=6 if !retained.is_empty() => { + let client = clients[(next() % clients.len() as u64) as usize]; + let count = 1 + (next() % 4) as usize; + let ids = (0..count) + .map(|_| retained[(next() % retained.len() as u64) as usize].clone()) + .collect::>(); + ack_counter += 1; + let mutation = + WorkspaceMutation::new(format!("fuzz-ack-{ack_counter}"), "test").unwrap(); + mux.ack_notifications(&mutation, None, client, &ids).unwrap(); + for id in &ids { + expected.get_mut(id).unwrap().insert(client.to_string()); + } + last_ack = Some((mutation, client.to_string(), ids)); + } + 7 => { + if let Some((mutation, client, ids)) = &last_ack { + let replay = mux.ack_notifications(mutation, None, client, ids).unwrap(); + assert!(replay.replayed, "step {step}: replay must not re-commit"); + } + } + 8 => { + drop(mux); + mux = open(); + } + _ => {} + } + for id in &retained { + let actual = mux.notification_read_by(id).into_iter().collect::>(); + assert_eq!(&actual, &expected[id], "step {step}: read set diverged for {id}"); + } + } + drop(mux); + let mux = open(); + let snapshot = crate::resource_api::public_session_snapshot(&mux).unwrap(); + for row in snapshot["notifications"].as_array().unwrap() { + let id = NotificationPublicId::parse(row["id"].as_str().unwrap()).unwrap(); + let read_by = row["read_by"] + .as_array() + .unwrap() + .iter() + .map(|value| value.as_str().unwrap().to_string()) + .collect::>(); + assert_eq!(read_by, expected[&id], "restart lost or invented a read mark for {id}"); + } + let _ = std::fs::remove_dir_all(&root); + } + #[test] fn replayed_agent_hook_events_do_not_rewrite_the_record() { let mux = test_mux(); diff --git a/cmux-tui/crates/cmux-tui-core/src/mux/public_projections.rs b/cmux-tui/crates/cmux-tui-core/src/mux/public_projections.rs index 4ac907c40a02..500022c5d55c 100644 --- a/cmux-tui/crates/cmux-tui-core/src/mux/public_projections.rs +++ b/cmux-tui/crates/cmux-tui-core/src/mux/public_projections.rs @@ -12,6 +12,7 @@ pub(super) struct RestoredPublicProjections { pub(super) agent_hook_fences: HashMap, pub(super) terminal_notifications: HashMap, pub(super) notification_ledger: VecDeque, + pub(super) notification_reads: HashMap>, } pub(super) fn restore_public_projections( @@ -22,6 +23,7 @@ pub(super) fn restore_public_projections( let default_colors = projections.terminal_defaults.unwrap_or_default(); let mut notification_ledger = VecDeque::with_capacity(projections.notifications.len()); let mut terminal_notifications = HashMap::new(); + let mut notification_reads = HashMap::new(); for (index, notification) in projections.notifications.into_iter().enumerate() { let numeric_id = u64::try_from(index).context("notification count exceeds uint64")?.saturating_add(1); @@ -45,6 +47,12 @@ pub(super) fn restore_public_projections( ); } } + if !notification.read_by.is_empty() { + notification_reads.insert( + notification.id.clone(), + notification.read_by.into_iter().collect::>(), + ); + } notification_ledger.push_back(ResourceNotification { id: notification.id, title: notification.title, @@ -128,6 +136,7 @@ pub(super) fn restore_public_projections( agent_hook_fences, terminal_notifications, notification_ledger, + notification_reads, }) } @@ -217,6 +226,7 @@ mod tests { terminal_id: Some(terminal.clone()), created_at_ms: 1, unread: true, + read_by: vec![], }], agents: vec![RegistryAgentProjection { id: AgentPublicId::parse("agent_00000000000000000000000000000001").unwrap(), @@ -254,6 +264,7 @@ mod tests { terminal_id: None, created_at_ms: 2, unread: true, + read_by: vec![], }], agents: Vec::new(), agent_hook_states: Vec::new(), @@ -278,6 +289,7 @@ mod tests { terminal_id: Some(terminal.clone()), created_at_ms: 3, unread: true, + read_by: vec![], }], agents: Vec::new(), agent_hook_states: Vec::new(), diff --git a/cmux-tui/crates/cmux-tui-core/src/resource.rs b/cmux-tui/crates/cmux-tui-core/src/resource.rs index e799e9762687..e6f9485d4b70 100644 --- a/cmux-tui/crates/cmux-tui-core/src/resource.rs +++ b/cmux-tui/crates/cmux-tui-core/src/resource.rs @@ -357,6 +357,8 @@ pub enum ResourceOperation { NotificationList, #[serde(rename = "notification.create")] NotificationCreate, + #[serde(rename = "notification.ack")] + NotificationAck, #[serde(rename = "agent.list")] AgentList, #[serde(rename = "agent.report")] @@ -632,6 +634,7 @@ impl ResourceOperation { Self::BrowserClose => "browser.close", Self::NotificationList => "notification.list", Self::NotificationCreate => "notification.create", + Self::NotificationAck => "notification.ack", Self::AgentList => "agent.list", Self::AgentReport => "agent.report", Self::SidebarViewGet => "sidebar_view.get", diff --git a/cmux-tui/crates/cmux-tui-core/src/resource_api.rs b/cmux-tui/crates/cmux-tui-core/src/resource_api.rs index ad72aad48af9..78064719efd5 100644 --- a/cmux-tui/crates/cmux-tui-core/src/resource_api.rs +++ b/cmux-tui/crates/cmux-tui-core/src/resource_api.rs @@ -714,6 +714,7 @@ pub(crate) fn public_session_snapshot_with_journal_head( .as_ref() .and_then(|terminal_id| mux.terminal_notification(terminal_id)) .is_some_and(|notification| notification.unread), + "read_by": notification.read_by, }); if let Some(terminal_id) = notification.terminal_id { snapshot["terminal_id"] = json!(terminal_id); diff --git a/cmux-tui/crates/cmux-tui-core/src/resource_router.rs b/cmux-tui/crates/cmux-tui-core/src/resource_router.rs index 917d4d4ebf84..0ac43d780e02 100644 --- a/cmux-tui/crates/cmux-tui-core/src/resource_router.rs +++ b/cmux-tui/crates/cmux-tui-core/src/resource_router.rs @@ -881,6 +881,7 @@ fn dispatch_resource_request( )) } ResourceOperation::NotificationCreate => create_notification(mux, request), + ResourceOperation::NotificationAck => ack_notifications(mux, request), _ => unreachable!("operation_owner classifies snapshot operations exhaustively"), }, OperationOwner::Connection => Err(ResourceError::operation_failed( @@ -921,7 +922,8 @@ const fn operation_owner(operation: ResourceOperation) -> OperationOwner { | ResourceOperation::BrowserList | ResourceOperation::BrowserGet | ResourceOperation::NotificationList - | ResourceOperation::NotificationCreate => OperationOwner::Snapshot, + | ResourceOperation::NotificationCreate + | ResourceOperation::NotificationAck => OperationOwner::Snapshot, ResourceOperation::WorkspaceList | ResourceOperation::WorkspaceGet | ResourceOperation::WorkspaceCreate @@ -1355,20 +1357,19 @@ fn execute_notification_effect( terminal_id.clone(), created_at_ms, ); - let mut value = json!({ - "id":notification_id, - "session_id":session_id, - "title":title, - "body":body, - "level":level.as_str(), - "created_at_ms":created_at_ms.to_string(), - "unread":surface - .and_then(|surface| mux.surface_notification(surface)) - .is_some_and(|notification| notification.unread), - }); - if let Some(terminal_id) = terminal_id { - value["terminal_id"] = json!(terminal_id); - } + let value = mux.notification_snapshot_value( + &crate::ResourceNotification { + id: notification_id.clone(), + title: title.to_string(), + body: body.to_string(), + level, + terminal_id, + created_at_ms, + surface, + }, + &session_id, + &[], + ); let outcome = ResourceEffectOutcome::Success(value.clone()); let deltas = json!([{ "kind":"upsert", @@ -1390,9 +1391,50 @@ fn execute_notification_effect( return Err(indeterminate_error(idempotency_key, "notification.create")); } }; + mux.prune_evicted_notification_reads(); mutation_result(mux, value, revision, false) } +fn ack_notifications(mux: &Mux, request: ParsedResourceRequest) -> Result { + ensure_session_route(mux, &request.selectors)?; + let client_id = required_string(&request.fields, "client_id")?.to_string(); + crate::mux::validate_client_id(&client_id) + .map_err(|error| validation_error(&error.to_string(), json!({"client_id":client_id})))?; + let notifications = request + .fields + .get("notifications") + .and_then(Value::as_array) + .ok_or_else(|| validation_error("notifications must be an array", json!({})))? + .iter() + .map(|value| { + NotificationPublicId::parse( + value + .as_str() + .ok_or_else(|| validation_error("notification id must be a string", json!({})))? + .to_string(), + ) + }) + .collect::, ResourceError>>()?; + let mutation = crate::workspace_registry::WorkspaceMutation::new( + request + .envelope + .idempotency_key + .clone() + .expect("catalog-validated mutations have an idempotency key"), + "resource-api", + ) + .map_err(resource_operation_error)?; + let ack = mux + .ack_notifications( + &mutation, + expected_revision(&request.fields)?, + &client_id, + ¬ifications, + ) + .map_err(resource_operation_error)?; + mutation_result(mux, ack.result, ack.revision, ack.replayed) +} + fn indeterminate_error(idempotency_key: &str, operation: &str) -> ResourceError { ResourceError::new( "mutation.indeterminate", @@ -1630,7 +1672,7 @@ mod tests { #[test] fn every_catalog_operation_has_one_concrete_owner() { let operations = operation_catalog()["operations"].as_object().unwrap(); - assert_eq!(operations.len(), 125); + assert_eq!(operations.len(), 126); for name in operations.keys() { let operation: ResourceOperation = serde_json::from_value(Value::String(name.clone())).unwrap(); @@ -1649,7 +1691,7 @@ mod tests { #[test] fn every_catalog_operation_accepts_its_result_and_declared_error_fixtures() { let operations = operation_catalog()["operations"].as_object().unwrap(); - assert_eq!(operations.len(), 125); + assert_eq!(operations.len(), 126); for (name, descriptor) in operations { let operation: ResourceOperation = serde_json::from_value(Value::String(name.clone())).unwrap(); diff --git a/cmux-tui/crates/cmux-tui-core/src/workspace_registry/public_projection_store.rs b/cmux-tui/crates/cmux-tui-core/src/workspace_registry/public_projection_store.rs index 2d7d1be89cbd..67f67e50935c 100644 --- a/cmux-tui/crates/cmux-tui-core/src/workspace_registry/public_projection_store.rs +++ b/cmux-tui/crates/cmux-tui-core/src/workspace_registry/public_projection_store.rs @@ -32,6 +32,8 @@ pub struct RegistryNotificationProjection { pub terminal_id: Option, pub created_at_ms: u64, pub unread: bool, + /// Client ids that acknowledged this notification, sorted and unique. + pub read_by: Vec, } #[derive(Debug, Clone, PartialEq, Eq)] @@ -87,6 +89,11 @@ struct StoredNotification { terminal_id: Option, created_at_ms: WireDecimal, unread: bool, + /// Read marks at commit time are always empty; the durable truth is the + /// `resource_notification_reads` table, so this field is decoded and + /// ignored. + #[serde(default)] + read_by: Vec, #[serde(default)] extra: Option>, } @@ -258,6 +265,29 @@ impl WorkspaceRegistry { Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)) })? .collect::, _>>()?; + let mut reads = self.durable_notification_reads()?; + { + // Marks for notifications outside the retained window are dead + // weight after a restart (the in-memory prune queue did not + // survive). Drop them here so the table stays bounded. + let retained_ids = rows + .iter() + .filter_map(|(outcome_json, _)| { + serde_json::from_str::(outcome_json) + .ok() + .and_then(|value| value["value"]["id"].as_str().map(str::to_string)) + }) + .collect::>(); + let stale = + reads.keys().filter(|id| !retained_ids.contains(*id)).cloned().collect::>(); + for id in &stale { + self.connection.execute( + "DELETE FROM resource_notification_reads WHERE notification_id = ?1", + [id.as_str()], + )?; + reads.remove(id); + } + } let mut notifications = Vec::with_capacity(rows.len()); for (outcome_json, idempotency_key) in rows { let outcome: ResourceEffectOutcome = serde_json::from_str(&outcome_json) @@ -284,6 +314,8 @@ impl WorkspaceRegistry { self.session_id ); let _ = stored.extra; + let _ = stored.read_by; + let read_by = reads.remove(stored.id.as_str()).unwrap_or_default(); notifications.push(RegistryNotificationProjection { id: stored.id, title: stored.title, @@ -294,12 +326,41 @@ impl WorkspaceRegistry { .filter(|terminal_id| live_terminals.contains(terminal_id)), created_at_ms: stored.created_at_ms.get(), unread: stored.unread, + read_by, }); } notifications.reverse(); Ok(notifications) } + /// Read marks stored for one notification, for tests that verify pruning. + #[cfg(test)] + pub(crate) fn durable_notification_read_clients( + &self, + notification_id: &str, + ) -> anyhow::Result> { + Ok(self.durable_notification_reads()?.remove(notification_id).unwrap_or_default()) + } + + /// Per-client read marks keyed by notification id, each list sorted and + /// unique. Rows for notifications the ledger evicted are pruned at the + /// next acknowledgement, so this stays bounded. + fn durable_notification_reads(&self) -> anyhow::Result>> { + let mut statement = self.connection.prepare( + "SELECT notification_id, client_id + FROM resource_notification_reads + ORDER BY notification_id ASC, client_id ASC", + )?; + let rows = statement + .query_map([], |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)))? + .collect::, _>>()?; + let mut reads: HashMap> = HashMap::new(); + for (notification_id, client_id) in rows { + reads.entry(notification_id).or_default().push(client_id); + } + Ok(reads) + } + fn durable_agents( &self, terminal: Option<&TerminalPublicId>, diff --git a/cmux-tui/crates/cmux-tui-core/src/workspace_registry/resource_store.rs b/cmux-tui/crates/cmux-tui-core/src/workspace_registry/resource_store.rs index f6469c5e3aba..c9f6f2bbf744 100644 --- a/cmux-tui/crates/cmux-tui-core/src/workspace_registry/resource_store.rs +++ b/cmux-tui/crates/cmux-tui-core/src/workspace_registry/resource_store.rs @@ -1,6 +1,6 @@ use super::*; use crate::JournalIngress; -use crate::resource::TerminalPublicId; +use crate::resource::{NotificationPublicId, TerminalPublicId}; use serde_json::json; /// Completed pure mutations keep a finite exactly-once replay window. Pruning @@ -203,6 +203,13 @@ pub(super) fn create_resource_schema(transaction: &Transaction<'_>) -> anyhow::R attempt INTEGER NOT NULL CHECK(attempt >= 0), PRIMARY KEY(producer_id, origin, idempotency_key) ); + CREATE TABLE IF NOT EXISTS resource_notification_reads ( + notification_id TEXT NOT NULL, + client_id TEXT NOT NULL, + read_at_ms INTEGER NOT NULL CHECK(read_at_ms >= 0), + committed_revision INTEGER NOT NULL CHECK(committed_revision >= 0), + PRIMARY KEY(notification_id, client_id) + ); DROP TRIGGER IF EXISTS resource_agent_projection_terminal_tombstone; CREATE INDEX IF NOT EXISTS resource_mutations_by_operation_revision ON resource_mutations(operation, committed_revision DESC); @@ -991,6 +998,127 @@ impl WorkspaceRegistry { Ok(ResourcePatchCommit { revision, result: result.clone(), replayed: false }) } + /// Record that `client_id` has read `acknowledged` notifications and + /// publish the refreshed notification rows as one resource revision. + /// Only the requested marks are written; eviction pruning is a separate + /// exact step (`prune_notification_reads`) driven by committed creates. + #[allow(clippy::too_many_arguments)] + pub(crate) fn commit_notification_ack( + &mut self, + mutation: &WorkspaceMutation, + fingerprint: &Value, + expected_revision: Option, + client_id: &str, + acknowledged: &[NotificationPublicId], + read_at_ms: u64, + result: &Value, + deltas: &Value, + ) -> anyhow::Result { + const OPERATION: &str = "notification.ack"; + validate_identifier("mutation id", &mutation.id)?; + validate_identifier("mutation origin", &mutation.origin)?; + let fingerprint = canonical_json(fingerprint)?; + let result_json = canonical_json(result)?; + let tx = self.connection.transaction()?; + if let Some(replayed) = resource_patch_replay(&tx, mutation, OPERATION, &fingerprint)? { + return Ok(replayed); + } + let previous_revision = transaction_resource_revision(&tx)?; + if let Some(expected) = expected_revision + && expected != previous_revision + { + anyhow::bail!( + "resource revision conflict: expected {expected}, current {previous_revision}" + ); + } + let revision = previous_revision + .checked_add(1) + .ok_or_else(|| anyhow::anyhow!("resource revision exhausted"))?; + let sqlite_revision = + i64::try_from(revision).context("resource revision exceeds SQLite range")?; + let sqlite_read_at = + i64::try_from(read_at_ms).context("notification read time exceeds SQLite range")?; + for notification_id in acknowledged { + tx.execute( + "INSERT INTO resource_notification_reads( + notification_id, client_id, read_at_ms, committed_revision + ) VALUES(?1, ?2, ?3, ?4) + ON CONFLICT(notification_id, client_id) DO NOTHING", + params![notification_id.as_str(), client_id, sqlite_read_at, sqlite_revision], + )?; + } + tx.execute( + "UPDATE meta SET value = ?1 WHERE key = 'resource_revision'", + [revision.to_string()], + )?; + tx.execute( + "INSERT INTO resource_mutations( + origin, idempotency_key, operation, fingerprint, result_json, committed_revision + ) VALUES(?1, ?2, ?3, ?4, ?5, ?6)", + params![ + mutation.origin, + mutation.id, + OPERATION, + fingerprint, + result_json, + sqlite_revision, + ], + )?; + append_resource_journal_record( + &tx, + revision, + previous_revision, + &mutation.origin, + &mutation.id, + OPERATION, + None, + result, + deltas, + )?; + prune_resource_mutations(&tx)?; + tx.commit()?; + Ok(ResourcePatchCommit { revision, result: result.clone(), replayed: false }) + } + + /// Delete read marks for `candidates` that the committed notification + /// receipts no longer retain. Returns the candidates still retained, so + /// the caller keeps them queued. Retention here is the same query the + /// projection rebuild uses, so memory and disk agree after a restart. + pub(crate) fn prune_notification_reads( + &mut self, + candidates: &[NotificationPublicId], + ) -> anyhow::Result> { + let tx = self.connection.transaction()?; + let mut remaining = Vec::new(); + for candidate in candidates { + let retained: bool = tx.query_row( + "SELECT EXISTS( + SELECT 1 FROM ( + SELECT json_extract(outcome_json, '$.value.id') AS id + FROM resource_effect_receipts + WHERE operation = 'notification.create' + AND state = 'committed' + AND json_extract(outcome_json, '$.kind') = 'success' + ORDER BY committed_revision DESC, idempotency_key DESC + LIMIT 256 + ) WHERE id = ?1 + )", + [candidate.as_str()], + |row| row.get(0), + )?; + if retained { + remaining.push(candidate.clone()); + } else { + tx.execute( + "DELETE FROM resource_notification_reads WHERE notification_id = ?1", + [candidate.as_str()], + )?; + } + } + tx.commit()?; + Ok(remaining) + } + pub fn terminal_resource_id( &self, terminal_id: &str, diff --git a/cmux-tui/crates/cmux-tui/src/cli/command.rs b/cmux-tui/crates/cmux-tui/src/cli/command.rs index e9ddb670f340..ca9f7ae3d7dc 100644 --- a/cmux-tui/crates/cmux-tui/src/cli/command.rs +++ b/cmux-tui/crates/cmux-tui/src/cli/command.rs @@ -1332,6 +1332,35 @@ fn parse_notification(words: &[String], flags: &mut Flags) -> Result { + let mut params = Map::new(); + let client_id = flags.required("client")?; + if client_id.is_empty() + || client_id.len() > 128 + || !client_id.bytes().all(|byte| byte.is_ascii_graphic()) + { + return Err(UsageError::new( + "--client must be 1 to 128 printable ASCII bytes without spaces", + )); + } + params.insert("client_id".into(), Value::String(client_id)); + if ids.is_empty() { + return Err(UsageError::new("notification ack needs at least one notification ID")); + } + if ids.len() > 256 { + return Err(UsageError::new( + "notification ack accepts at most 256 notification IDs", + )); + } + for id in ids { + validate_prefixed_id("notification", "notification", id)?; + } + params.insert( + "notifications".into(), + Value::Array(ids.iter().map(|id| Value::String((*id).to_string())).collect()), + ); + request(ResourceOperation::NotificationAck, &selectors, flags, params) + } _ => usage("notification action"), } } @@ -4520,11 +4549,21 @@ mod tests { "sidebar_view.resize", ), (vec!["sidebar", "view", "reload", "--view", VIEW], "sidebar_view.reload"), + ( + vec![ + "notification", + "ack", + "notification_00000000000000000000000000000041", + "--client", + "mac-1", + ], + "notification.ack", + ), ]; - assert_eq!(cases.len(), 118); + assert_eq!(cases.len(), 119); let catalog = operation_catalog(); - assert_eq!(catalog["operations"].as_object().unwrap().len(), 125); + assert_eq!(catalog["operations"].as_object().unwrap().len(), 126); let mut seen = std::collections::BTreeSet::new(); let mut covered_fields = BTreeMap::<&str, std::collections::BTreeSet>::new(); for (args, expected) in &cases { diff --git a/cmux-tui/scripts/test_check_resource_api_boundary.py b/cmux-tui/scripts/test_check_resource_api_boundary.py index 9f920ec1de1b..db1c04fd70ef 100644 --- a/cmux-tui/scripts/test_check_resource_api_boundary.py +++ b/cmux-tui/scripts/test_check_resource_api_boundary.py @@ -680,7 +680,7 @@ def test_live_catalog_counts_and_local_endpoint_scope_are_frozen(self) -> None: catalog = json.loads( (SCRIPT.parents[1] / "spec/resource-operations-v2.json").read_text(encoding="utf-8") ) - self.assertEqual(len(catalog["operations"]), 125) + self.assertEqual(len(catalog["operations"]), 126) self.assertEqual(len(catalog["local_operations"]), 6) self.assertEqual( set(catalog["types"]["MachineSnapshot"]["fields"]), diff --git a/cmux-tui/spec/cli.md b/cmux-tui/spec/cli.md index a764c69b8de6..95eeac703251 100644 --- a/cmux-tui/spec/cli.md +++ b/cmux-tui/spec/cli.md @@ -314,6 +314,7 @@ browser key|text|attach|close browser mouse|wheel --pointer-frame-seq notification list|create +notification ack --client ... agent list|report pairing request list pairing request respond @@ -326,6 +327,16 @@ provider authority install ``` +`notification ack --client ...` records that one client +install has read the listed notifications. `--client` is the durable client +id (1 to 128 printable ASCII bytes) that the client also reports through +`client-focus`. Read state is per client: every notification row carries +`read_by`, the sorted client ids that acknowledged it, and a second client +keeps its own unread state. The shared `unread` marker on the console tree is +unchanged by an acknowledgement. Ids the bounded ledger no longer retains are +returned under `unknown`, not rejected, so a late acknowledgement after +eviction is complete. + `terminal output read` returns a bounded plain-text window of the terminal's journaled output stream: `{text, start_offset, next_offset, complete}`. Offsets are `terminal.output` stream byte offsets; pass a previous diff --git a/cmux-tui/spec/inventory.json b/cmux-tui/spec/inventory.json index 803cead59366..97eb244e0082 100644 --- a/cmux-tui/spec/inventory.json +++ b/cmux-tui/spec/inventory.json @@ -31,6 +31,7 @@ "frontend_projection.put", "machine.get", "machine.list", + "notification.ack", "notification.create", "notification.list", "pairing_request.list", diff --git a/cmux-tui/spec/resource-api-v2.json b/cmux-tui/spec/resource-api-v2.json index cfd3c40618dd..a3957239d218 100644 --- a/cmux-tui/spec/resource-api-v2.json +++ b/cmux-tui/spec/resource-api-v2.json @@ -66,6 +66,7 @@ "frontend_projection.put", "machine.get", "machine.list", + "notification.ack", "notification.create", "notification.list", "pairing_request.list", @@ -221,6 +222,7 @@ "browser.navigate", "browser.reload", "frontend_projection.put", + "notification.ack", "notification.create", "pairing_request.resolve", "pane.close", @@ -330,6 +332,7 @@ "browser.navigate", "browser.reload", "frontend_projection.put", + "notification.ack", "notification.create", "pairing_request.resolve", "pane.close", diff --git a/cmux-tui/spec/resource-api-v2.md b/cmux-tui/spec/resource-api-v2.md index e007bca4a185..baf89a1bbe30 100644 --- a/cmux-tui/spec/resource-api-v2.md +++ b/cmux-tui/spec/resource-api-v2.md @@ -399,7 +399,7 @@ defines the catalog format. Unknown parameter and result fields are rejected. | Class | Operations | | --- | --- | | read | `agent.list`, `browser.get`, `browser.list`, `client.get`, `client.list`, `frontend_projection.get`, `machine.get`, `machine.list`, `notification.list`, `pairing_request.list`, `pane.get`, `pane.list`, `pane.neighbor.get`, `screen.get`, `screen.layout.export`, `screen.list`, `session.creation.resolve`, `session.get`, `session.journal.checkpoint.list`, `session.journal.hook.list`, `session.journal.producer.list`, `session.journal.restore.preview`, `session.journal.segment.list`, `session.list`, `session.ping`, `session.snapshot`, `sidebar_view.get`, `tab.get`, `tab.list`, `terminal.copy`, `terminal.get`, `terminal.history.read`, `terminal.list`, `terminal.output_read`, `terminal.process.get`, `terminal.screen.read`, `terminal.state.read`, `terminal.wait`, `terminal.wait_exit`, `workspace.get`, `workspace.list` | -| mutation | `agent.report`, `browser.activate`, `browser.back`, `browser.close`, `browser.forward`, `browser.input.key`, `browser.input.mouse`, `browser.input.text`, `browser.input.wheel`, `browser.navigate`, `browser.reload`, `frontend_projection.put`, `notification.create`, `pairing_request.resolve`, `pane.close`, `pane.create`, `pane.focus`, `pane.focus_direction`, `pane.rename`, `pane.run`, `pane.split`, `pane.split_ratio.set`, `pane.swap`, `pane.viewport_width.set`, `pane.zoom`, `screen.close`, `screen.create`, `screen.focus`, `screen.layout.undo`, `screen.rename`, `session.journal.append`, `session.journal.checkpoint.create`, `session.journal.hook.put`, `session.journal.producer.put`, `session.journal.segment.seal`, `session.open`, `session.reload_config`, `session.shutdown`, `session.terminal_defaults.update`, `session.window.title.clear`, `session.window.title.set`, `sidebar_view.ensure`, `sidebar_view.input`, `sidebar_view.reload`, `sidebar_view.resize`, `tab.close`, `tab.create_browser`, `tab.create_terminal`, `tab.focus`, `tab.move`, `tab.rename`, `terminal.close`, `terminal.history.clear`, `terminal.input.focus`, `terminal.input.keys`, `terminal.input.mouse`, `terminal.input.write`, `terminal.move`, `terminal.project`, `terminal.viewport.scroll`, `workspace.close`, `workspace.create`, `workspace.focus`, `workspace.layout.apply`, `workspace.move`, `workspace.rename`, `workspace.run` | +| mutation | `agent.report`, `browser.activate`, `browser.back`, `browser.close`, `browser.forward`, `browser.input.key`, `browser.input.mouse`, `browser.input.text`, `browser.input.wheel`, `browser.navigate`, `browser.reload`, `frontend_projection.put`, `notification.ack`, `notification.create`, `pairing_request.resolve`, `pane.close`, `pane.create`, `pane.focus`, `pane.focus_direction`, `pane.rename`, `pane.run`, `pane.split`, `pane.split_ratio.set`, `pane.swap`, `pane.viewport_width.set`, `pane.zoom`, `screen.close`, `screen.create`, `screen.focus`, `screen.layout.undo`, `screen.rename`, `session.journal.append`, `session.journal.checkpoint.create`, `session.journal.hook.put`, `session.journal.producer.put`, `session.journal.segment.seal`, `session.open`, `session.reload_config`, `session.shutdown`, `session.terminal_defaults.update`, `session.window.title.clear`, `session.window.title.set`, `sidebar_view.ensure`, `sidebar_view.input`, `sidebar_view.reload`, `sidebar_view.resize`, `tab.close`, `tab.create_browser`, `tab.create_terminal`, `tab.focus`, `tab.move`, `tab.rename`, `terminal.close`, `terminal.history.clear`, `terminal.input.focus`, `terminal.input.keys`, `terminal.input.mouse`, `terminal.input.write`, `terminal.move`, `terminal.project`, `terminal.viewport.scroll`, `workspace.close`, `workspace.create`, `workspace.focus`, `workspace.layout.apply`, `workspace.move`, `workspace.rename`, `workspace.run` | | stream_open | `browser.attach`, `session.events`, `session.journal.subscribe`, `sidebar_view.attach`, `terminal.attach` | | connection_control | `browser.viewer.release`, `browser.viewer.resize`, `client.cell_pixels.set`, `client.detach`, `client.metadata.update`, `client.sizing.release`, `client.sizing.set`, `request.cancel`, `stream.cancel`, `terminal.renderer_grant.create`, `terminal.viewer.release`, `terminal.viewer.resize` | | local | `sidebar_plugin.install`, `sidebar_plugin.list`, `sidebar_plugin.remove`, `sidebar_plugin.update`, `sidebar_plugin.use`, `sidebar_plugin.use_builtin` | diff --git a/cmux-tui/spec/resource-operations-v2.json b/cmux-tui/spec/resource-operations-v2.json index fdb1722264a1..0603163744ef 100644 --- a/cmux-tui/spec/resource-operations-v2.json +++ b/cmux-tui/spec/resource-operations-v2.json @@ -2453,6 +2453,21 @@ "name": "boolean" } }, + "read_by": { + "required": true, + "type": { + "kind": "array", + "min_items": 0, + "max_items": 256, + "items": { + "kind": "primitive", + "name": "string", + "min_length": 1, + "max_length": 128 + } + }, + "description": "Client ids that acknowledged this notification through notification.ack. Sorted, unique. Per-client read state; the shared unread marker is unaffected." + }, "extra": { "required": false, "type": { @@ -2467,6 +2482,47 @@ }, "extra": false }, + "NotificationAckResult": { + "kind": "object", + "fields": { + "client_id": { + "required": true, + "type": { + "kind": "primitive", + "name": "string", + "min_length": 1, + "max_length": 128 + } + }, + "acknowledged": { + "required": true, + "type": { + "kind": "array", + "min_items": 0, + "max_items": 256, + "items": { + "kind": "resource_id", + "resource": "notification" + } + }, + "description": "Notification ids now recorded as read by client_id, including ids that were already read." + }, + "unknown": { + "required": true, + "type": { + "kind": "array", + "min_items": 0, + "max_items": 256, + "items": { + "kind": "resource_id", + "resource": "notification" + } + }, + "description": "Requested ids the session no longer retains. They are not an error: a bounded ledger may have evicted them." + } + }, + "extra": false + }, "AgentSnapshot": { "kind": "object", "fields": { @@ -7240,6 +7296,73 @@ "validation.invalid" ] }, + "notification.ack": { + "class": "mutation", + "idempotency": "required", + "target": "notification", + "ancestors": [ + "machine", + "session" + ], + "params": { + "selectors": { + "machine": "required", + "session": "required" + }, + "fields": { + "client_id": { + "required": true, + "type": { + "kind": "primitive", + "name": "string", + "min_length": 1, + "max_length": 128 + }, + "description": "Durable identity of the acknowledging client install. ASCII graphic bytes only." + }, + "notifications": { + "required": true, + "type": { + "kind": "array", + "min_items": 1, + "max_items": 256, + "items": { + "kind": "resource_id", + "resource": "notification" + } + } + }, + "expected_revision": { + "required": false, + "type": { + "kind": "primitive", + "name": "decimal" + }, + "description": "Optimistic concurrency cursor revision. Omit for no revision precondition." + } + }, + "extra": false + }, + "result": { + "kind": "apply", + "name": "MutationResult", + "arguments": [ + { + "kind": "ref", + "name": "NotificationAckResult" + } + ] + }, + "errors": [ + "idempotency.conflict", + "operation.failed", + "revision.conflict", + "selector.ambiguous", + "selector.invalid", + "selector.not_found", + "validation.invalid" + ] + }, "notification.create": { "class": "mutation", "idempotency": "required", diff --git a/cmux-tui/spec/resource-operations-v2.md b/cmux-tui/spec/resource-operations-v2.md index 9ef47704a3b8..97bcb1ec5e28 100644 --- a/cmux-tui/spec/resource-operations-v2.md +++ b/cmux-tui/spec/resource-operations-v2.md @@ -6,14 +6,14 @@ selectors, fields, results, errors, constraints, or stream types. ## Transported operations -`cmux.protocol/2` transports 125 operations for exactly one local mux +`cmux.protocol/2` transports 126 operations for exactly one local mux session. Cross-machine aggregation and provider lifecycle require a later broker protocol. | Class | Count | Semantics | | --- | ---: | --- | | `read` | 41 | Reads state and forbids an idempotency key | -| `mutation` | 67 | Requires an idempotency key and returns a mutation result | +| `mutation` | 68 | Requires an idempotency key and returns a mutation result | | `stream_open` | 5 | Opens a connection-owned typed stream | | `connection_control` | 12 | Changes only connection-local state | @@ -33,7 +33,7 @@ correlation, and idempotency metadata. | `client` | 7 | `client.cell_pixels.set`, `client.detach`, `client.get`, `client.list`, `client.metadata.update`, `client.sizing.release`, `client.sizing.set` | | `frontend_projection` | 2 | `frontend_projection.get`, `frontend_projection.put` | | `machine` | 2 | `machine.get`, `machine.list` | -| `notification` | 2 | `notification.create`, `notification.list` | +| `notification` | 3 | `notification.ack`, `notification.create`, `notification.list` | | `pairing_request` | 2 | `pairing_request.list`, `pairing_request.resolve` | | `pane` | 14 | `pane.close`, `pane.create`, `pane.focus`, `pane.focus_direction`, `pane.get`, `pane.list`, `pane.neighbor.get`, `pane.rename`, `pane.run`, `pane.split`, `pane.split_ratio.set`, `pane.swap`, `pane.viewport_width.set`, `pane.zoom` | | `request` | 1 | `request.cancel` | diff --git a/docs/cloud-cmux-tui-daemon.md b/docs/cloud-cmux-tui-daemon.md index 2951408ef078..504df20449b9 100644 --- a/docs/cloud-cmux-tui-daemon.md +++ b/docs/cloud-cmux-tui-daemon.md @@ -540,6 +540,52 @@ orthogonal: it routes model credentials, not compute, and is configured inside the machine the same way as locally. The `skills/cmux-cloud-vm` skill teaches this policy to Claude Code, Codex, OpenCode, and Pi. +## Notifications: the VM is the source of truth + +A Cloud machine's notifications live in its cmux-tui daemon, and every client +(the Mac app, a second Mac, an iPhone, the in-VM TUI) derives its unread state +from that one ledger. The Mac never runs a listener and the VM never dials the +Mac; the existing Mac-to-VM state feed carries notifications like any other +resource. + +Sources. An agent hook transition that deserves attention posts a durable +`notification.create` effect from inside the journal fold: `turn.completed` +(info, "Claude finished"), `approval.requested`, `question.requested`, +`plan_review.requested` (warning), and `error.reported` (error). The key is +derived from the journal sequence, so a crash between the notification commit +and the agent-report commit replays the notification on retry instead of +posting it twice, and the hook fence still advances once. Prompt and message +text is redacted before the journal accepts it, so the body carries only the +tool name an approval waits on. The legacy `notify` verb and `cmux-tui +notification create` use the same durable path. All three survive a daemon +restart and are rebuilt from committed effect receipts. + +Delivery. Every notification is a row in the `notifications` collection of +the public session snapshot and an `upsert` delta on the `session current +events` feed. The Mac's per-machine link already resumes that feed from its +`(generation, revision)` cursor with bounded recovery, so a notification posted +while the link was down arrives on reconnect through the same catch-up path +as a renamed tab. No journal subscription and no second stream are involved. + +Read state is per client. Each row carries `read_by`, the sorted client ids +that acknowledged it. `notification.ack {client_id, notifications[]}` records +the marks in `resource_notification_reads`, publishes the refreshed rows as one +revision, and replays under its idempotency key. Ids the 256-entry ledger no +longer retains come back under `unknown` rather than as an error, and their +read rows are pruned in the same transaction. The shared `unread` marker on +the console tree (the TUI's own tab dot) is unchanged by an acknowledgement: +it answers "does this terminal need attention on the shared console", while +`read_by` answers "has this client install seen it". Two Macs attached to one +machine therefore keep independent dots, and a Mac reattaching after a +reinstall with a new client id starts unread, which is the safe direction. + +Client ids are the durable per-install identity already used for focus +memory (`client-focus`), 1 to 128 printable ASCII bytes. The Mac app derives +one from its `vm-tui-devices.json` record for the machine. + +CLI: `cmux-tui notification list` prints rows with `read_by`; +`cmux-tui notification ack --client ...` acknowledges. + ## Surface catalog Terminals, VNC screens and browsers are *resources*; panes and workspaces are