diff --git a/crates/server/src/http.rs b/crates/server/src/http.rs index 16777d14..9f879a27 100644 --- a/crates/server/src/http.rs +++ b/crates/server/src/http.rs @@ -424,6 +424,10 @@ pub struct AppState { /// only ever diffs against its own last read or two, and a miss degrades /// gracefully to the full tree. pub element_snapshots: Arc>, + /// The last snapshot each named owner (`X-Phone-Owner`) was handed. An + /// observed action that names no baseline diffs against it, so a client + /// gets only what changed without having to carry the token itself. + pub owner_snapshots: Arc>>, /// What the driving agent did in the current app, for `registry` hints, /// `flow_suggestion` and `GET /agent/flow/draft` (see `crate::flows`). pub flow_trail: Arc>, @@ -4816,6 +4820,35 @@ fn remember_element_snapshot( } } +/// Owners whose last snapshot is remembered; past this the map is cleared +/// (owners are a handful of agent sessions, not an unbounded set). +const OWNER_SNAPSHOTS_CAP: usize = 64; + +fn named_owner(headers: &HeaderMap) -> Option<&str> { + match owner_claim_from_headers(headers) { + Ok(OwnerClaim::Named(name) | OwnerClaim::Takeover(name)) => Some(name), + _ => None, + } +} + +/// Remember `snapshot` as the last screen this request's owner saw. +fn note_owner_snapshot(state: &AppState, headers: &HeaderMap, snapshot: &str) { + let Some(owner) = named_owner(headers) else { + return; + }; + let mut owners = recover(state.owner_snapshots.lock()); + if owners.len() >= OWNER_SNAPSHOTS_CAP && !owners.contains_key(owner) { + owners.clear(); + } + owners.insert(owner.to_string(), snapshot.to_string()); +} + +/// The last snapshot this request's owner saw, if any. +fn owner_last_snapshot(state: &AppState, headers: &HeaderMap) -> Option { + let owner = named_owner(headers)?; + recover(state.owner_snapshots.lock()).get(owner).cloned() +} + fn lookup_element_snapshot( state: &AppState, snapshot: &str, @@ -8937,6 +8970,7 @@ async fn agent_actions( if let Some((snapshot, rows)) = observed { let rows = Arc::new(rows); remember_element_snapshot(&state, &snapshot, &rows); + note_owner_snapshot(&state, &headers, &snapshot); result["snapshot"] = serde_json::json!(snapshot); result["elements"] = serde_json::json!(&*rows); } @@ -9714,6 +9748,11 @@ async fn agent_input( }), Some((snapshot, rows)) => { remember_element_snapshot(&state, &snapshot, &rows); + // No baseline named: the last screen this owner + // saw, so a client gets the change without + // carrying the token. Read before noting the new one. + let owner_baseline = owner_last_snapshot(&state, &headers); + note_owner_snapshot(&state, &headers, &snapshot); let baseline = query .since .as_deref() @@ -9721,6 +9760,7 @@ async fn agent_input( .or_else(|| { value.get("snapshot").and_then(serde_json::Value::as_str) }) + .or(owner_baseline.as_deref()) .and_then(|since| { lookup_element_snapshot(&state, since) .map(|baseline| (since, baseline)) @@ -10467,6 +10507,7 @@ async fn agent_elements( }; let rows = Arc::new(rows); remember_element_snapshot(&state, &snapshot, &rows); + note_owner_snapshot(&state, &headers, &snapshot); // Additive, read-only usability signals over the same rows (visual-fallback // design §1.3): the client decides AX-vs-vision policy; the daemon only // reports. Computed before `screen` is consumed by JSON conversion. diff --git a/crates/server/src/main.rs b/crates/server/src/main.rs index 841de33c..f1af9049 100644 --- a/crates/server/src/main.rs +++ b/crates/server/src/main.rs @@ -597,6 +597,7 @@ fn serve() -> Result<()> { live_streams: Arc::new(std::sync::atomic::AtomicUsize::new(0)), mjpeg_stream_activity: Arc::new(Mutex::new(std::collections::HashMap::new())), element_snapshots: Arc::new(Mutex::new(std::collections::VecDeque::new())), + owner_snapshots: Arc::new(Mutex::new(std::collections::HashMap::new())), flow_trail: Arc::new(Mutex::new(server::flows::FlowTrail::default())), agent_focus: Arc::new(Mutex::new(server::focus::AgentFocus::load( server::focus::state_path(), diff --git a/crates/server/tests/agent_input_settle.rs b/crates/server/tests/agent_input_settle.rs index 67bec60a..74e3331f 100644 --- a/crates/server/tests/agent_input_settle.rs +++ b/crates/server/tests/agent_input_settle.rs @@ -1009,3 +1009,59 @@ fn the_settled_screen_answers_the_next_screenshot_without_a_capture() { assert_eq!(frames.load(Ordering::Acquire), after_action + 2); }); } + +/// An observed action that names no baseline diffs against the last screen +/// the same owner was handed, so a plain HTTP client gets only the change. +#[test] +fn an_observed_action_diffs_against_the_owners_last_read() { + block(async { + let wda = mock_wda(move |request, _| { + if is_session(request) { + return Some((Duration::ZERO, SESSION.to_string())); + } + if is_mutation(request) { + return Some((Duration::ZERO, r#"{"value":null}"#.to_string())); + } + if is_source(request) { + return Some((Duration::ZERO, simple_tree("搜索"))); + } + Some(( + Duration::ZERO, + r#"{"value":{"error":"no such alert","message":"no alert"}}"#.to_string(), + )) + }); + let state = build_state_with_wda(wda.url()); + let send = |method: &'static str, uri: &'static str, owner: Option<&'static str>| { + let state = state.clone(); + async move { + let mut builder = Request::builder() + .method(method) + .uri(uri) + .header("x-phone-control", "1"); + if let Some(owner) = owner { + builder = builder.header("x-phone-owner", owner); + } + let body = if method == "POST" { + builder = builder.header(header::CONTENT_TYPE, "application/json"); + Body::from(r#"{"type":"home"}"#) + } else { + Body::empty() + }; + let response = server::http::router(state) + .oneshot(builder.body(body).unwrap()) + .await + .unwrap(); + let bytes = response.into_body().collect().await.unwrap().to_bytes(); + serde_json::from_slice::(&bytes).unwrap() + } + }; + + let read = send("GET", "/agent/elements", Some("agent-a")).await; + let snapshot = read["snapshot"].as_str().unwrap().to_string(); + + let observed = send("POST", "/agent/input?return=delta", Some("agent-a")).await; + assert_eq!(observed["baseline"], snapshot, "{observed}"); + assert!(observed["delta"].is_object(), "{observed}"); + assert!(observed.get("elements").is_none(), "{observed}"); + }); +} diff --git a/crates/server/tests/fixtures/app_state.rs b/crates/server/tests/fixtures/app_state.rs index 99817816..d649e4ae 100644 --- a/crates/server/tests/fixtures/app_state.rs +++ b/crates/server/tests/fixtures/app_state.rs @@ -56,6 +56,9 @@ fn fixture_app_state(password: Option<&str>) -> Arc { element_snapshots: std::sync::Arc::new(std::sync::Mutex::new( std::collections::VecDeque::new(), )), + owner_snapshots: std::sync::Arc::new(std::sync::Mutex::new( + std::collections::HashMap::new(), + )), hold_until: std::sync::Arc::new(std::sync::Mutex::new(None)), owner: std::sync::Arc::new(std::sync::Mutex::new(None)), owner_lease_secs: 300,