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();