diff --git a/crates/solstone-linux/src/private_link.rs b/crates/solstone-linux/src/private_link.rs index 72f3e84..f2764d8 100644 --- a/crates/solstone-linux/src/private_link.rs +++ b/crates/solstone-linux/src/private_link.rs @@ -1695,7 +1695,7 @@ pub(crate) async fn start_private_link_owner( async fn finish_owner_start( mut session: PrivateLinkSession, ) -> Result { - let capability = session.capability(INGEST_PATH.to_owned()); + let capability = session.capability(); if session.opener.generation() == 0 { match capability.report_unauthorized(0).await { RepairOutcome::Repaired { .. } | RepairOutcome::AlreadySuperseded { .. } => {} @@ -1781,7 +1781,7 @@ pub(crate) async fn start_registered_private_link_for_test( } impl PrivateLinkSession { - pub(crate) fn capability(&self, _ingest_path: String) -> PrivateLinkCapability { + pub(crate) fn capability(&self) -> PrivateLinkCapability { PrivateLinkCapability { inner: Arc::new(PrivateLinkCapabilityInner { client: self.client.clone(), @@ -3658,7 +3658,7 @@ 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/devices/ingest".to_owned()); + let capability = session.capability(); let second_generation = publish_observer_registration( &session, &ObserverState { @@ -3792,7 +3792,7 @@ mod tests { let session = start_private_link_session(temp.path(), peer.credential(), "stream") .await .unwrap(); - let capability = session.capability("/ignored-v2-ingest-path".to_owned()); + let capability = session.capability(); assert!(matches!( capability .ingest(multipart::Form::new().text("envelope", "{}")) @@ -3843,7 +3843,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/devices/ingest".to_owned()); + let capability = session.capability(); for day in ["", "2026010", "202601011", "202601?1", "../20260101"] { assert!(matches!( capability.segments_day(day).await, @@ -3861,7 +3861,7 @@ mod tests { async fn typed_unauthorized_report_is_only_recovery_surface() { let peer = PrivateLinkPeer::start().await; let (_temp, session) = start_peer_session(&peer).await; - let capability = session.capability("/ingest".to_owned()); + let capability = session.capability(); assert!(matches!( capability.report_unauthorized(0).await, RepairOutcome::AlreadySuperseded { generation: 1 } @@ -3970,7 +3970,7 @@ mod tests { let form = reqwest::multipart::Form::new() .part("files", reqwest::multipart::Part::bytes(body.clone())); let upload = tokio::spawn({ - let capability = session.capability("/ingest".to_owned()); + let capability = session.capability(); async move { capability.ingest(form).await } }); peer.wait_for_request_staged_at_least(spl_core::mux::INITIAL_WINDOW) @@ -4011,10 +4011,7 @@ mod tests { .await; assert!(response.starts_with(b"HTTP/1.1 400")); assert!(matches!( - session - .capability("/ingest".to_owned()) - .segments_day("20260101") - .await, + session.capability().segments_day("20260101").await, LinkOutcome::Success { .. } )); assert_eq!(peer.requests().len(), 1); @@ -4275,10 +4272,7 @@ mod tests { ) .unwrap(); assert!(matches!( - session - .capability("/ingest".to_owned()) - .segments_day("20260101") - .await, + session.capability().segments_day("20260101").await, LinkOutcome::Success { .. } )); assert_eq!(peer.accepted_carriers(), 1); diff --git a/crates/solstone-linux/src/private_link_test_peer.rs b/crates/solstone-linux/src/private_link_test_peer.rs index 4d2f033..4bff229 100644 --- a/crates/solstone-linux/src/private_link_test_peer.rs +++ b/crates/solstone-linux/src/private_link_test_peer.rs @@ -457,8 +457,7 @@ fn next_response( return match queued { Some(QueuedResponse::Static(response)) => response, Some(QueuedResponse::DayCustody(fixture)) => { - let response = - fixture.response_for(crate::test_support::DayCustodyLeg::Manifest, None); + let response = fixture.response_for(crate::test_support::DayCustodyLeg::Manifest); if !fixture.stops_after(crate::test_support::DayCustodyLeg::Manifest) { day_manifests.push_back(fixture); } @@ -467,11 +466,13 @@ fn next_response( None => plain_response(500, Vec::new()), }; } - if let Some(day) = request.path.strip_prefix("/app/devices/ingest/manifest/") + if request + .path + .strip_prefix("/app/devices/ingest/manifest/") + .is_some() && let Some(fixture) = day_manifests.pop_front() { - let response = - fixture.response_for(crate::test_support::DayCustodyLeg::DayManifest, Some(day)); + let response = fixture.response_for(crate::test_support::DayCustodyLeg::DayManifest); if !fixture.stops_after(crate::test_support::DayCustodyLeg::DayManifest) { segment_lists.push_back(fixture); } @@ -483,7 +484,7 @@ fn next_response( .is_some() && let Some(fixture) = segment_lists.pop_front() { - let response = fixture.response_for(crate::test_support::DayCustodyLeg::Segments, None); + let response = fixture.response_for(crate::test_support::DayCustodyLeg::Segments); return plain_response(response.0, response.1); } pop_static_response(state) diff --git a/crates/solstone-linux/src/run.rs b/crates/solstone-linux/src/run.rs index dd03fd9..502767a 100644 --- a/crates/solstone-linux/src/run.rs +++ b/crates/solstone-linux/src/run.rs @@ -1402,7 +1402,7 @@ mod tests { let session = start_private_link_session(temp.path(), peer.credential(), "stream") .await .unwrap(); - let capability = session.capability("/app/devices/ingest".to_owned()); + let capability = session.capability(); let config = Config { config_dir: temp.path().to_path_buf(), stream: "stream".to_owned(), @@ -1646,7 +1646,7 @@ mod tests { }; let client = Arc::new(UploadClient::new( &config, - session.capability("/app/devices/ingest".to_owned()), + session.capability(), "host", "linux", "1", @@ -1689,7 +1689,7 @@ mod tests { reqwest::multipart::Part::stream(reqwest::Body::wrap_stream(stream)), ); let request = tokio::spawn({ - let capability = session.capability("/app/devices/ingest".to_owned()); + let capability = session.capability(); 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 b6cc06a..242c4e5 100644 --- a/crates/solstone-linux/src/sync.rs +++ b/crates/solstone-linux/src/sync.rs @@ -595,7 +595,7 @@ impl SyncWorker { Ok(files) => files, Err(error) => { tracing::warn!(%error, path = %segment_dir.display(), "Failed to enumerate segment files"); - self.last_error_type = Some(ErrorType::Transient); + self.last_error_type = Some(ErrorType::Client); self.last_error_code = None; return false; } @@ -887,11 +887,19 @@ fn segment_custody_proven(segment_dir: &Path, entry: &ListingEntry) -> bool { else { return false; }; - matches!(remote.status.as_deref(), Some("present" | "processed")) - && remote - .size - .is_some_and(|size| local.metadata().is_ok_and(|meta| meta.len() == size)) - && sha256_file(local).ok().as_deref() == remote.sha256.as_deref() + if !matches!(remote.status.as_deref(), Some("present" | "processed")) { + return false; + } + let Some(size) = remote.size else { + return false; + }; + if !local.metadata().is_ok_and(|meta| meta.len() == size) { + return false; + } + let Some(remote_sha256) = remote.sha256.as_deref() else { + return false; + }; + sha256_file(local).is_ok_and(|local_sha256| local_sha256 == remote_sha256) }) } @@ -1188,7 +1196,10 @@ mod tests { if let Some(status) = status { file["status"] = json!(status); } - custody(vec![json!({"key":key,"observed":true,"files":[file]})]) + custody_for_day( + "20260101", + vec![json!({"key":key,"observed":true,"files":[file]})], + ) } async fn test_worker( @@ -1271,7 +1282,7 @@ mod tests { .unwrap(); let client = Arc::new(UploadClient::new( &config, - session.capability("/app/devices/ingest".to_owned()), + session.capability(), "host", "linux", "test", @@ -1332,7 +1343,7 @@ mod tests { let (server, mut worker) = test_worker( &temp, vec![ - (200, custody(Vec::new())), + (200, custody_for_day("20270115", Vec::new())), (200, remote.clone()), (200, json!({"status":"ok","segment":"120000_300"})), (200, remote), @@ -1355,7 +1366,7 @@ mod tests { let (server, mut worker) = test_worker( &temp, vec![ - (200, custody(Vec::new())), + (200, custody_for_day("20270115", Vec::new())), (200, remote.clone()), (200, remote), ], @@ -1476,6 +1487,66 @@ mod tests { .await; } + #[tokio::test] + async fn present_status_with_absent_sha_uploads_and_cleanup_keeps() { + upload_then_cleanup_keeps(custody_for_day( + "20260101", + vec![json!({ + "key": "120000_300", + "observed": true, + "files": [{ + "name": "screen.webm", + "size": 6, + "status": "present", + }], + })], + )) + .await; + } + + #[tokio::test] + async fn present_status_with_absent_sha_and_unreadable_file_quarantines() { + use std::os::unix::fs::PermissionsExt; + + let temp = tempfile::tempdir().unwrap(); + let segment = create_segment(&temp, "120000_300", b"screen"); + let file = segment.join("screen.webm"); + fs::set_permissions(&file, fs::Permissions::from_mode(0o000)).unwrap(); + assert!(sha256_file(&file).is_err()); + let remote = custody_for_day( + "20260101", + vec![json!({ + "key": "120000_300", + "observed": true, + "files": [{ + "name": "screen.webm", + "size": 6, + "status": "present", + }], + })], + ); + let (server, mut worker) = test_worker( + &temp, + vec![ + (200, custody_for_day("20270115", Vec::new())), + (200, remote), + ], + 7, + ) + .await; + + worker.sync_pass(true).await; + + let failed = segment.with_file_name("120000_300.failed"); + assert!(!segment.exists()); + assert!(failed.exists()); + assert!(!failed.join(SERVER_KEY_FILENAME).exists()); + assert!(!worker.synced_days.contains("20260101")); + assert_eq!(upload_hits(&server), 0); + worker.cleanup_synced_segments().await; + assert!(failed.exists()); + } + #[tokio::test] async fn processed_status_with_absent_size_uploads() { let sha = format!("{:x}", Sha256::digest(b"screen")); @@ -1547,6 +1618,13 @@ mod tests { upload_then_cleanup_keeps(listing("120000_300", "screen.webm", None, &sha)).await; } + #[tokio::test] + async fn missing_custody_status_uploads_and_cleanup_keeps() { + let sha = format!("{:x}", Sha256::digest(b"screen")); + upload_then_cleanup_keeps(listing("120000_300", "screen.webm", Some("missing"), &sha)) + .await; + } + // tests/test_sync.py::test_unknown_status_uploads_and_cleanup_keeps #[tokio::test] async fn unknown_status_uploads_and_cleanup_keeps() { @@ -1569,14 +1647,17 @@ mod tests { let temp = tempfile::tempdir().unwrap(); let segment = create_segment(&temp, "120000_300", b"screen"); fs::write(segment.join("audio.flac"), b"audio").unwrap(); - let remote = custody(vec![json!({"key":"120000_300", "observed":true, "files":[ - {"name":"screen.webm","size":6,"status":"present","sha256":format!("{:x}",Sha256::digest(b"screen"))}, - {"name":"audio.flac","size":5,"status":"processed","sha256":format!("{:x}",Sha256::digest(b"audio"))} - ]})]); + let remote = custody_for_day( + "20260101", + vec![json!({"key":"120000_300", "observed":true, "files":[ + {"name":"screen.webm","size":6,"status":"present","sha256":format!("{:x}",Sha256::digest(b"screen"))}, + {"name":"audio.flac","size":5,"status":"processed","sha256":format!("{:x}",Sha256::digest(b"audio"))} + ]})], + ); let (server, mut worker) = test_worker( &temp, vec![ - (200, custody(Vec::new())), + (200, custody_for_day("20270115", Vec::new())), (200, remote.clone()), (200, remote), ], @@ -1628,19 +1709,33 @@ mod tests { assert!(segment.exists()); } - // AC: upload enumeration failures are transient and never send a partial request. + // AC: a file that cannot be statted is quarantined and never sends a partial request. #[tokio::test] - async fn unreadable_segment_directory_upload_is_retried() { - use std::os::unix::fs::PermissionsExt; + async fn unstatable_file_is_quarantined_without_upload() { + use std::os::unix::fs::symlink; + let temp = tempfile::tempdir().unwrap(); let segment = create_segment(&temp, "120000_300", b"screen"); - fs::set_permissions(&segment, fs::Permissions::from_mode(0o000)).unwrap(); + let file = segment.join("screen.webm"); + fs::remove_file(&file).unwrap(); + symlink(segment.join("missing.webm"), &file).unwrap(); assert!(eligible_files(&segment).is_err()); - let (server, mut worker) = test_worker(&temp, vec![], -1).await; - assert!(!worker.upload_segment("20260101", &segment).await); - fs::set_permissions(&segment, fs::Permissions::from_mode(0o700)).unwrap(); - assert_eq!(worker.last_error_type, Some(ErrorType::Transient)); - assert!(server.requests().is_empty()); + let (server, mut worker) = test_worker( + &temp, + vec![ + (200, custody_for_day("20270115", Vec::new())), + (200, custody_for_day("20260101", Vec::new())), + ], + -1, + ) + .await; + + worker.sync_pass(true).await; + + assert!(!segment.exists()); + assert!(segment.with_file_name("120000_300.failed").exists()); + assert_eq!(worker.last_error_type, Some(ErrorType::Client)); + assert_eq!(upload_hits(&server), 0); } #[tokio::test] @@ -1676,13 +1771,13 @@ mod tests { let (server, mut worker) = test_worker( &temp, vec![ - (200, custody(Vec::new())), - (200, custody(Vec::new())), + (200, custody_for_day("20270115", Vec::new())), + (200, custody_for_day("20260101", Vec::new())), ( 200, json!({"status":"duplicate","existing_segment":"existing_300"}), ), - (200, custody(Vec::new())), + (200, custody_for_day("20270115", Vec::new())), (200, held), ], -1, @@ -1704,16 +1799,19 @@ mod tests { let temp = tempfile::tempdir().unwrap(); let segment = create_segment(&temp, "120000_300", b"screen"); let sha = sha256_file(&segment.join("screen.webm")).unwrap(); - let held = custody(vec![ - json!({"key":"120000_301","original_key":"120000_300","observed":true,"files":[{"name":"screen.webm","size":6,"status":"present","sha256":sha}]}), - ]); + let held = custody_for_day( + "20260101", + vec![ + json!({"key":"120000_301","original_key":"120000_300","observed":true,"files":[{"name":"screen.webm","size":6,"status":"present","sha256":sha}]}), + ], + ); let (server, mut worker) = test_worker( &temp, vec![ - (200, custody(Vec::new())), - (200, custody(Vec::new())), + (200, custody_for_day("20270115", Vec::new())), + (200, custody_for_day("20260101", Vec::new())), (200, json!({"status":"ok","segment":"120000_301"})), - (200, custody(Vec::new())), + (200, custody_for_day("20270115", Vec::new())), (200, held), ], -1, @@ -2028,18 +2126,21 @@ mod tests { let media = segment.join("screen.webm"); fs::write(&media, b"screen").unwrap(); let sha = sha256_file(&media).unwrap(); - let held = custody(vec![json!({ - "key": "120000_300", - "observed": true, - "files": [{ - "name": "screen.webm", - "size": 6, - "status": "present", - "sha256": sha - }] - })]); + let held = custody_for_day( + "20260101", + vec![json!({ + "key": "120000_300", + "observed": true, + "files": [{ + "name": "screen.webm", + "size": 6, + "status": "present", + "sha256": sha + }] + })], + ); let server = MockServer::new(vec![ - (200, custody(Vec::new())), + (200, custody_for_day("20270115", Vec::new())), (200, held.clone()), (200, held), ]) @@ -3044,7 +3145,10 @@ mod tests { async fn query_failures_recover_to_connected() { let temp = tempfile::tempdir().unwrap(); let mut responses = (0..5).map(|_| (500, json!({}))).collect::>(); - responses.extend([(200, json!({"days": {}})), (200, custody(Vec::new()))]); + responses.extend([ + (200, json!({"days": {}})), + (200, custody_for_day("20270115", Vec::new())), + ]); let (server, mut worker) = test_worker(&temp, responses, -1).await; let clock = Arc::new(MutableClock::new(1_800_000_000.0, 100.0)); worker.clock = clock.clone(); @@ -3238,7 +3342,11 @@ mod tests { ); let (_server, mut worker) = test_worker( &temp, - vec![(200, custody(Vec::new())), (200, held.clone()), (200, held)], + vec![ + (200, custody_for_day("20270115", Vec::new())), + (200, held.clone()), + (200, held), + ], 7, ) .await; @@ -3340,7 +3448,7 @@ mod tests { }); let client = Arc::new(UploadClient::new( &config, - session.capability("/app/devices/ingest".into()), + session.capability(), "host", "linux", "test", @@ -3407,7 +3515,7 @@ mod tests { }); let client = Arc::new(UploadClient::new( &config, - session.capability("/app/devices/ingest".into()), + session.capability(), "host", "linux", "test", @@ -3844,7 +3952,7 @@ mod tests { #[tokio::test] async fn completion_trigger_starts_pass() { let temp = tempfile::tempdir().unwrap(); - let server = MockServer::new(vec![(200, custody(Vec::new()))]).await; + let server = MockServer::new(vec![(200, custody_for_day("20270115", Vec::new()))]).await; let config = Config { base_dir: temp.path().to_path_buf(), config_dir: temp.path().join("config"), @@ -3925,9 +4033,9 @@ mod tests { &sha256_file(&media).unwrap(), ); let responses = vec![ - (200, custody(Vec::new())), + (200, custody_for_day("20270115", Vec::new())), (200, held.clone()), - (200, custody(Vec::new())), + (200, custody_for_day("20270115", Vec::new())), (200, custody_for_day("20270116", Vec::new())), (200, held), ]; @@ -3975,8 +4083,8 @@ mod tests { let (server, mut worker) = test_worker( &temp, vec![ - (200, custody(Vec::new())), - (200, custody(Vec::new())), + (200, custody_for_day("20270115", Vec::new())), + (200, custody_for_day("20260101", Vec::new())), (200, custody_for_day("20250101", Vec::new())), ], -1, @@ -4183,7 +4291,7 @@ mod tests { .unwrap(); let client = Arc::new(UploadClient::new( &config, - session.capability("/app/devices/ingest".to_owned()), + session.capability(), "host", "linux", "test", @@ -4245,7 +4353,7 @@ mod tests { .unwrap(); let client = Arc::new(UploadClient::new( &config, - session.capability("/app/devices/ingest".to_owned()), + session.capability(), "host", "linux", "test", @@ -4315,7 +4423,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/devices/ingest".to_owned()), + session.capability(), "host", "linux", "test", diff --git a/crates/solstone-linux/src/test_support.rs b/crates/solstone-linux/src/test_support.rs index 725af09..46319e6 100644 --- a/crates/solstone-linux/src/test_support.rs +++ b/crates/solstone-linux/src/test_support.rs @@ -119,11 +119,7 @@ impl DayCustodyFixture { self } - pub(crate) fn response_for( - &self, - leg: DayCustodyLeg, - requested_day: Option<&str>, - ) -> (u16, Vec) { + pub(crate) fn response_for(&self, leg: DayCustodyLeg) -> (u16, Vec) { if let Some((failed_leg, status, bytes)) = &self.failed && *failed_leg == leg { @@ -138,20 +134,15 @@ impl DayCustodyFixture { DayCustodyLeg::Manifest => { let mut days = serde_json::Map::new(); if !self.absent { - // Sync tests exercise both their fixed "today" and retained capture day. - // The peer fills the following legs from the request path, so this declares - // both candidates without coupling ordered fixtures to call order. - for day in [&self.day, "20260101", "20270115"] { - days.insert( - day.to_owned(), - serde_json::json!({"segments": self.items.len()}), - ); - } + days.insert( + self.day.clone(), + serde_json::json!({"segments": self.items.len()}), + ); } serde_json::json!({"days": days}) } DayCustodyLeg::DayManifest => serde_json::json!({ - "day": self.day_manifest_day.as_deref().or(requested_day).unwrap_or(&self.day), + "day": self.day_manifest_day.as_deref().unwrap_or(&self.day), "version": self.version, "segments": {}, }), @@ -341,7 +332,7 @@ impl LinkedMockServer { } pub(crate) fn capability(&self) -> PrivateLinkCapability { - self.session.capability("/app/devices/ingest".to_owned()) + self.session.capability() } pub(crate) fn enqueue_day_custody(&self, fixture: DayCustodyFixture) { diff --git a/crates/solstone-linux/src/upload.rs b/crates/solstone-linux/src/upload.rs index eaa8606..a81838a 100644 --- a/crates/solstone-linux/src/upload.rs +++ b/crates/solstone-linux/src/upload.rs @@ -1037,7 +1037,10 @@ fn parse_segments_envelope(body: Value, status_code: u16) -> DayCustody { let Some(items) = object.get("items").and_then(Value::as_array) else { return custody_failure(ErrorType::Incompatible, Some(status_code)); }; - if protocol_version != 3 || total != items.len() as u64 { + if protocol_version != 3 { + return custody_failure(ErrorType::Incompatible, Some(status_code)); + } + if total != items.len() as u64 { return DayCustody { day_present: true, items: Vec::new(), @@ -1167,7 +1170,7 @@ mod tests { .unwrap(); let client = UploadClient::new( &config, - session.capability("/app/devices/ingest".into()), + session.capability(), "host-a", "linux", "0.1.0", @@ -1208,7 +1211,7 @@ mod tests { .unwrap(); let client = UploadClient::new( &config, - session.capability("/app/devices/ingest".into()), + session.capability(), "host", "linux", "test", @@ -1559,7 +1562,7 @@ mod tests { .unwrap(); let client = UploadClient::new( &config, - session.capability("/app/devices/ingest".into()), + session.capability(), "host-a", "linux", "0.1.0", @@ -1993,7 +1996,8 @@ mod tests { ); } - // tests/test_upload.py::test_upload_segment_uses_bearer_and_keyless_route + // Python ancestor: tests/test_upload.py::test_upload_segment_uses_bearer_and_keyless_route. + // V3 uses certificate-only identity; this test asserts those identity headers are absent. // tests/test_upload.py::test_upload_segment_declares_content_types // AC: multipart received bytes preserve fields, filenames, content types, and file content #[tokio::test] @@ -2384,14 +2388,12 @@ mod tests { } #[tokio::test] - async fn v3_bad_segments_protocol_and_malformed_legs_are_incompatible() { + async fn v3_wrong_segments_protocol_and_malformed_legs_are_incompatible() { let protocol = fetch_custody_fixture( DayCustodyFixture::new("20260101", Vec::new()).with_segments_protocol_version(2), ) .await; - assert!(protocol.day_present); - assert!(!protocol.proof_available); - assert!(protocol.error_type.is_none()); + assert_eq!(protocol.error_type, Some(ErrorType::Incompatible)); let malformed = fetch_custody_fixture( DayCustodyFixture::new("20260101", Vec::new()) @@ -2465,7 +2467,7 @@ mod tests { .unwrap(); let client = UploadClient::new( &config, - session.capability("/app/devices/ingest".to_owned()), + session.capability(), "host", "linux", "test", @@ -2633,7 +2635,7 @@ mod tests { }; let client = UploadClient::new( &config, - session.capability("/app/devices/ingest".to_owned()), + session.capability(), "host", "linux", "1", @@ -2718,7 +2720,7 @@ mod tests { .unwrap(); let client = UploadClient::new( &config, - session.capability("/app/devices/ingest".into()), + session.capability(), "host-a", "linux", "0.1.0",