From 43b0e32271752f3ccdfcdd5bf8e5e7e0993e015e Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Fri, 31 Jul 2026 13:34:48 -0600 Subject: [PATCH] reach the journal away from the home network and retry registration Use separate LAN and relay transport clients, dialled sequentially inside one admission-guarded future. Bound the LAN leg at 5 seconds only when a relay client exists, because that deadline exists solely to reserve caller budget for the alternate leg; credentials with no relay origin keep today's unbounded ladder. This survives the untimed TLS-handshake case that ladder shortening alone cannot fix, without raising any caller budget. Only the relay client holds the refreshable device token, preserving exactly one refresh single-flight. Wire the existing ensure_registered path into the sync worker's unregistered branch. After a failed owner start, the first attempt fires on the first wake, whether a notification or the 60-second periodic tick, because owner start does not stamp last_recovery_attempt; after that, the 300-second recovery cooldown is the bound. Keep the terminal signals distinct: RepairOutcome::GuardRefused at owner start still ends startup with RegistrationInvalid, while UploadClient::is_revoked(), latched by a 403, is checked at the sync call site because ensure_registered deliberately does not consult revocation. Co-Authored-By: Claude Opus 5 (1M context) --- CHANGELOG.md | 1 + crates/solstone-linux/src/private_link.rs | 353 +++++++++++++++++- .../src/private_link_test_peer.rs | 11 + crates/solstone-linux/src/sync.rs | 231 +++++++++++- crates/solstone-linux/src/upload.rs | 5 + 5 files changed, 585 insertions(+), 16 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 7c87834..4cc0de6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,7 @@ and this project adheres to Semantic Versioning. ## [Unreleased] ### Fixed +- sol now reaches your journal even when your home network is out of reach, and no longer waits on a local connection that will never answer. if the first connection at startup does not get through, sol keeps trying on its own instead of staying quiet until you restart it. - sol no longer shows sync status left behind by a previous run while a new run is still starting. - a single refusal from your journal no longer stops uploads for the rest of the session. one refused request used to close sending for as long as sol kept running, with nothing shown anywhere. a refusal sol can recover from now backs off and retries on the existing schedule, and a device you have removed from your journal still stops for good. - sol now recovers on its own when your journal no longer recognises this machine. sol used to keep running with nothing arriving and nothing said about it. it now renews its identity with your journal once, resumes where it left off, and writes down what it did so you can confirm afterwards that it was a recovery. a machine you have removed from your journal is never brought back this way. diff --git a/crates/solstone-linux/src/private_link.rs b/crates/solstone-linux/src/private_link.rs index 61880be..8f54eff 100644 --- a/crates/solstone-linux/src/private_link.rs +++ b/crates/solstone-linux/src/private_link.rs @@ -45,6 +45,7 @@ pub(crate) const PRIVATE_STATE_READY_LOCK_FILENAME: &str = const MAX_PAIR_LINK_BYTES: u64 = 4096; pub(crate) const MAX_REQUEST_BODY_BYTES: u64 = 64 * 1024 * 1024; const LOOPBACK_CONNECT_TIMEOUT: Duration = Duration::from_secs(5); +const LAN_CARRIER_TIMEOUT: Duration = Duration::from_secs(5); const BOOTSTRAP_TIMEOUT: Duration = Duration::from_secs(30); const REGISTRATION_TIMEOUT: Duration = Duration::from_secs(30); const INGEST_TIMEOUT: Duration = Duration::from_secs(300); @@ -953,7 +954,8 @@ struct AuthEpoch { } struct PrivateLinkOpener { - transport: Arc, + lan_transport: Option>, + relay_transport: Option>, auth: RwLock>, admission: tokio::sync::Mutex<()>, transport_unavailable: Arc, @@ -962,12 +964,14 @@ struct PrivateLinkOpener { impl PrivateLinkOpener { fn new( - transport: TransportClient, + lan_transport: Option, + relay_transport: Option, transport_unavailable: Arc, facts: LinkFacts, ) -> Self { Self { - transport: Arc::new(transport), + lan_transport: lan_transport.map(Arc::new), + relay_transport: relay_transport.map(Arc::new), auth: RwLock::new(Arc::new(AuthEpoch { generation: 0, state: OpenerAuth::Unregistered, @@ -1028,8 +1032,9 @@ impl PrivateLinkOpener { "linked transport unavailable".into(), )); } - // Pinned client.rs:265 and :296 are the only refresh callers, both - // inside dial_carrier_over_relay reached by TransportClient::dial_carrier. + // The relay transport is the only client that holds a device token. Pinned + // client.rs:265 and :296 refresh only inside its dial_carrier relay path, so + // this guard keeps a refreshed carrier behind synchronous token persistence. // There is no timer, background, or live-stream refresh path. let _admission = self.admission.lock().await; if self.transport_unavailable.load(Ordering::Acquire) { @@ -1064,7 +1069,40 @@ impl CarrierOpener for PrivateLinkOpener { &self, ) -> Pin> + Send + '_>> { Box::pin(async move { - let carrier = self.admit_dial(self.transport.dial_carrier()).await?; + let carrier = self + .admit_dial(async { + let mut lan_error = None; + if let Some(lan) = &self.lan_transport { + // Reserve time for the alternate relay leg only when it exists. A + // LAN-only credential keeps the pinned client's behavior unchanged. + let result = if self.relay_transport.is_some() { + match tokio::time::timeout(LAN_CARRIER_TIMEOUT, lan.dial_carrier()) + .await + { + Ok(result) => result, + Err(_) => Err(TransportError::Io(io::Error::new( + io::ErrorKind::TimedOut, + "LAN carrier dial timed out", + ))), + } + } else { + lan.dial_carrier().await + }; + match result { + Ok(carrier) => return Ok(carrier), + Err(error) => lan_error = Some(error), + } + } + if let Some(relay) = &self.relay_transport { + match relay.dial_carrier().await { + Ok(carrier) => return Ok(carrier), + Err(error) if lan_error.is_none() => return Err(error), + Err(_) => {} + } + } + Err(lan_error.unwrap_or(TransportError::NoEndpoint)) + }) + .await?; self.facts.publish(LinkFact::CarrierProven); Ok(carrier) }) @@ -1999,14 +2037,41 @@ async fn start_private_link_session_inner( transport_unavailable.clone(), facts.clone(), ); - let transport = if credential.endpoints.is_empty() { - TransportClient::new_relay_only(credential, Some(hook)) + let lan_transport = if credential.endpoints.is_empty() { + None } else { - TransportClient::new(credential, Some(hook)) - } - .map_err(|_| PrivateStateError::BridgeUnavailable)?; + let mut lan_credential = credential.clone(); + lan_credential.relay_origin = None; + lan_credential.device_token = None; + lan_credential.device_token_expires_at = None; + Some( + TransportClient::new(lan_credential, None) + .map_err(|_| PrivateStateError::BridgeUnavailable)?, + ) + }; + let relay_transport = if lan_transport.is_none() { + Some( + TransportClient::new_relay_only(credential.clone(), Some(hook)) + .map_err(|_| PrivateStateError::BridgeUnavailable)?, + ) + } else if credential + .relay_origin + .as_deref() + .is_some_and(|origin| !origin.is_empty()) + && credential + .device_token + .as_deref() + .is_some_and(|token| !token.is_empty()) + { + let mut relay_credential = credential.clone(); + relay_credential.endpoints.clear(); + TransportClient::new_relay_only(relay_credential, Some(hook)).ok() + } else { + None + }; let opener = Arc::new(PrivateLinkOpener::new( - transport, + lan_transport, + relay_transport, transport_unavailable, facts.clone(), )); @@ -2156,6 +2221,7 @@ mod tests { use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, net::TcpStream, + sync::Notify, }; async fn raw_local_request(port: u16, request: String) -> Vec { @@ -2287,7 +2353,8 @@ mod tests { persist_observer(temp.path(), &state).unwrap(); let peer = PrivateLinkPeer::start().await; let opener = PrivateLinkOpener::new( - TransportClient::new(peer.credential(), None).unwrap(), + Some(TransportClient::new(peer.credential(), None).unwrap()), + None, Arc::new(AtomicBool::new(false)), LinkFacts::default(), ); @@ -2651,7 +2718,8 @@ mod tests { let transport_unavailable = Arc::new(AtomicBool::new(false)); let facts = LinkFacts::default(); let opener = Arc::new(PrivateLinkOpener::new( - TransportClient::new(peer.credential(), None).unwrap(), + Some(TransportClient::new(peer.credential(), None).unwrap()), + None, transport_unavailable.clone(), facts.clone(), )); @@ -2715,7 +2783,8 @@ mod tests { async fn real_transport_dial_waits_for_admission_guard() { let peer = PrivateLinkPeer::start().await; let opener = Arc::new(PrivateLinkOpener::new( - TransportClient::new(peer.credential(), None).unwrap(), + Some(TransportClient::new(peer.credential(), None).unwrap()), + None, Arc::new(AtomicBool::new(false)), LinkFacts::default(), )); @@ -2731,6 +2800,260 @@ mod tests { peer.shutdown().await; } + #[tokio::test(start_paused = true)] + async fn stalled_lan_tls_reaches_relay_within_registration_budget() { + let temp = tempfile::tempdir().unwrap(); + let peer = PrivateLinkPeer::start().await; + let lan = tokio::net::TcpListener::bind(("127.0.0.1", 0)) + .await + .unwrap(); + let lan_port = lan.local_addr().unwrap().port(); + let lan_accepted = Arc::new(AtomicUsize::new(0)); + let lan_count = lan_accepted.clone(); + let stalled = tokio::spawn(async move { + let (_stream, _) = lan.accept().await.unwrap(); + lan_count.fetch_add(1, Ordering::SeqCst); + std::future::pending::<()>().await; + }); + let relay = tokio::net::TcpListener::bind(("127.0.0.1", 0)) + .await + .unwrap(); + let relay_origin = format!("http://{}", relay.local_addr().unwrap()); + let relay_accepted = Arc::new(AtomicUsize::new(0)); + let relay_count = relay_accepted.clone(); + let relay_task = tokio::spawn(async move { + let (mut stream, _) = relay.accept().await.unwrap(); + relay_count.fetch_add(1, Ordering::SeqCst); + let request = read_http_head(&mut stream).await; + assert!(request.starts_with(b"GET /session/dial?instance=")); + stream + .write_all( + b"HTTP/1.1 503 Service Unavailable\r\nContent-Length: 0\r\nConnection: close\r\n\r\n", + ) + .await + .unwrap(); + }); + let mut credential = peer.credential(); + credential.endpoints = vec![EndpointAddr { + host: "127.0.0.1".into(), + port: lan_port, + }]; + credential.relay_origin = Some(relay_origin); + credential.device_token = Some(test_jwt(i64::MAX / 2)); + let root = temp.path().to_path_buf(); + let owner = + tokio::spawn( + async move { start_private_link_owner(&root, credential, "stream").await }, + ); + while lan_accepted.load(Ordering::SeqCst) == 0 { + tokio::task::yield_now().await; + } + tokio::time::advance(LAN_CARRIER_TIMEOUT).await; + for _ in 0..100 { + if relay_accepted.load(Ordering::SeqCst) != 0 { + break; + } + tokio::task::yield_now().await; + } + assert_eq!( + relay_accepted.load(Ordering::SeqCst), + 1, + "relay connection count after stalled LAN deadline" + ); + relay_task.await.unwrap(); + let owner = owner.await.unwrap().unwrap(); + assert_eq!(lan_accepted.load(Ordering::SeqCst), 1); + owner.shutdown().await.unwrap(); + stalled.abort(); + peer.shutdown().await; + } + + #[tokio::test(start_paused = true)] + async fn delayed_lan_remains_preferred_over_relay() { + let temp = tempfile::tempdir().unwrap(); + let tls_gate = Arc::new(Notify::new()); + let peer = PrivateLinkPeer::start_with_delayed_tls(tls_gate.clone()).await; + peer.enqueue_response( + 200, + serde_json::json!({ + "key":"K", "name":"stream", "prefix":"prefix", + "ingest_url":"/app/observer/ingest", "protocol_version":2 + }) + .to_string(), + ); + let relay = tokio::net::TcpListener::bind(("127.0.0.1", 0)) + .await + .unwrap(); + let relay_origin = format!("http://{}", relay.local_addr().unwrap()); + let relay_accepted = Arc::new(AtomicUsize::new(0)); + let relay_count = relay_accepted.clone(); + let relay_task = tokio::spawn(async move { + if relay.accept().await.is_ok() { + relay_count.fetch_add(1, Ordering::SeqCst); + } + }); + let mut credential = peer.credential(); + credential.relay_origin = Some(relay_origin); + credential.device_token = Some(test_jwt(i64::MAX / 2)); + let root = temp.path().to_path_buf(); + let owner = + tokio::spawn( + async move { start_private_link_owner(&root, credential, "stream").await }, + ); + while peer.accepted_carriers() == 0 { + tokio::task::yield_now().await; + } + tokio::time::advance(Duration::from_secs(1)).await; + tls_gate.notify_one(); + let owner = owner.await.unwrap().unwrap(); + assert_eq!(peer.accepted_carriers(), 1); + assert_eq!(relay_accepted.load(Ordering::SeqCst), 0); + owner.shutdown().await.unwrap(); + relay_task.abort(); + peer.shutdown().await; + } + + #[tokio::test] + async fn no_relay_origin_credential_keeps_lan_path() { + let temp = tempfile::tempdir().unwrap(); + let peer = PrivateLinkPeer::start().await; + peer.enqueue_response( + 200, + serde_json::json!({ + "key":"K", "name":"stream", "prefix":"prefix", + "ingest_url":"/app/observer/ingest", "protocol_version":2 + }) + .to_string(), + ); + let owner = start_private_link_owner(temp.path(), peer.credential(), "stream") + .await + .unwrap(); + assert_eq!(peer.accepted_carriers(), 1); + owner.shutdown().await.unwrap(); + peer.shutdown().await; + } + + #[tokio::test] + async fn owner_start_guard_refusal_is_terminal() { + let temp = tempfile::tempdir().unwrap(); + let peer = PrivateLinkPeer::start().await; + peer.enqueue_response( + 403, + serde_json::json!({"reason_code":"local_request_only"}).to_string(), + ); + let result = start_private_link_owner(temp.path(), peer.credential(), "stream").await; + assert!(matches!( + result, + Err(PrivateStateError::RegistrationInvalid) + )); + assert_eq!(peer.requests().len(), 1); + peer.shutdown().await; + } + + #[tokio::test] + async fn relay_origin_without_device_token_uses_lan_only() { + let temp = tempfile::tempdir().unwrap(); + let peer = PrivateLinkPeer::start().await; + peer.enqueue_response( + 200, + serde_json::json!({ + "key":"K", "name":"stream", "prefix":"prefix", + "ingest_url":"/app/observer/ingest", "protocol_version":2 + }) + .to_string(), + ); + let relay = tokio::net::TcpListener::bind(("127.0.0.1", 0)) + .await + .unwrap(); + let mut credential = peer.credential(); + credential.relay_origin = Some(format!("http://{}", relay.local_addr().unwrap())); + credential.device_token = None; + let owner = start_private_link_owner(temp.path(), credential, "stream") + .await + .unwrap(); + assert_eq!(peer.accepted_carriers(), 1); + assert!( + tokio::time::timeout(Duration::ZERO, relay.accept()) + .await + .is_err() + ); + owner.shutdown().await.unwrap(); + peer.shutdown().await; + } + + #[tokio::test(start_paused = true)] + async fn relay_failure_preserves_lan_timeout_error() { + let peer = PrivateLinkPeer::start().await; + let lan = tokio::net::TcpListener::bind(("127.0.0.1", 0)) + .await + .unwrap(); + let lan_port = lan.local_addr().unwrap().port(); + let lan_accepted = Arc::new(AtomicUsize::new(0)); + let lan_count = lan_accepted.clone(); + let stalled = tokio::spawn(async move { + let (_stream, _) = lan.accept().await.unwrap(); + lan_count.fetch_add(1, Ordering::SeqCst); + std::future::pending::<()>().await; + }); + let relay = tokio::net::TcpListener::bind(("127.0.0.1", 0)) + .await + .unwrap(); + let relay_origin = format!("http://{}", relay.local_addr().unwrap()); + let relay_accepted = Arc::new(AtomicUsize::new(0)); + let relay_count = relay_accepted.clone(); + let relay_task = tokio::spawn(async move { + let (mut stream, _) = relay.accept().await.unwrap(); + relay_count.fetch_add(1, Ordering::SeqCst); + let _ = read_http_head(&mut stream).await; + stream + .write_all( + b"HTTP/1.1 503 Service Unavailable\r\nContent-Length: 0\r\nConnection: close\r\n\r\n", + ) + .await + .unwrap(); + }); + let mut lan_credential = peer.credential(); + lan_credential.endpoints = vec![EndpointAddr { + host: "127.0.0.1".into(), + port: lan_port, + }]; + lan_credential.relay_origin = None; + lan_credential.device_token = None; + let mut relay_credential = peer.credential(); + relay_credential.endpoints.clear(); + relay_credential.relay_origin = Some(relay_origin); + relay_credential.device_token = Some(test_jwt(i64::MAX / 2)); + let opener = Arc::new(PrivateLinkOpener::new( + Some(TransportClient::new(lan_credential, None).unwrap()), + Some(TransportClient::new_relay_only(relay_credential, None).unwrap()), + Arc::new(AtomicBool::new(false)), + LinkFacts::default(), + )); + let dial = tokio::spawn({ + let opener = opener.clone(); + async move { opener.dial_carrier().await } + }); + while lan_accepted.load(Ordering::SeqCst) == 0 { + tokio::task::yield_now().await; + } + tokio::time::advance(LAN_CARRIER_TIMEOUT).await; + relay_task.await.unwrap(); + let error = match dial.await.unwrap() { + Ok(_) => panic!("both carrier legs unexpectedly succeeded"), + Err(error) => error, + }; + match error { + TransportError::Io(error) => { + assert_eq!(error.kind(), io::ErrorKind::TimedOut); + assert_eq!(error.to_string(), "LAN carrier dial timed out"); + } + other => panic!("expected LAN timeout, got {other}"), + } + assert_eq!(relay_accepted.load(Ordering::SeqCst), 1); + stalled.abort(); + peer.shutdown().await; + } + impl Pairer for FakePairer { fn pair<'a>( &'a self, diff --git a/crates/solstone-linux/src/private_link_test_peer.rs b/crates/solstone-linux/src/private_link_test_peer.rs index 4a3802a..af3ae2f 100644 --- a/crates/solstone-linux/src/private_link_test_peer.rs +++ b/crates/solstone-linux/src/private_link_test_peer.rs @@ -79,6 +79,14 @@ pub(crate) struct PrivateLinkPeer { impl PrivateLinkPeer { pub(crate) async fn start() -> Self { + Self::start_with_tls_gate(None).await + } + + pub(crate) async fn start_with_delayed_tls(gate: Arc) -> Self { + Self::start_with_tls_gate(Some(gate)).await + } + + async fn start_with_tls_gate(tls_gate: Option>) -> Self { let listener = TcpListener::bind(("127.0.0.1", 0)).await.unwrap(); let (credential, acceptor) = credential_and_acceptor(listener.local_addr().unwrap().port()); let state = PeerState { @@ -96,6 +104,9 @@ impl PrivateLinkPeer { let task = tokio::spawn(async move { while let Ok((stream, _)) = listener.accept().await { task_state.accepted.fetch_add(1, Ordering::SeqCst); + if let Some(gate) = &tls_gate { + gate.notified().await; + } let Ok(tls) = acceptor.accept(stream).await else { continue; }; diff --git a/crates/solstone-linux/src/sync.rs b/crates/solstone-linux/src/sync.rs index 621c3ba..30edf7b 100644 --- a/crates/solstone-linux/src/sync.rs +++ b/crates/solstone-linux/src/sync.rs @@ -372,7 +372,12 @@ impl SyncWorker { } break; } - if !self.client.is_registered() { + // ensure_registered deliberately ignores revocation; this production call site + // stops the retry loop for a terminally revoked credential. + if !self.client.is_registered() + && (self.client.is_revoked() + || !self.client.ensure_registered(&mut self.config).await) + { if !self.registration_refused { tracing::warn!("Sync refused: observer is not registered"); self.registration_refused = true; @@ -4132,6 +4137,230 @@ mod tests { assert!(server.requests().len() >= 5); } + // AC7: an unregistered worker retries once and falls through to sync on success. + #[tokio::test] + async fn unregistered_worker_retries_registration_and_syncs_on_success() { + let temp = tempfile::tempdir().unwrap(); + let peer = PrivateLinkPeer::start().await; + peer.enqueue_response(200, registration("REGISTERED-KEY", "desktop").to_string()); + peer.enqueue_response(200, json!({"items":[],"total":0}).to_string()); + let config = Config { + stream: "desktop".to_owned(), + base_dir: temp.path().join("data"), + config_dir: temp.path().join("config"), + ..Config::default() + }; + let session = start_private_link_session(&config.config_dir, peer.credential(), "desktop") + .await + .unwrap(); + let client = Arc::new(UploadClient::new( + &config, + session.capability("/app/observer/ingest".to_owned()), + "host", + "linux", + "test", + Arc::new(FixedClock { + wall: 1_800_000_000.0, + mono: 100.0, + }), + )); + let notify = Arc::new(Notify::new()); + let running = Arc::new(AtomicBool::new(true)); + let mut worker = SyncWorker::new( + config, + client.clone(), + Arc::new(FixedClock { + wall: 1_800_000_000.0, + mono: 100.0, + }), + SyncControl { + notify: notify.clone(), + pending_trigger: Arc::new(AtomicBool::new(false)), + running: running.clone(), + }, + Arc::new(Mutex::new(SyncFacts::default())), + Arc::new(AtomicU8::new(0)), + ); + let task = tokio::spawn(async move { worker.run().await }); + notify.notify_one(); + let _ = tokio::time::timeout(Duration::from_secs(1), peer.wait_for_requests(2)).await; + assert_eq!( + peer.requests().len(), + 2, + "registration and same-tick sync requests: registered={}, paths={:?}", + client.is_registered(), + peer.requests() + .iter() + .map(|request| &request.path) + .collect::>() + ); + assert!(client.is_registered()); + running.store(false, Ordering::Release); + notify.notify_one(); + task.await.unwrap(); + session.shutdown().await.unwrap(); + peer.shutdown().await; + } + + #[tokio::test] + async fn revoked_worker_never_attempts_registration() { + let temp = tempfile::tempdir().unwrap(); + let peer = PrivateLinkPeer::start().await; + let config = Config { + stream: "desktop".to_owned(), + base_dir: temp.path().join("data"), + config_dir: temp.path().join("config"), + ..Config::default() + }; + let session = start_private_link_session(&config.config_dir, peer.credential(), "desktop") + .await + .unwrap(); + let client = Arc::new(UploadClient::new( + &config, + session.capability("/app/observer/ingest".to_owned()), + "host", + "linux", + "test", + Arc::new(FixedClock { + wall: 1_800_000_000.0, + mono: 0.0, + }), + )); + client.revoke_for_test(); + let notify = Arc::new(Notify::new()); + let running = Arc::new(AtomicBool::new(true)); + let mut worker = SyncWorker::new( + config, + client, + Arc::new(FixedClock { + wall: 1_800_000_000.0, + mono: 0.0, + }), + SyncControl { + notify: notify.clone(), + pending_trigger: Arc::new(AtomicBool::new(false)), + running: running.clone(), + }, + Arc::new(Mutex::new(SyncFacts::default())), + Arc::new(AtomicU8::new(0)), + ); + let output = Arc::new(Mutex::new(Vec::new())); + let writer = Buffer(output.clone()); + let subscriber = tracing_subscriber::fmt() + .without_time() + .with_ansi(false) + .with_writer(move || writer.clone()) + .finish(); + let task = tokio::spawn(async move { worker.run().await }.with_subscriber(subscriber)); + notify.notify_one(); + for _ in 0..100 { + tokio::task::yield_now().await; + } + running.store(false, Ordering::Release); + notify.notify_one(); + task.await.unwrap(); + assert_eq!(peer.requests().len(), 0, "revoked registration requests"); + let output = String::from_utf8(output.lock().unwrap().clone()).unwrap(); + assert_eq!(output.matches("Sync refused").count(), 1); + session.shutdown().await.unwrap(); + peer.shutdown().await; + } + + #[tokio::test(start_paused = true)] + async fn unregistered_worker_registration_retry_respects_cooldown() { + let temp = tempfile::tempdir().unwrap(); + let peer = PrivateLinkPeer::start().await; + peer.enqueue_response(503, Vec::new()); + peer.enqueue_response(200, registration("REGISTERED-KEY", "desktop").to_string()); + peer.enqueue_response(200, json!({"items":[],"total":0}).to_string()); + let config = Config { + stream: "desktop".to_owned(), + base_dir: temp.path().join("data"), + config_dir: temp.path().join("config"), + ..Config::default() + }; + let session = start_private_link_session(&config.config_dir, peer.credential(), "desktop") + .await + .unwrap(); + 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()), + "host", + "linux", + "test", + clock.clone(), + )); + let notify = Arc::new(Notify::new()); + let running = Arc::new(AtomicBool::new(true)); + let mut worker = SyncWorker::new( + config, + client, + clock.clone(), + SyncControl { + notify: notify.clone(), + pending_trigger: Arc::new(AtomicBool::new(false)), + running: running.clone(), + }, + Arc::new(Mutex::new(SyncFacts::default())), + Arc::new(AtomicU8::new(0)), + ); + let task = tokio::spawn(async move { worker.run().await }); + notify.notify_one(); + for _ in 0..10_000 { + if !peer.requests().is_empty() { + break; + } + tokio::task::yield_now().await; + } + assert_eq!(peer.requests().len(), 1, "first registration request"); + tokio::time::advance(Duration::from_millis(1)).await; + + clock.set_mono(60.0); + for _ in 0..3 { + notify.notify_one(); + tokio::task::yield_now().await; + } + tokio::time::advance(Duration::from_secs(60)).await; + assert_eq!( + peer.requests() + .iter() + .filter(|request| request.path == "/app/observer/register") + .count(), + 1, + "registration requests during cooldown" + ); + + tokio::time::advance(Duration::from_secs(240)).await; + clock.set_mono(300.0); + notify.notify_one(); + for _ in 0..10_000 { + if peer + .requests() + .iter() + .filter(|request| request.path == "/app/observer/register") + .count() + >= 2 + { + break; + } + tokio::task::yield_now().await; + } + assert_eq!( + peer.requests() + .iter() + .filter(|request| request.path == "/app/observer/register") + .count(), + 2, + "registration requests after cooldown" + ); + running.store(false, Ordering::Release); + notify.notify_one(); + task.await.unwrap(); + session.shutdown().await.unwrap(); + peer.shutdown().await; + } + // AC: an unregistered worker logs refusal once, stays idle, and performs no HTTP. #[tokio::test] async fn unregistered_worker_refuses_once_without_http() { diff --git a/crates/solstone-linux/src/upload.rs b/crates/solstone-linux/src/upload.rs index 534a422..2ccb55d 100644 --- a/crates/solstone-linux/src/upload.rs +++ b/crates/solstone-linux/src/upload.rs @@ -191,6 +191,11 @@ impl UploadClient { self.inner.revoked.load(Ordering::Acquire) } + #[cfg(test)] + pub(crate) fn revoke_for_test(&self) { + self.inner.revoked.store(true, Ordering::Release); + } + pub fn is_registered(&self) -> bool { if let Some(capability) = self.inner.capability() { return capability.is_registered(); -- 2.51.2