From fd99075f312ba58b519a6ae7383d150a42cd4a67 Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Tue, 18 Aug 2026 11:11:21 -0600 Subject: [PATCH] fix(private-link): speak /app/devices and migrate stored ingest_url MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The three production request paths now use the devices spelling: register (RegistrationCoordinator::repair), the finish_owner_start capability ingest path, and EVENT_PATH. Day-segment listing follows for free because list_day derives from the ingest path. REGISTER_PATH and INGEST_PATH sit alongside the existing EVENT_PATH const. load_observer rewrites a persisted ingest_url that is exactly /app/observer/ingest to /app/devices/ingest, returns the rewritten value, and writes it back best-effort through write_observer_durably. The rewrite sits after observer_is_valid, so invalid and malformed state stay read-only, and a failed write leaves the predecessor bytes intact without failing session start. No /app/observer alias, fallback, or re-register. protocol_version stays 2. Observer-role identifiers, headers, and the observer.json filename are unchanged. vendor/observer-client-contract/ is unchanged — its register response fixtures still carry the predecessor ingest_url, which is correct payload data. --- crates/solstone-linux/src/cli.rs | 2 +- .../src/observer_contract_tests.rs | 12 +- crates/solstone-linux/src/private_link.rs | 143 ++++++++++++++++-- crates/solstone-linux/src/run.rs | 10 +- crates/solstone-linux/src/sync.rs | 34 ++--- crates/solstone-linux/src/test_support.rs | 10 +- crates/solstone-linux/src/upload.rs | 50 +++--- 7 files changed, 186 insertions(+), 75 deletions(-) diff --git a/crates/solstone-linux/src/cli.rs b/crates/solstone-linux/src/cli.rs index c86895e..be25aba 100644 --- a/crates/solstone-linux/src/cli.rs +++ b/crates/solstone-linux/src/cli.rs @@ -1083,7 +1083,7 @@ mod tests { let source = include_str!("cli.rs"); assert!(source.contains("setup_with_stream(")); assert!(!source.contains(&["UploadClient", "::new("].concat())); - assert!(!source.contains(&["/app/observer", "/register"].concat())); + assert!(!source.contains(&["/app/devices", "/register"].concat())); assert!(!source.contains(&[".bearer_", "auth("].concat())); } diff --git a/crates/solstone-linux/src/observer_contract_tests.rs b/crates/solstone-linux/src/observer_contract_tests.rs index f7562ab..95438c7 100644 --- a/crates/solstone-linux/src/observer_contract_tests.rs +++ b/crates/solstone-linux/src/observer_contract_tests.rs @@ -431,7 +431,7 @@ impl LinkedHarness { peer.credential(), "desktop", "K", - "/app/observer/ingest", + "/app/devices/ingest", ) .await; let client = client(&config(temp), owner.capability()); @@ -581,7 +581,7 @@ async fn assert_upload_contract( let request = &requests[0]; let body = String::from_utf8_lossy(&request.body); assert_eq!(request.method, "POST"); - assert_eq!(request.path, "/app/observer/ingest"); + assert_eq!(request.path, "/app/devices/ingest"); assert!(header(request, "authorization").is_some_and(|value| value.starts_with("Bearer "))); assert_eq!(body.matches("name=\"day\"").count(), 1); assert_eq!(body.matches("name=\"segment\"").count(), 1); @@ -700,7 +700,7 @@ async fn assert_listing_contract( let requests = harness.peer.requests(); let request = &requests[0]; assert_eq!(request.method, "GET"); - assert_eq!(request.path, "/app/observer/ingest/segments/20260618"); + assert_eq!(request.path, "/app/devices/ingest/segments/20260618"); assert_eq!(header(request, "x-solstone-protocol-version"), Some("2")); let authorization = header(request, "authorization"); assert!(authorization.is_some_and(|value| value.starts_with("Bearer "))); @@ -839,7 +839,7 @@ async fn assert_event_and_register( let request = &requests[0]; assert_eq!( (request.method.as_str(), request.path.as_str()), - ("POST", "/app/observer/ingest/event") + ("POST", "/app/devices/ingest/event") ); assert_eq!(header(request, "content-type"), Some("application/json")); assert_eq!( @@ -868,7 +868,7 @@ async fn assert_event_and_register( peer.credential(), label, "K", - "/app/observer/ingest", + "/app/devices/ingest", ) .await; assert!(matches!( @@ -882,7 +882,7 @@ async fn assert_event_and_register( let body: Value = serde_json::from_slice(&request.body).unwrap(); assert_eq!( (request.method.as_str(), request.path.as_str()), - ("POST", "/app/observer/register") + ("POST", "/app/devices/register") ); assert_eq!(body, fixtures[register_request]["payload"]); assert!(header(request, "authorization").is_none()); diff --git a/crates/solstone-linux/src/private_link.rs b/crates/solstone-linux/src/private_link.rs index a387a23..c2a7ac9 100644 --- a/crates/solstone-linux/src/private_link.rs +++ b/crates/solstone-linux/src/private_link.rs @@ -55,7 +55,9 @@ pub(crate) const OBSERVER_HEADER_NAME: &str = "x-solstone-observer"; pub(crate) const PROTOCOL_VERSION_HEADER_NAME: &str = "x-solstone-protocol-version"; const REGISTRATION_MARKER_HEADER_NAME: &str = "x-solstone-linux-registration-route"; const REGISTRATION_MARKER_HEADER_VALUE: &str = "1"; -const EVENT_PATH: &str = "/app/observer/ingest/event"; +const EVENT_PATH: &str = "/app/devices/ingest/event"; +const REGISTER_PATH: &str = "/app/devices/register"; +const INGEST_PATH: &str = "/app/devices/ingest"; #[derive(Clone, Deserialize, Eq, PartialEq, Serialize)] #[serde(deny_unknown_fields)] @@ -748,6 +750,7 @@ pub(crate) fn load_observer( credential_instance_id: &str, expected_name: &str, origin: &Url, + fault: &dyn DurableWriteFault, ) -> Result, PrivateStateError> { let Some(bytes) = read_private_file( &config_root.join(OBSERVER_FILENAME), @@ -756,11 +759,15 @@ pub(crate) fn load_observer( else { return Ok(None); }; - let observer = serde_json::from_slice::(&bytes) + let mut observer = serde_json::from_slice::(&bytes) .map_err(|_| PrivateStateError::MalformedObserver)?; if !observer_is_valid(&observer, credential_instance_id, expected_name, origin) { return Ok(None); } + if observer.ingest_url == "/app/observer/ingest" { + observer.ingest_url = INGEST_PATH.to_owned(); + let _ = write_observer_durably(config_root, &observer, fault); + } Ok(Some(observer)) } @@ -1435,7 +1442,7 @@ impl RegistrationCoordinator { "stream_type": "desktop", "version": self.version, }); - let Ok(url) = confine_path(&self.origin, "/app/observer/register") else { + let Ok(url) = confine_path(&self.origin, REGISTER_PATH) else { return RepairOutcome::InvalidRegistration; }; let response = match self @@ -1601,7 +1608,7 @@ impl PrivateLinkOwner { #[cfg(test)] async fn register(&self, body: &serde_json::Value) -> LinkOutcome { - let Ok(url) = confine_path(&self.capability.inner.origin, "/app/observer/register") else { + let Ok(url) = confine_path(&self.capability.inner.origin, REGISTER_PATH) else { return LinkOutcome::LocalRejected { status: StatusCode::BAD_REQUEST, }; @@ -1648,7 +1655,7 @@ pub(crate) async fn start_private_link_owner( async fn finish_owner_start( mut session: PrivateLinkSession, ) -> Result { - let capability = session.capability("/app/observer/ingest".to_owned()); + let capability = session.capability(INGEST_PATH.to_owned()); if session.opener.generation() == 0 { match capability.report_unauthorized(0).await { RepairOutcome::Repaired { .. } | RepairOutcome::AlreadySuperseded { .. } => {} @@ -1748,7 +1755,7 @@ impl PrivateLinkSession { #[cfg(test)] pub(crate) async fn register_for_test(&self, body: &serde_json::Value) -> LinkOutcome { - let Ok(url) = confine_path(&self.origin, "/app/observer/register") else { + let Ok(url) = confine_path(&self.origin, REGISTER_PATH) else { return LinkOutcome::LocalRejected { status: StatusCode::BAD_REQUEST, }; @@ -2030,6 +2037,7 @@ async fn start_private_link_session_inner( .iter() .map(|endpoint| endpoint.host.clone()) .collect(); + let persistence_fault = Arc::clone(&options.persistence_fault); let (token_persistence, hook) = TokenPersistence::new( config_root.clone(), credential.clone(), @@ -2153,6 +2161,7 @@ async fn start_private_link_session_inner( &credential_instance_id, &expected_name, &origin, + persistence_fault.as_ref(), ) { Ok(Some(observer)) => { opener.set_registered(&observer)?; @@ -2363,6 +2372,7 @@ mod tests { "instance", "stream", &Url::parse("http://127.0.0.1:1").unwrap(), + &NoWriteFault, ) .unwrap(); if let Some(observer) = loaded { @@ -3108,7 +3118,7 @@ mod tests { assert_eq!(load_credential(temp.path()).unwrap(), Some(credential())); let origin = Url::parse("http://127.0.0.1:1234").unwrap(); assert!( - load_observer(temp.path(), "instance", "stream", &origin) + load_observer(temp.path(), "instance", "stream", &origin, &NoWriteFault) .unwrap() .is_some() ); @@ -3123,6 +3133,42 @@ mod tests { ); } + #[test] + fn load_observer_rewrites_predecessor_ingest_url_exactly() { + let origin = Url::parse("http://127.0.0.1:1234").unwrap(); + + let predecessor = tempfile::tempdir().unwrap(); + ensure_private_directory(predecessor.path()).unwrap(); + persist_observer(predecessor.path(), &observer("/app/observer/ingest")).unwrap(); + let rewritten = load_observer( + predecessor.path(), + "instance", + "stream", + &origin, + &NoWriteFault, + ) + .unwrap() + .unwrap(); + assert_eq!(rewritten.ingest_url, "/app/devices/ingest"); + let persisted: ObserverState = + serde_json::from_slice(&fs::read(predecessor.path().join(OBSERVER_FILENAME)).unwrap()) + .unwrap(); + assert_eq!(persisted.ingest_url, "/app/devices/ingest"); + + let devices = tempfile::tempdir().unwrap(); + ensure_private_directory(devices.path()).unwrap(); + persist_observer(devices.path(), &observer("/app/devices/ingest")).unwrap(); + let prior_bytes = fs::read(devices.path().join(OBSERVER_FILENAME)).unwrap(); + let loaded = load_observer(devices.path(), "instance", "stream", &origin, &NoWriteFault) + .unwrap() + .unwrap(); + assert_eq!(loaded.ingest_url, "/app/devices/ingest"); + assert_eq!( + fs::read(devices.path().join(OBSERVER_FILENAME)).unwrap(), + prior_bytes + ); + } + #[test] fn malformed_state_errors_are_distinct() { let temp = tempfile::tempdir().unwrap(); @@ -3134,7 +3180,7 @@ mod tests { fs::write(temp.path().join(OBSERVER_FILENAME), b"{").unwrap(); let origin = Url::parse("http://127.0.0.1:1").unwrap(); assert!(matches!( - load_observer(temp.path(), "instance", "stream", &origin), + load_observer(temp.path(), "instance", "stream", &origin, &NoWriteFault), Err(PrivateStateError::MalformedObserver) )); } @@ -3173,7 +3219,7 @@ mod tests { let bytes = serde_json::to_vec(&state).unwrap(); fs::write(temp.path().join(OBSERVER_FILENAME), &bytes).unwrap(); assert!( - load_observer(temp.path(), "instance", "stream", &origin) + load_observer(temp.path(), "instance", "stream", &origin, &NoWriteFault) .unwrap() .is_none() ); @@ -3204,7 +3250,7 @@ mod tests { let bytes = serde_json::to_vec(&observer(path)).unwrap(); fs::write(temp.path().join(OBSERVER_FILENAME), &bytes).unwrap(); assert!( - load_observer(temp.path(), "instance", "stream", &origin) + load_observer(temp.path(), "instance", "stream", &origin, &NoWriteFault) .unwrap() .is_none(), "{path}" @@ -3233,6 +3279,7 @@ mod tests { "instance", "stream", &Url::parse("http://127.0.0.1:1").unwrap(), + &NoWriteFault, ) .map(|_| ()) }; @@ -3525,7 +3572,7 @@ mod tests { .unwrap(); assert_eq!( session - .request(Method::POST, "/app/observer/register") + .request(Method::POST, "/app/devices/register") .unwrap() .header( REGISTRATION_MARKER_HEADER_NAME, @@ -3541,7 +3588,7 @@ mod tests { &session, &ObserverState { credential_instance_id, - ..observer("/app/observer/ingest") + ..observer("/app/devices/ingest") }, ) .unwrap(); @@ -3559,7 +3606,7 @@ mod tests { let requests = peer.requests(); assert_eq!(requests.len(), 2); assert_eq!(requests[0].method, "POST"); - assert_eq!(requests[0].path, "/app/observer/register"); + assert_eq!(requests[0].path, "/app/devices/register"); assert!(requests[0].body.is_empty()); assert_eq!( requests[0] @@ -3579,13 +3626,13 @@ mod tests { assert!(requests[1].body.is_empty()); assert_registered_auth(&requests[1], "observer-key"); assert_eq!(peer.accepted_carriers(), 1); - let capability = session.capability("/app/observer/ingest".to_owned()); + let capability = session.capability("/app/devices/ingest".to_owned()); let second_generation = publish_observer_registration( &session, &ObserverState { credential_instance_id: session.credential_instance_id.clone(), key: "replacement-key".into(), - ..observer("/app/observer/ingest") + ..observer("/app/devices/ingest") }, ) .unwrap(); @@ -3690,7 +3737,7 @@ mod tests { async fn capability_rejects_admin_path_query_and_route_substitution() { let peer = PrivateLinkPeer::start().await; let (_temp, session) = start_peer_session(&peer).await; - let capability = session.capability("/app/observer/ingest".to_owned()); + let capability = session.capability("/app/devices/ingest".to_owned()); for day in ["", "2026010", "202601011", "202601?1", "../20260101"] { assert!(matches!( capability.list_day(day).await, @@ -4257,6 +4304,7 @@ mod tests { &restarted.credential_instance_id, "stream", &restarted.origin, + &NoWriteFault, ) .unwrap(); assert_eq!( @@ -4509,6 +4557,69 @@ mod tests { assert_token_failure(DurableWriteStage::Fsync).await; } + #[tokio::test] + async fn session_start_rewrites_predecessor_observer_best_effort() { + let peer = PrivateLinkPeer::start().await; + let predecessor = ObserverState { + credential_instance_id: peer.credential().instance_id, + ..observer("/app/observer/ingest") + }; + + let unfaulted = tempfile::tempdir().unwrap(); + persist_observer(unfaulted.path(), &predecessor).unwrap(); + let session = start_private_link_session(unfaulted.path(), peer.credential(), "stream") + .await + .unwrap(); + assert_eq!(session.opener.generation(), 1); + let persisted: ObserverState = + serde_json::from_slice(&fs::read(unfaulted.path().join(OBSERVER_FILENAME)).unwrap()) + .unwrap(); + assert_eq!(persisted.ingest_url, "/app/devices/ingest"); + session.shutdown().await.unwrap(); + + let faulted = tempfile::tempdir().unwrap(); + persist_observer(faulted.path(), &predecessor).unwrap(); + let prior_bytes = fs::read(faulted.path().join(OBSERVER_FILENAME)).unwrap(); + let session = start_private_link_session_inner( + faulted.path(), + peer.credential(), + "stream", + SessionStartOptions { + persistence_fault: Arc::new(RecordingFault { + stages: Arc::new(Mutex::new(Vec::new())), + fail: Some(DurableWriteStage::Write), + }), + ..SessionStartOptions::default() + }, + ) + .await + .unwrap(); + assert_eq!(session.opener.generation(), 1); + assert_eq!( + fs::read(faulted.path().join(OBSERVER_FILENAME)).unwrap(), + prior_bytes + ); + let loaded = load_observer( + faulted.path(), + &session.credential_instance_id, + "stream", + &session.origin, + &RecordingFault { + stages: Arc::new(Mutex::new(Vec::new())), + fail: Some(DurableWriteStage::Write), + }, + ) + .unwrap() + .unwrap(); + assert_eq!(loaded.ingest_url, "/app/devices/ingest"); + assert_eq!( + fs::read(faulted.path().join(OBSERVER_FILENAME)).unwrap(), + prior_bytes + ); + session.shutdown().await.unwrap(); + peer.shutdown().await; + } + #[tokio::test] async fn owner_shutdown_joins_bridge_and_closes_listener() { let temp = tempfile::tempdir().unwrap(); diff --git a/crates/solstone-linux/src/run.rs b/crates/solstone-linux/src/run.rs index 89e387e..e2376b9 100644 --- a/crates/solstone-linux/src/run.rs +++ b/crates/solstone-linux/src/run.rs @@ -1353,7 +1353,7 @@ mod tests { let session = start_private_link_session(temp.path(), peer.credential(), "stream") .await .unwrap(); - let capability = session.capability("/app/observer/ingest".to_owned()); + let capability = session.capability("/app/devices/ingest".to_owned()); let config = Config { config_dir: temp.path().to_path_buf(), stream: "stream".to_owned(), @@ -1380,7 +1380,7 @@ mod tests { assert!( peer.requests() .iter() - .all(|request| request.path == "/app/observer/register") + .all(|request| request.path == "/app/devices/register") ); response_gate.store(true, Ordering::Release); peer.notify_response_gates(); @@ -1585,7 +1585,7 @@ mod tests { key: "K".to_owned(), prefix: "prefix".to_owned(), name: "stream".to_owned(), - ingest_url: "/app/observer/ingest".to_owned(), + ingest_url: "/app/devices/ingest".to_owned(), protocol_version: 2, }, ) @@ -1597,7 +1597,7 @@ mod tests { }; let client = Arc::new(UploadClient::new( &config, - session.capability("/app/observer/ingest".to_owned()), + session.capability("/app/devices/ingest".to_owned()), "host", "linux", "1", @@ -1640,7 +1640,7 @@ mod tests { reqwest::multipart::Part::stream(reqwest::Body::wrap_stream(stream)), ); let request = tokio::spawn({ - let capability = session.capability("/app/observer/ingest".to_owned()); + let capability = session.capability("/app/devices/ingest".to_owned()); async move { capability.ingest(form).await } }); assert_real_observer_ticks_advance(); diff --git a/crates/solstone-linux/src/sync.rs b/crates/solstone-linux/src/sync.rs index 30edf7b..034e2cb 100644 --- a/crates/solstone-linux/src/sync.rs +++ b/crates/solstone-linux/src/sync.rs @@ -1275,14 +1275,14 @@ mod tests { key: "STALE-KEY-FULL".to_owned(), prefix: "prefix".to_owned(), name: "desktop".to_owned(), - ingest_url: "/app/observer/ingest".to_owned(), + ingest_url: "/app/devices/ingest".to_owned(), protocol_version: 2, }, ) .unwrap(); let client = Arc::new(UploadClient::new( &config, - session.capability("/app/observer/ingest".to_owned()), + session.capability("/app/devices/ingest".to_owned()), "host", "linux", "test", @@ -1333,7 +1333,7 @@ mod tests { server .logged_requests() .iter() - .filter(|request| request.uri == "/app/observer/ingest") + .filter(|request| request.uri == "/app/devices/ingest") .count() } @@ -2709,7 +2709,7 @@ mod tests { assert_eq!( peer.requests() .iter() - .filter(|request| request.path == "/app/observer/register") + .filter(|request| request.path == "/app/devices/register") .count(), 1 ); @@ -2756,7 +2756,7 @@ mod tests { assert_eq!( peer.requests() .iter() - .filter(|request| request.path == "/app/observer/register") + .filter(|request| request.path == "/app/devices/register") .count(), 1 ); @@ -3086,7 +3086,7 @@ mod tests { let capped_steps = (14_400.0 / CIRCUIT_COOLDOWN_MAX).ceil() as usize; let bound = CIRCUIT_THRESHOLD_TRANSIENT as usize + ramp_steps + capped_steps + 1; assert!( - server.request_count("/app/observer/ingest/segments/") <= bound, + server.request_count("/app/devices/ingest/segments/") <= bound, "listing retries exceeded conservative bound {bound}" ); assert_eq!(worker.circuit_cooldown, CIRCUIT_COOLDOWN_MAX); @@ -3303,7 +3303,7 @@ mod tests { key: "K".to_owned(), prefix: "prefix".into(), name: "host".into(), - ingest_url: "/app/observer/ingest".into(), + ingest_url: "/app/devices/ingest".into(), protocol_version: 2, }, ) @@ -3314,7 +3314,7 @@ mod tests { }); let client = Arc::new(UploadClient::new( &config, - session.capability("/app/observer/ingest".into()), + session.capability("/app/devices/ingest".into()), "host", "linux", "test", @@ -3370,7 +3370,7 @@ mod tests { key: "K".to_owned(), prefix: "prefix".into(), name: "host".into(), - ingest_url: "/app/observer/ingest".into(), + ingest_url: "/app/devices/ingest".into(), protocol_version: 2, }, ) @@ -3381,7 +3381,7 @@ mod tests { }); let client = Arc::new(UploadClient::new( &config, - session.capability("/app/observer/ingest".into()), + session.capability("/app/devices/ingest".into()), "host", "linux", "test", @@ -4155,7 +4155,7 @@ mod tests { .unwrap(); let client = Arc::new(UploadClient::new( &config, - session.capability("/app/observer/ingest".to_owned()), + session.capability("/app/devices/ingest".to_owned()), "host", "linux", "test", @@ -4217,7 +4217,7 @@ mod tests { .unwrap(); let client = Arc::new(UploadClient::new( &config, - session.capability("/app/observer/ingest".to_owned()), + session.capability("/app/devices/ingest".to_owned()), "host", "linux", "test", @@ -4285,7 +4285,7 @@ mod tests { let clock = Arc::new(MutableClock::new(1_800_000_000.0, 0.0)); let client = Arc::new(UploadClient::new( &config, - session.capability("/app/observer/ingest".to_owned()), + session.capability("/app/devices/ingest".to_owned()), "host", "linux", "test", @@ -4325,7 +4325,7 @@ mod tests { assert_eq!( peer.requests() .iter() - .filter(|request| request.path == "/app/observer/register") + .filter(|request| request.path == "/app/devices/register") .count(), 1, "registration requests during cooldown" @@ -4338,7 +4338,7 @@ mod tests { if peer .requests() .iter() - .filter(|request| request.path == "/app/observer/register") + .filter(|request| request.path == "/app/devices/register") .count() >= 2 { @@ -4349,7 +4349,7 @@ mod tests { assert_eq!( peer.requests() .iter() - .filter(|request| request.path == "/app/observer/register") + .filter(|request| request.path == "/app/devices/register") .count(), 2, "registration requests after cooldown" @@ -4466,7 +4466,7 @@ mod tests { let upload = server .requests() .into_iter() - .find(|request| request.uri == "/app/observer/ingest") + .find(|request| request.uri == "/app/devices/ingest") .expect("final segment upload"); assert!(String::from_utf8_lossy(&upload.body).contains("final screen")); let facts = load_facts(&config.state_dir()); diff --git a/crates/solstone-linux/src/test_support.rs b/crates/solstone-linux/src/test_support.rs index c29726c..f88ee4d 100644 --- a/crates/solstone-linux/src/test_support.rs +++ b/crates/solstone-linux/src/test_support.rs @@ -178,7 +178,7 @@ impl LinkedMockServer { key: "K".to_owned(), prefix: "prefix".to_owned(), name: "desktop".to_owned(), - ingest_url: "/app/observer/ingest".to_owned(), + ingest_url: "/app/devices/ingest".to_owned(), protocol_version: 2, }, ) @@ -187,7 +187,7 @@ impl LinkedMockServer { } pub(crate) fn capability(&self) -> PrivateLinkCapability { - self.session.capability("/app/observer/ingest".to_owned()) + self.session.capability("/app/devices/ingest".to_owned()) } pub(crate) fn requests(&self) -> Vec { @@ -438,12 +438,12 @@ mod tests { #[tokio::test] async fn counts_requests_by_uri_substring() { let server = MockServer::new(vec![(200, json!({}))]).await; - reqwest::get(format!("{}/app/observer/segments?day=20260101", server.url)) + reqwest::get(format!("{}/app/devices/segments?day=20260101", server.url)) .await .unwrap(); wait_for_requests(&server, 1).await; - assert_eq!(server.request_count("/app/observer/segments"), 1); - assert_eq!(server.request_count("/app/observer/ingest"), 0); + assert_eq!(server.request_count("/app/devices/segments"), 1); + assert_eq!(server.request_count("/app/devices/ingest"), 0); } // AC: shared HTTP harness can emit test-controlled streaming chunks. diff --git a/crates/solstone-linux/src/upload.rs b/crates/solstone-linux/src/upload.rs index 2ccb55d..a20363a 100644 --- a/crates/solstone-linux/src/upload.rs +++ b/crates/solstone-linux/src/upload.rs @@ -972,14 +972,14 @@ mod tests { key: "K".to_owned(), prefix: "prefix".into(), name: "host-a".into(), - ingest_url: "/app/observer/ingest".into(), + ingest_url: "/app/devices/ingest".into(), protocol_version: 2, }, ) .unwrap(); let client = UploadClient::new( &config, - session.capability("/app/observer/ingest".into()), + session.capability("/app/devices/ingest".into()), "host-a", "linux", "0.1.0", @@ -1013,14 +1013,14 @@ mod tests { key: "STALE-KEY".into(), prefix: "prefix".into(), name: "desktop".into(), - ingest_url: "/app/observer/ingest".into(), + ingest_url: "/app/devices/ingest".into(), protocol_version: 2, }, ) .unwrap(); let client = UploadClient::new( &config, - session.capability("/app/observer/ingest".into()), + session.capability("/app/devices/ingest".into()), "host", "linux", "test", @@ -1277,14 +1277,14 @@ mod tests { key: "K123456789".into(), prefix: "prefix".into(), name: "host-a".into(), - ingest_url: "/app/observer/ingest".into(), + ingest_url: "/app/devices/ingest".into(), protocol_version: 2, }, ) .unwrap(); let client = UploadClient::new( &config, - session.capability("/app/observer/ingest".into()), + session.capability("/app/devices/ingest".into()), "host-a", "linux", "0.1.0", @@ -1321,7 +1321,7 @@ mod tests { assert_eq!(requests.len(), 5); assert_eq!( (requests[0].method.as_str(), requests[0].path.as_str()), - ("POST", "/app/observer/register") + ("POST", "/app/devices/register") ); assert_eq!( requests[0] @@ -1349,7 +1349,7 @@ mod tests { ); assert_eq!( (requests[1].method.as_str(), requests[1].path.as_str()), - ("POST", "/app/observer/ingest") + ("POST", "/app/devices/ingest") ); assert_eq!( requests[1] @@ -1410,7 +1410,7 @@ mod tests { ); assert_eq!( (requests[2].method.as_str(), requests[2].path.as_str()), - ("GET", "/app/observer/ingest/segments/20260101") + ("GET", "/app/devices/ingest/segments/20260101") ); assert_eq!( requests[2] @@ -1444,7 +1444,7 @@ mod tests { ] { assert_eq!( (request.method.as_str(), request.path.as_str()), - ("POST", "/app/observer/ingest/event") + ("POST", "/app/devices/ingest/event") ); assert_eq!( serde_json::from_slice::(&request.body).unwrap(), @@ -1501,8 +1501,8 @@ mod tests { assert!(!client.relay_event("observe", "status", Map::new()).await); let requests = peer.requests(); assert_eq!(requests.len(), 2); - assert_eq!(requests[0].path, "/app/observer/ingest/event"); - assert_eq!(requests[1].path, "/app/observer/register"); + assert_eq!(requests[0].path, "/app/devices/ingest/event"); + assert_eq!(requests[1].path, "/app/devices/register"); session.shutdown().await.unwrap(); peer.shutdown().await; } @@ -1533,7 +1533,7 @@ mod tests { assert_eq!( peer.requests() .iter() - .filter(|request| request.path == "/app/observer/register") + .filter(|request| request.path == "/app/devices/register") .count(), 1 ); @@ -1577,7 +1577,7 @@ mod tests { assert_eq!( peer.requests() .iter() - .filter(|request| request.path == "/app/observer/register") + .filter(|request| request.path == "/app/devices/register") .count(), 1 ); @@ -1603,7 +1603,7 @@ mod tests { assert_eq!( peer.requests() .iter() - .filter(|request| request.path == "/app/observer/register") + .filter(|request| request.path == "/app/devices/register") .count(), 1 ); @@ -1612,7 +1612,7 @@ mod tests { assert_eq!( peer.requests() .iter() - .filter(|request| request.path == "/app/observer/register") + .filter(|request| request.path == "/app/devices/register") .count(), 2 ); @@ -1754,7 +1754,7 @@ mod tests { ); let request = &server.requests()[0]; assert_eq!(request.method, "POST"); - assert_eq!(request.uri, "/app/observer/ingest"); + assert_eq!(request.uri, "/app/devices/ingest"); assert_eq!(request.headers["authorization"], "Bearer K"); assert_eq!(request.headers[OBSERVER_PROTOCOL_VERSION_HEADER], "2"); let body = String::from_utf8_lossy(&request.body); @@ -1984,7 +1984,7 @@ mod tests { .await ); let request = &server.requests()[0]; - assert_eq!(request.uri, "/app/observer/ingest/event"); + assert_eq!(request.uri, "/app/devices/ingest/event"); assert_eq!(request.headers["authorization"], "Bearer K"); let body: Value = serde_json::from_slice(&request.body).unwrap(); assert_eq!( @@ -2014,7 +2014,7 @@ mod tests { assert_eq!(result.segments, Some(vec![])); assert!(result.legacy); let request = &server.requests()[0]; - assert_eq!(request.uri, "/app/observer/ingest/segments/20260101"); + assert_eq!(request.uri, "/app/devices/ingest/segments/20260101"); assert_eq!(request.headers["authorization"], "Bearer K"); assert_eq!(request.headers[OBSERVER_PROTOCOL_VERSION_HEADER], "2"); } @@ -2130,7 +2130,7 @@ mod tests { .unwrap(); let client = UploadClient::new( &config, - session.capability("/app/observer/ingest".to_owned()), + session.capability("/app/devices/ingest".to_owned()), "host", "linux", "test", @@ -2290,7 +2290,7 @@ mod tests { }; let client = UploadClient::new( &config, - session.capability("/app/observer/ingest".to_owned()), + session.capability("/app/devices/ingest".to_owned()), "host", "linux", "1", @@ -2368,14 +2368,14 @@ mod tests { key: "K".to_owned(), prefix: "prefix".into(), name: "host-a".into(), - ingest_url: "/app/observer/ingest".into(), + ingest_url: "/app/devices/ingest".into(), protocol_version: 2, }, ) .unwrap(); let client = UploadClient::new( &config, - session.capability("/app/observer/ingest".into()), + session.capability("/app/devices/ingest".into()), "host-a", "linux", "0.1.0", @@ -2507,7 +2507,7 @@ mod tests { assert_eq!( peer.requests() .iter() - .filter(|request| request.path == "/app/observer/register") + .filter(|request| request.path == "/app/devices/register") .count(), 1 ); @@ -2525,7 +2525,7 @@ mod tests { assert_eq!( peer.requests() .iter() - .filter(|request| request.path == "/app/observer/register") + .filter(|request| request.path == "/app/devices/register") .count(), 1 ); -- 2.51.2