Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
41 changes: 41 additions & 0 deletions crates/server/src/http.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Mutex<ElementSnapshotCache>>,
/// 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<Mutex<std::collections::HashMap<String, String>>>,
/// 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<Mutex<crate::flows::FlowTrail>>,
Expand Down Expand Up @@ -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<String> {
let owner = named_owner(headers)?;
recover(state.owner_snapshots.lock()).get(owner).cloned()
}

fn lookup_element_snapshot(
state: &AppState,
snapshot: &str,
Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -9714,13 +9748,19 @@ 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()
.filter(|since| !since.is_empty())
.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))
Expand Down Expand Up @@ -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.
Expand Down
1 change: 1 addition & 0 deletions crates/server/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand Down
56 changes: 56 additions & 0 deletions crates/server/tests/agent_input_settle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<serde_json::Value>(&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}");
});
}
3 changes: 3 additions & 0 deletions crates/server/tests/fixtures/app_state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,9 @@ fn fixture_app_state(password: Option<&str>) -> Arc<AppState> {
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,
Expand Down
Loading