diff --git a/crates/solstone-linux/src/cli.rs b/crates/solstone-linux/src/cli.rs index d2ce2e0..06d8d33 100644 --- a/crates/solstone-linux/src/cli.rs +++ b/crates/solstone-linux/src/cli.rs @@ -883,7 +883,7 @@ mod tests { } #[test] - fn setup_pairs_over_private_link_and_never_uses_legacy_registration() { + fn setup_source_policy_uses_private_link_and_excludes_legacy_registration() { let source = include_str!("cli.rs"); assert!(source.contains("setup_with_stream(")); assert!(!source.contains(&["UploadClient", "::new("].concat())); diff --git a/crates/solstone-linux/src/linked_authority_policy_tests.rs b/crates/solstone-linux/src/linked_authority_policy_tests.rs index 01f1ce3..502a0da 100644 --- a/crates/solstone-linux/src/linked_authority_policy_tests.rs +++ b/crates/solstone-linux/src/linked_authority_policy_tests.rs @@ -77,10 +77,49 @@ fn marked_items(lines: &[&str], marker: &str) -> Vec { marked } -fn constructs_origin_relative_url(line: &str) -> bool { - let journal_route = line.contains("/app/") || line.contains("/api/"); - journal_route - && (line.contains("format!(") || line.contains(".join(") || line.contains("push_str(")) +fn origin_relative_constructions(source: &str) -> Vec { + let bytes = source.as_bytes(); + let mut findings = Vec::new(); + for needle in ["format!(", ".join(", "push_str("] { + let mut offset = 0; + while let Some(relative) = source[offset..].find(needle) { + let start = offset + relative; + let mut depth = 0_i32; + let mut quoted = false; + let mut escaped = false; + let mut end = start; + for (index, byte) in bytes[start..].iter().copied().enumerate() { + end = start + index + 1; + if quoted { + if escaped { + escaped = false; + } else if byte == b'\\' { + escaped = true; + } else if byte == b'"' { + quoted = false; + } + continue; + } + match byte { + b'"' => quoted = true, + b'(' => depth += 1, + b')' => { + depth -= 1; + if depth == 0 { + break; + } + } + _ => {} + } + } + let construction = &source[start..end]; + if construction.contains("/app/") || construction.contains("/api/") { + findings.push(start); + } + offset = end.max(start + needle.len()); + } + } + findings } #[test] @@ -99,18 +138,19 @@ fn active_production_paths_have_no_legacy_direct_authority() { let lines = production_prefix(&source).lines().collect::>(); let l3_marked = marked_items(&lines, L3_MARKER); let test_only = marked_items(&lines, "#[cfg(test)]"); - for (index, line) in lines.iter().enumerate() { + for byte in origin_relative_constructions(production_prefix(&source)) { + let line = production_prefix(&source)[..byte] + .lines() + .count() + .saturating_sub(1); assert!( - !constructs_origin_relative_url(line) - || path - .file_name() - .is_some_and(|name| name == "private_link.rs") - || l3_marked[index] - || test_only[index], + l3_marked[line] || test_only[line], "{}:{} constructs an origin-relative Journal URL", path.display(), - index + 1 + line + 1 ); + } + for (index, line) in lines.iter().enumerate() { for needle in forbidden { assert!( !line.contains(needle) || l3_marked[index] || test_only[index], @@ -131,6 +171,17 @@ fn active_production_paths_have_no_legacy_direct_authority() { } } +#[test] +fn origin_construction_scan_handles_split_invocations_and_routes() { + for source in [ + "format!(\n\"{origin}/app/observer\"\n)", + "format!(\n\"{}{}\",\norigin,\n\"/app/observer\"\n)", + "origin\n.join(\n\"/app/observer\"\n)", + ] { + assert_eq!(origin_relative_constructions(source).len(), 1); + } +} + #[test] fn l3_authority_is_confined_to_explicit_unreachable_surfaces() { let root = source_root(); diff --git a/crates/solstone-linux/src/private_link.rs b/crates/solstone-linux/src/private_link.rs index 8da0bf8..42a667e 100644 --- a/crates/solstone-linux/src/private_link.rs +++ b/crates/solstone-linux/src/private_link.rs @@ -621,7 +621,12 @@ impl LinkFacts { LinkFact::ConfigSanitationFailed => state.config_sanitation_failed = true, LinkFact::ListenerReady => state.listener_ready = true, LinkFact::CarrierProven => state.carrier_proven = true, - LinkFact::ObserverRegistered => state.observer_registered = true, + LinkFact::ObserverRegistered => { + state.observer_registered = true; + if !state.token_persistence_failure { + state.transport_unavailable = false; + } + } LinkFact::TransportUnavailable => state.transport_unavailable = true, LinkFact::TerminalRevocation => state.terminal_revocation = true, LinkFact::TokenPersistenceFailure => state.token_persistence_failure = true, @@ -1116,6 +1121,15 @@ impl PrivateLinkOwner { self.session.shutdown().await } + #[cfg(test)] + pub(crate) async fn shutdown_with_join_probe( + self, + joined: Arc, + release: Arc, + ) -> Result<(), PrivateStateError> { + self.session.shutdown_with_join_probe(joined, release).await + } + #[cfg(test)] fn loopback_addr(&self) -> std::net::SocketAddr { format!( @@ -1173,18 +1187,12 @@ async fn finish_owner_start( if session.opener.generation() == 0 { match capability.report_unauthorized(0).await { RepairOutcome::Repaired { .. } | RepairOutcome::AlreadySuperseded { .. } => {} - RepairOutcome::PersistenceFailed => { - return Err(PrivateStateError::Io { - operation: PrivateIoOperation::Persist, - source: io::Error::other("registration persistence failed"), - }); + RepairOutcome::PersistenceFailed | RepairOutcome::TransportUnavailable => { + capability.facts().publish(LinkFact::TransportUnavailable); } RepairOutcome::GuardRefused { .. } | RepairOutcome::InvalidRegistration => { return Err(PrivateStateError::RegistrationInvalid); } - RepairOutcome::TransportUnavailable => { - return Err(PrivateStateError::BridgeUnavailable); - } } } Ok(PrivateLinkOwner { @@ -1310,7 +1318,20 @@ impl PrivateLinkSession { } pub(crate) async fn shutdown(self) -> Result<(), PrivateStateError> { + self.shutdown_inner(None).await + } + + async fn shutdown_inner( + self, + #[cfg(test)] join_probe: Option<(Arc, Arc)>, + #[cfg(not(test))] _join_probe: Option<()>, + ) -> Result<(), PrivateStateError> { let status = self.handle.shutdown_and_wait().await; + #[cfg(test)] + if let Some((joined, release)) = join_probe { + joined.notify_one(); + release.notified().await; + } if status.listener_active || status.active_requests != 0 { return Err(PrivateStateError::ShutdownFailed); } @@ -1320,6 +1341,15 @@ impl PrivateLinkSession { Ok(()) } + #[cfg(test)] + async fn shutdown_with_join_probe( + self, + joined: Arc, + release: Arc, + ) -> Result<(), PrivateStateError> { + self.shutdown_inner(Some((joined, release))).await + } + #[cfg(test)] fn publish_observer(&self, observer: &ObserverState) -> Result { persist_and_publish_observer( @@ -1539,8 +1569,12 @@ async fn start_private_link_session_inner( transport_unavailable.clone(), facts.clone(), ); - let transport = TransportClient::new(credential, Some(hook)) - .map_err(|_| PrivateStateError::BridgeUnavailable)?; + let transport = if credential.endpoints.is_empty() { + TransportClient::new_relay_only(credential, Some(hook)) + } else { + TransportClient::new(credential, Some(hook)) + } + .map_err(|_| PrivateStateError::BridgeUnavailable)?; let opener = Arc::new(PrivateLinkOpener::new( transport, transport_unavailable, @@ -1701,6 +1735,60 @@ mod tests { response } + fn base64url_no_pad(input: &[u8]) -> String { + const TABLE: &[u8; 64] = + b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-_"; + let mut output = String::new(); + let mut index = 0; + while index + 3 <= input.len() { + let chunk = ((input[index] as u32) << 16) + | ((input[index + 1] as u32) << 8) + | input[index + 2] as u32; + output.push(TABLE[((chunk >> 18) & 0x3f) as usize] as char); + output.push(TABLE[((chunk >> 12) & 0x3f) as usize] as char); + output.push(TABLE[((chunk >> 6) & 0x3f) as usize] as char); + output.push(TABLE[(chunk & 0x3f) as usize] as char); + index += 3; + } + match input.len() - index { + 1 => { + let chunk = (input[index] as u32) << 16; + output.push(TABLE[((chunk >> 18) & 0x3f) as usize] as char); + output.push(TABLE[((chunk >> 12) & 0x3f) as usize] as char); + } + 2 => { + let chunk = ((input[index] as u32) << 16) | ((input[index + 1] as u32) << 8); + output.push(TABLE[((chunk >> 18) & 0x3f) as usize] as char); + output.push(TABLE[((chunk >> 12) & 0x3f) as usize] as char); + output.push(TABLE[((chunk >> 6) & 0x3f) as usize] as char); + } + _ => {} + } + output + } + + fn test_jwt(exp: i64) -> String { + let claims = format!(r#"{{"iat":{},"exp":{exp}}}"#, exp - 3600); + format!( + "{}.{}.sig", + base64url_no_pad(b"{}"), + base64url_no_pad(claims.as_bytes()) + ) + } + + async fn read_http_head(stream: &mut tokio::net::TcpStream) -> Vec { + let mut request = Vec::new(); + let mut buffer = [0_u8; 1024]; + while !request.windows(4).any(|window| window == b"\r\n\r\n") { + let count = stream.read(&mut buffer).await.unwrap(); + if count == 0 { + break; + } + request.extend_from_slice(&buffer[..count]); + } + request + } + fn credential() -> Credential { Credential { client_key_pem: "client-key".into(), @@ -3096,21 +3184,62 @@ mod tests { #[tokio::test] async fn sanitation_precedes_pairer_bridge_carrier_and_peer() { let temp = tempfile::tempdir().unwrap(); + let traps = LegacyNetworkTraps::bind(); + let peer = PrivateLinkPeer::start().await; fs::write( temp.path().join("config.json"), - r#"{"server_url":"http://127.0.0.1:9","key":"secret","stream":"stream"}"#, + format!( + r#"{{"server_url":"{}","key":"secret","stream":"stream"}}"#, + traps.configured_origin() + ), ) .unwrap(); let calls = Arc::new(AtomicUsize::new(0)); let pairer = SanitizedConfigPairer { config_path: temp.path().join("config.json"), calls: calls.clone(), - result: credential(), + result: peer.credential(), }; setup_with_pairer(&pairer, temp.path(), "device", Cursor::new(b"pair")) .await .unwrap(); assert_eq!(calls.load(Ordering::SeqCst), 1); + + let config_path = temp.path().join("config.json"); + let mut reacquired: serde_json::Value = + serde_json::from_slice(&fs::read(&config_path).unwrap()).unwrap(); + reacquired["server_url"] = serde_json::json!(traps.configured_origin()); + reacquired["key"] = serde_json::json!("reintroduced"); + fs::write(&config_path, serde_json::to_vec(&reacquired).unwrap()).unwrap(); + + peer.enqueue_response(200, br#"{"items":[],"total":0}"#.to_vec()); + let session = start_private_link_session(temp.path(), peer.credential(), "stream") + .await + .unwrap(); + let sanitized: serde_json::Value = + serde_json::from_slice(&fs::read(&config_path).unwrap()).unwrap(); + assert!(sanitized.get("server_url").is_none()); + assert!(sanitized.get("key").is_none()); + publish_observer_registration( + &session, + &ObserverState { + credential_instance_id: peer.credential().instance_id, + ..observer("/ingest") + }, + ) + .unwrap(); + assert!(matches!( + session + .capability("/ingest".to_owned()) + .list_day("20260101") + .await, + LinkOutcome::Success { .. } + )); + assert_eq!(peer.accepted_carriers(), 1); + assert_eq!(peer.requests().len(), 1); + traps.assert_zero_connections(); + session.shutdown().await.unwrap(); + peer.shutdown().await; } struct LegacyNetworkTraps { @@ -3248,7 +3377,7 @@ mod tests { } #[tokio::test] - async fn restart_mismatch_retries_without_manual_cleanup() { + async fn restart_mismatch_starts_unregistered_without_deleting_observer() { let temp = tempfile::tempdir().unwrap(); let peer = PrivateLinkPeer::start().await; let credential = peer.credential(); @@ -3316,7 +3445,7 @@ mod tests { } #[tokio::test] - async fn failed_repair_preserves_usable_prior_credential() { + async fn failed_two_file_publish_preserves_usable_prior_credential() { let peer = PrivateLinkPeer::start().await; peer.enqueue_response(200, b"{}".to_vec()); let (_temp, session) = start_peer_session(&peer).await; @@ -3348,7 +3477,7 @@ mod tests { } #[tokio::test] - async fn successful_repair_invalidates_prior_instance_observer() { + async fn credential_instance_mismatch_rejects_prior_observer() { let temp = tempfile::tempdir().unwrap(); persist_observer(temp.path(), &observer("/ingest")).unwrap(); let peer = PrivateLinkPeer::start().await; @@ -3363,7 +3492,7 @@ mod tests { } #[tokio::test] - async fn repaired_instance_registers_before_first_data_request() { + async fn published_registration_authenticates_first_data_request() { let peer = PrivateLinkPeer::start().await; peer.enqueue_response(200, b"{}".to_vec()); let temp = tempfile::tempdir().unwrap(); @@ -3550,6 +3679,69 @@ mod tests { peer.shutdown().await; } + #[tokio::test] + async fn refreshed_relay_token_is_durable_before_owner_shutdown_returns() { + let temp = tempfile::tempdir().unwrap(); + let peer = PrivateLinkPeer::start().await; + let listener = tokio::net::TcpListener::bind(("127.0.0.1", 0)) + .await + .unwrap(); + let relay_origin = format!("http://{}", listener.local_addr().unwrap()); + let refreshed = test_jwt( + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_secs() as i64 + + 3600, + ); + let relay_refreshed = refreshed.clone(); + let relay = tokio::spawn(async move { + let (mut refresh, _) = listener.accept().await.unwrap(); + let request = read_http_head(&mut refresh).await; + assert!( + request.starts_with(b"POST /token/refresh HTTP/1.1"), + "{}", + String::from_utf8_lossy(&request) + ); + let body = serde_json::json!({"device_token": relay_refreshed}).to_string(); + refresh + .write_all( + format!( + "HTTP/1.1 200 OK\r\nContent-Length: {}\r\nContent-Type: application/json\r\nConnection: close\r\n\r\n{body}", + body.len() + ) + .as_bytes(), + ) + .await + .unwrap(); + let (mut dial, _) = listener.accept().await.unwrap(); + let _ = read_http_head(&mut dial).await; + dial.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.clear(); + credential.relay_origin = Some(relay_origin); + credential.device_token = Some(test_jwt(1)); + credential.device_token_expires_at = Some(1); + let owner = start_private_link_owner(temp.path(), credential, "stream") + .await + .unwrap(); + relay.await.unwrap(); + let persisted = load_credential(temp.path()).unwrap().unwrap(); + assert_eq!(persisted.device_token.as_deref(), Some(refreshed.as_str())); + owner.shutdown().await.unwrap(); + let persisted_after_shutdown = load_credential(temp.path()).unwrap().unwrap(); + assert_eq!( + persisted_after_shutdown.device_token.as_deref(), + Some(refreshed.as_str()) + ); + peer.shutdown().await; + } + #[tokio::test] async fn executed_session_surfaces_do_not_disclose_secrets() { let peer = PrivateLinkPeer::start().await; diff --git a/crates/solstone-linux/src/run.rs b/crates/solstone-linux/src/run.rs index 205e465..0e0b751 100644 --- a/crates/solstone-linux/src/run.rs +++ b/crates/solstone-linux/src/run.rs @@ -1060,11 +1060,13 @@ mod tests { async fn drive_real_link_start( temp: &tempfile::TempDir, transport_enabled: bool, - ) -> (Result, LinkFactState) { + ) -> ( + Result, + LinkFactState, + Arc, + ) { let legacy_origin = MockServer::new(Vec::new()).await; let default_listener = OpportunisticDefaultListenerTrap::bind(); - let pending = temp.path().join("pending.segment"); - std::fs::write(&pending, b"pending").unwrap(); let config = Config { config_dir: temp.path().to_path_buf(), stream: "stream".to_owned(), @@ -1089,10 +1091,9 @@ mod tests { )); assert_real_observer_ticks_advance(); let result = start.await.unwrap(); - assert!(pending.exists()); assert!(legacy_origin.requests().is_empty()); default_listener.assert_zero_connections(); - (result, upload.link_fact_state().unwrap()) + (result, upload.link_fact_state().unwrap(), upload) } fn assert_real_observer_ticks_advance() { @@ -1108,7 +1109,7 @@ mod tests { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn missing_credentials_observer_ticks_without_transport() { let temp = tempfile::tempdir().unwrap(); - let (result, facts) = drive_real_link_start(&temp, true).await; + let (result, facts, _upload) = drive_real_link_start(&temp, true).await; assert!(result.is_err()); assert!(facts.pairing_required); } @@ -1117,7 +1118,7 @@ mod tests { async fn malformed_credentials_observer_ticks_without_transport() { let temp = tempfile::tempdir().unwrap(); std::fs::write(temp.path().join(CREDENTIALS_FILENAME), b"{").unwrap(); - let (result, facts) = drive_real_link_start(&temp, true).await; + let (result, facts, _upload) = drive_real_link_start(&temp, true).await; assert!(result.is_err()); assert!(facts.private_state_invalid); } @@ -1136,7 +1137,7 @@ mod tests { }) .to_string(), ); - let (result, facts) = drive_real_link_start(&temp, true).await; + let (result, facts, _upload) = drive_real_link_start(&temp, true).await; let owner = result.unwrap(); assert!(facts.observer_registered); let registration: serde_json::Value = @@ -1154,9 +1155,10 @@ mod tests { let peer = PrivateLinkPeer::start().await; persist_credential(temp.path(), &peer.credential()).unwrap(); peer.shutdown().await; - let (result, facts) = drive_real_link_start(&temp, true).await; - assert!(result.is_err()); + let (result, facts, _upload) = drive_real_link_start(&temp, true).await; + let owner = result.expect("retryable carrier failure must retain the owner"); assert!(facts.transport_unavailable); + owner.shutdown().await.unwrap(); } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] @@ -1165,10 +1167,23 @@ mod tests { let peer = PrivateLinkPeer::start().await; persist_credential(temp.path(), &peer.credential()).unwrap(); peer.enqueue_response(503, Vec::new()); - let (result, facts) = drive_real_link_start(&temp, true).await; - assert!(result.is_err()); + let (result, facts, upload) = drive_real_link_start(&temp, true).await; + let owner = result.expect("retryable registration failure must retain the owner"); assert!(facts.transport_unavailable); assert_eq!(peer.requests().len(), 1); + enqueue_registration(&peer); + let mut config = Config { + config_dir: temp.path().to_path_buf(), + stream: "stream".to_owned(), + ..Config::default() + }; + assert!(upload.ensure_registered(&mut config).await); + assert_eq!(peer.requests().len(), 2); + assert!(upload.is_registered()); + let recovered = upload.link_fact_state().unwrap(); + assert!(recovered.observer_registered); + assert!(!recovered.transport_unavailable); + owner.shutdown().await.unwrap(); peer.shutdown().await; } @@ -1214,17 +1229,8 @@ mod tests { } #[tokio::test] - async fn concurrent_initial_demand_performs_one_registration() { + async fn concurrent_ensure_registered_waiters_share_one_attempt_and_result() { assert_concurrent_initial_registration(200, true).await; - } - - #[tokio::test] - async fn initial_registration_waiters_share_one_success() { - assert_concurrent_initial_registration(200, true).await; - } - - #[tokio::test] - async fn initial_registration_waiters_share_one_unavailable_result() { assert_concurrent_initial_registration(503, false).await; } @@ -1242,6 +1248,8 @@ mod tests { }) .to_string(), ); + let response_gate = Arc::new(AtomicBool::new(false)); + peer.gate_next_response_nonblocking(response_gate.clone()); let session = start_private_link_session(temp.path(), peer.credential(), "stream") .await .unwrap(); @@ -1267,6 +1275,15 @@ mod tests { client.ensure_registered(&mut config).await })); } + peer.wait_for_requests(1).await; + assert!(demands.iter().all(|demand| !demand.is_finished())); + assert!( + peer.requests() + .iter() + .all(|request| request.path == "/app/observer/register") + ); + response_gate.store(true, Ordering::Release); + peer.notify_response_gates(); for demand in demands { assert_eq!(demand.await.unwrap(), expected); } @@ -1287,7 +1304,18 @@ mod tests { PrivateStateLock::acquire(temp.path()), Err(PrivateStateError::LockContended) )); - owner.shutdown().await.unwrap(); + let joined = Arc::new(tokio::sync::Notify::new()); + let release = Arc::new(tokio::sync::Notify::new()); + let shutdown = + tokio::spawn(owner.shutdown_with_join_probe(joined.clone(), release.clone())); + joined.notified().await; + assert!(!shutdown.is_finished()); + assert!(matches!( + PrivateStateLock::acquire(temp.path()), + Err(PrivateStateError::LockContended) + )); + release.notify_one(); + shutdown.await.unwrap().unwrap(); let lock = PrivateStateLock::acquire(temp.path()).unwrap(); drop(lock); peer.shutdown().await; @@ -1305,7 +1333,17 @@ mod tests { PrivateStateLock::acquire(temp.path()), Err(PrivateStateError::LockContended) )); - owner.shutdown().await.unwrap(); + let joined = Arc::new(tokio::sync::Notify::new()); + let release = Arc::new(tokio::sync::Notify::new()); + let shutdown = + tokio::spawn(owner.shutdown_with_join_probe(joined.clone(), release.clone())); + joined.notified().await; + assert!(matches!( + PrivateStateLock::acquire(temp.path()), + Err(PrivateStateError::LockContended) + )); + release.notify_one(); + shutdown.await.unwrap().unwrap(); assert!(PrivateStateLock::acquire(temp.path()).is_ok()); peer.shutdown().await; } @@ -1350,9 +1388,9 @@ mod tests { } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] - async fn sanitation_failure_keeps_observer_ticks_advancing_and_exposes_fact() { + async fn disabled_transport_keeps_observer_ticks_advancing_and_exposes_sanitation_fact() { let temp = tempfile::tempdir().unwrap(); - let (result, facts) = drive_real_link_start(&temp, false).await; + let (result, facts, _upload) = drive_real_link_start(&temp, false).await; assert!(result.is_err()); assert!(facts.config_sanitation_failed); } @@ -1393,7 +1431,7 @@ mod tests { } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] - async fn large_backpressured_upload_does_not_stop_observer_ticks() { + async fn slow_large_upload_response_does_not_stop_observer_ticks() { let temp = tempfile::tempdir().unwrap(); let peer = PrivateLinkPeer::start().await; peer.enqueue_response(200, br#"{"status":"ok","segment":"large"}"#.to_vec()); diff --git a/crates/solstone-linux/src/sync.rs b/crates/solstone-linux/src/sync.rs index 12d255a..8da747c 100644 --- a/crates/solstone-linux/src/sync.rs +++ b/crates/solstone-linux/src/sync.rs @@ -2909,7 +2909,7 @@ mod tests { } #[tokio::test] - async fn slow_linked_response_does_not_block_capture_or_delete_unproven_segment() { + async fn slow_linked_response_does_not_delete_unproven_segment() { let temp = tempfile::tempdir().unwrap(); let segment = create_segment(&temp, "120000_300", b"screen"); let legacy = MockServer::new(vec![]).await; @@ -2978,9 +2978,6 @@ mod tests { } } => {} } - let capture_ticks = AtomicUsize::new(0); - capture_ticks.fetch_add(1, Ordering::SeqCst); - assert_eq!(capture_ticks.load(Ordering::SeqCst), 1); assert!(segment.exists()); gate.notify_one(); cleanup.await;