diff --git a/PROGRESS.md b/PROGRESS.md index 98207b9cb..5eee4bc89 100644 --- a/PROGRESS.md +++ b/PROGRESS.md @@ -77,6 +77,10 @@ input cap plus one-byte overflow refusal, class-only stderr failures, and no key material on argv. `spl service` intentionally remains a fixed `spl: unavailable` / exit-69 interim path: it is not service composition and must never be treated as cutover completion. +- **Founder hold (2026-07-31):** browser-specific HPKE and relay support are being removed + from product scope. No lane work is deleted or reverted pending the supervisor's replacement + boundary. The in-progress U2 owned-dependency continuation is checkpointed as buildable but + unaccepted; it must be assessed against that new scope rather than completed by inertia. - Checkpoint gates after the correction-driven units: `cargo fmt --all -- --check`, strict combined clippy, and combined locked tests are green (85 SPL + 12 HPKE); both HPKE and SPL libraries pass the explicit `aarch64-apple-ios` check without an exclusion; `cargo deny` diff --git a/core/Cargo.lock b/core/Cargo.lock index 787dd6899..f39ac7233 100644 --- a/core/Cargo.lock +++ b/core/Cargo.lock @@ -334,6 +334,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a69dedd701da44b0536442edf09c81a64b0ab97a7a4a5e3d1971f00027cbc63d" dependencies = [ "const-oid", + "pem-rfc7468", "zeroize", ] @@ -419,6 +420,7 @@ dependencies = [ "group", "hkdf 0.13.0", "hybrid-array", + "pem-rfc7468", "pkcs8", "rand_core 0.10.1", "sec1", @@ -1037,6 +1039,15 @@ dependencies = [ "serde_core", ] +[[package]] +name = "pem-rfc7468" +version = "1.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a6305423e0e7738146434843d1694d621cce767262b2a86910beab705e4493d9" +dependencies = [ + "base64ct", +] + [[package]] name = "percent-encoding" version = "2.3.2" diff --git a/core/crates/solstone-core-spl/Cargo.toml b/core/crates/solstone-core-spl/Cargo.toml index 0a639765c..5b9ea9a2c 100644 --- a/core/crates/solstone-core-spl/Cargo.toml +++ b/core/crates/solstone-core-spl/Cargo.toml @@ -14,7 +14,7 @@ bytes = { workspace = true } flate2 = "1.1.9" futures-util = { version = "0.3.33", default-features = false, features = ["sink", "std"] } hpke = { workspace = true } -p256 = { version = "0.14.0", default-features = false, features = ["pkcs8", "std"] } +p256 = { version = "0.14.0", default-features = false, features = ["pem", "pkcs8", "std"] } regex = "1.12.3" serde_json = { workspace = true } solstone-core-spl-hpke = { workspace = true } diff --git a/core/crates/solstone-core-spl/src/authorized_client_ledger.rs b/core/crates/solstone-core-spl/src/authorized_client_ledger.rs index 0b70673c6..95b823e16 100644 --- a/core/crates/solstone-core-spl/src/authorized_client_ledger.rs +++ b/core/crates/solstone-core-spl/src/authorized_client_ledger.rs @@ -6,7 +6,10 @@ use std::fs; use std::path::{Path, PathBuf}; -use crate::{AuthorizedClients, BrowserUploadAuthorization, parse_authorized_clients}; +use crate::{ + AuthorizedClients, BrowserLedger, LedgerError, LedgerRow, LedgerStatus, + parse_authorized_clients, +}; /// An owned result from a browser authorization lookup. /// @@ -57,38 +60,69 @@ impl AuthorizedClientLedger { /// Resolves a browser sender's authorization from a freshly read ledger. #[must_use] pub fn lookup_browser(&self, sender_fingerprint: &[u8; 32]) -> BrowserLedgerLookup { - let ledger = self.read_fresh(); - match ledger.browser_upload_authorization(sender_fingerprint) { - BrowserUploadAuthorization::LedgerUnavailable => BrowserLedgerLookup::LedgerUnavailable, - BrowserUploadAuthorization::NotAuthorized => BrowserLedgerLookup::NotAuthorized, - BrowserUploadAuthorization::Incomplete => BrowserLedgerLookup::Incomplete, - BrowserUploadAuthorization::Authorized { - pubkey_spki, - observer_handle, - } => BrowserLedgerLookup::Authorized { - pubkey_spki: pubkey_spki.to_owned(), - observer_handle: observer_handle.to_owned(), + let fingerprint = fingerprint_key(sender_fingerprint); + match self.lookup(&fingerprint) { + Err(_) => BrowserLedgerLookup::LedgerUnavailable, + Ok(None) => BrowserLedgerLookup::NotAuthorized, + Ok(Some(row)) => match ( + row.pubkey_spki_hex.filter(|value| !value.is_empty()), + row.observer_handle.filter(|value| !value.is_empty()), + ) { + (Some(pubkey_spki), Some(observer_handle)) => BrowserLedgerLookup::Authorized { + pubkey_spki, + observer_handle, + }, + _ => BrowserLedgerLookup::Incomplete, }, } } - fn read_fresh(&self) -> AuthorizedClients { - // Read metadata on every call before opening the file. Unlike a - // last-good cache, a failed stat or read cannot preserve authorization. - if fs::metadata(&self.path) + fn read_fresh(&self) -> Result { + // The Python pairing process is the only writer. Stat before every + // lookup, then read from disk instead of retaining a last-good cache: + // a write made between browser offers is visible immediately, and a + // disappeared or unreadable file revokes authorization fail-closed. + fs::metadata(&self.path) .and_then(|metadata| metadata.modified()) - .is_err() - { - return AuthorizedClients::unavailable(); + .map_err(|_| LedgerError::Unavailable)?; + let bytes = fs::read(&self.path).map_err(|_| LedgerError::Unavailable)?; + let ledger = parse_authorized_clients(&bytes); + if ledger.status() == LedgerStatus::Unavailable { + return Err(LedgerError::Malformed); } + Ok(ledger) + } +} - match fs::read(&self.path) { - Ok(bytes) => parse_authorized_clients(&bytes), - Err(_) => AuthorizedClients::unavailable(), +impl BrowserLedger for AuthorizedClientLedger { + fn lookup(&self, fingerprint: &str) -> Result, LedgerError> { + let ledger = self.read_fresh()?; + let Some(entry) = ledger.entries().get(fingerprint) else { + return Ok(None); + }; + if entry.kind != "browser" { + return Ok(None); } + + Ok(Some(LedgerRow { + pubkey_spki_hex: entry.pubkey_spki.clone(), + observer_handle: entry.observer_handle.clone(), + })) } } +fn fingerprint_key(sender_fingerprint: &[u8; 32]) -> String { + const HEX: &[u8; 16] = b"0123456789abcdef"; + + let mut fingerprint = String::with_capacity("sha256:".len() + sender_fingerprint.len() * 2); + fingerprint.push_str("sha256:"); + for &byte in sender_fingerprint { + fingerprint.push(char::from(HEX[usize::from(byte >> 4)])); + fingerprint.push(char::from(HEX[usize::from(byte & 0x0f)])); + } + fingerprint +} + #[cfg(test)] mod tests { use std::error::Error; @@ -98,6 +132,7 @@ mod tests { use std::time::{SystemTime, UNIX_EPOCH}; use super::{AuthorizedClientLedger, BrowserLedgerLookup}; + use crate::{BrowserLedger, LedgerRow}; static NEXT_TEST_DIRECTORY: AtomicUsize = AtomicUsize::new(0); @@ -140,6 +175,13 @@ mod tests { ledger.lookup_browser(&fingerprint), BrowserLedgerLookup::Incomplete ); + assert_eq!( + ledger.lookup(&format!("sha256:{}", "a5".repeat(32)))?, + Some(LedgerRow { + pubkey_spki_hex: Some("30aa".to_owned()), + observer_handle: None, + }) + ); fs::write( journal.ledger_path(), @@ -152,6 +194,13 @@ mod tests { observer_handle: "observer-handle".to_owned(), } ); + assert_eq!( + ledger.lookup(&format!("sha256:{}", "a5".repeat(32)))?, + Some(LedgerRow { + pubkey_spki_hex: Some("30aa".to_owned()), + observer_handle: Some("observer-handle".to_owned()), + }) + ); Ok(()) } diff --git a/core/crates/solstone-core-spl/src/blob_receive.rs b/core/crates/solstone-core-spl/src/blob_receive.rs index 6ba22d010..8f3750309 100644 --- a/core/crates/solstone-core-spl/src/blob_receive.rs +++ b/core/crates/solstone-core-spl/src/blob_receive.rs @@ -16,12 +16,7 @@ //! receiver outcomes. This preserves that known U2 health-observability hole: //! the sole event here is the required `admission_saturated` accounting event. -use std::{ - future::Future, - pin::Pin, - sync::Arc, - time::Duration, -}; +use std::{future::Future, pin::Pin, sync::Arc, time::Duration}; use bytes::Bytes; use serde_json::{Value, json}; @@ -231,6 +226,12 @@ pub enum BlobError { Acknowledgement, } +struct AuthorizedBlobOffer { + offer: solstone_core_spl_hpke::Offer, + sender_public_key_hex: String, + observer_handle: String, +} + /// Receives, authenticates, validates, and ingests one browser blob offer. /// /// `R` is the immutable read half and `W` the separate write half of the @@ -241,12 +242,34 @@ pub enum BlobError { pub async fn receive_blob( reader: &mut BufferedWsReader, sink: &mut W, - deps: &BlobDeps<'_>, + deps: &BlobDeps, + instance_id: [u8; 16], + gate: &BlobAdmissionGate, + emit: &dyn CallosumEmit, +) -> Result<(), BlobError> { + receive_blob_with_timing( + reader, + sink, + deps, + instance_id, + BlobReceiveTiming::default(), + gate, + emit, + ) + .await +} + +async fn receive_blob_with_timing( + reader: &mut BufferedWsReader, + sink: &mut W, + deps: &BlobDeps, + instance_id: [u8; 16], + timing: BlobReceiveTiming, gate: &BlobAdmissionGate, emit: &dyn CallosumEmit, ) -> Result<(), BlobError> { let header = match reader - .read_exactly_bounded(OFFER_LEN, deps.timing.offer_deadline) + .read_exactly_bounded(OFFER_LEN, timing.offer_deadline) .await { Ok(header) => header, @@ -260,19 +283,18 @@ pub async fn receive_blob( Err(BlobFrameError::AckKey) => return close_sink(sink).await, }; - let (sender_public_key_hex, observer_handle) = - match deps.ledger.lookup_browser(&offer.sender_fingerprint) { - BrowserLedgerLookup::LedgerUnavailable | BrowserLedgerLookup::NotAuthorized => { - return send_ready_then_close(sink, 0x01).await; - } - BrowserLedgerLookup::Incomplete => return send_ready_then_close(sink, 0x01).await, - BrowserLedgerLookup::Authorized { - pubkey_spki, - observer_handle, - } => (pubkey_spki, observer_handle), - }; - let sender = sender_fingerprint_key(&offer.sender_fingerprint); + let row = match deps.ledger.lookup(&sender) { + Err(_) | Ok(None) => return send_ready_then_close(sink, 0x01).await, + Ok(Some(row)) => row, + }; + let (Some(sender_public_key_hex), Some(observer_handle)) = ( + row.pubkey_spki_hex.filter(|value| !value.is_empty()), + row.observer_handle.filter(|value| !value.is_empty()), + ) else { + return send_ready_then_close(sink, 0x01).await; + }; + if !gate.try_acquire_sender(&sender) { emit.emit( "admission_saturated", @@ -288,9 +310,13 @@ pub async fn receive_blob( reader, sink, deps, - offer, - &sender_public_key_hex, - &observer_handle, + instance_id, + timing, + AuthorizedBlobOffer { + offer, + sender_public_key_hex, + observer_handle, + }, ) .await; gate.release_sender(&sender); @@ -300,32 +326,32 @@ pub async fn receive_blob( async fn receive_permitted( reader: &mut BufferedWsReader, sink: &mut W, - deps: &BlobDeps<'_>, - offer: solstone_core_spl_hpke::Offer, - sender_public_key_hex: &str, - observer_handle: &str, + deps: &BlobDeps, + instance_id: [u8; 16], + timing: BlobReceiveTiming, + authorized_offer: AuthorizedBlobOffer, ) -> Result<(), BlobError> { if let Err(error) = send_ready(sink, 0x00).await { return close_after_error(sink, error).await; } let encapsulated_key = match reader - .read_exactly_bounded(ENC_LEN, deps.timing.enc_deadline) + .read_exactly_bounded(ENC_LEN, timing.enc_deadline) .await { Ok(encapsulated_key) => encapsulated_key, Err(_) => return close_sink(sink).await, }; - let ciphertext_len = match usize::try_from(offer.ciphertext_len) { + let ciphertext_len = match usize::try_from(authorized_offer.offer.ciphertext_len) { Ok(ciphertext_len) => ciphertext_len, Err(_) => return close_after_error(sink, BlobError::CiphertextLength).await, }; let ciphertext = match reader .read_exactly_progress( ciphertext_len, - deps.timing.ciphertext_deadline, - deps.timing.ciphertext_progress_window, - deps.timing.ciphertext_min_bytes_per_window, + timing.ciphertext_deadline, + timing.ciphertext_progress_window, + timing.ciphertext_min_bytes_per_window, ) .await { @@ -333,28 +359,33 @@ async fn receive_permitted( Err(_) => return close_sink(sink).await, }; - let sender_public_key_der = match decode_sender_public_key(sender_public_key_hex) { - Ok(sender_public_key_der) => sender_public_key_der, - Err(error) => { - close_sink(sink).await?; - return Err(error); - } + let sender_public_key_der = + match decode_sender_public_key(&authorized_offer.sender_public_key_hex) { + Ok(sender_public_key_der) => sender_public_key_der, + Err(error) => { + close_sink(sink).await?; + return Err(error); + } + }; + let recipient_private_key = match deps.upload_key.private_key() { + Ok(private_key) => private_key, + Err(_) => return close_sink(sink).await, }; let prepared = match prepare_authenticated_blob( - offer, + authorized_offer.offer, &encapsulated_key, &ciphertext, - deps.recipient_private_key, + &recipient_private_key, &sender_public_key_der, - &deps.instance_id, + &instance_id, ) { Ok(prepared) => prepared, Err(_) => return close_sink(sink).await, }; let status = match deps - .ingestor - .ingest(&prepared.archive, observer_handle) + .ingest + .ingest(&prepared.archive, &authorized_offer.observer_handle) .await { Ok(status) => status, @@ -468,7 +499,7 @@ mod tests { io::Cursor, path::PathBuf, sync::{ - Mutex, + Arc, Mutex, atomic::{AtomicUsize, Ordering}, }, time::{SystemTime, UNIX_EPOCH}, @@ -490,7 +521,7 @@ mod tests { use super::{ BlobDeps, BlobIngest, BlobIngestError, BlobIngestFuture, BlobIngestStatus, CallosumEmit, - receive_blob, + KeyError, UploadKeySource, parse_convey_ingest_status, receive_blob, }; use crate::{ AuthorizedClientLedger, BlobAdmissionGate, BufferedWsReader, ValidatedBlobArchive, @@ -580,6 +611,26 @@ mod tests { status: BlobIngestStatus, } + struct FixedUploadKeySource(P256Secret); + + impl UploadKeySource for FixedUploadKeySource { + fn private_key(&self) -> Result { + Ok(self.0.clone()) + } + } + + fn blob_deps( + ledger: AuthorizedClientLedger, + private_key: P256Secret, + ingestor: impl BlobIngest + 'static, + ) -> BlobDeps { + BlobDeps::new( + Arc::new(ledger), + Arc::new(FixedUploadKeySource(private_key)), + Arc::new(ingestor), + ) + } + impl BlobIngest for StatusIngest { fn ingest<'a>( &'a self, @@ -590,6 +641,26 @@ mod tests { } } + struct ConveyResponseIngest(Value); + + impl BlobIngest for ConveyResponseIngest { + fn ingest<'a>( + &'a self, + _archive: &'a ValidatedBlobArchive, + _observer_handle: &'a str, + ) -> BlobIngestFuture<'a> { + Box::pin(future::ready(parse_convey_ingest_status(&self.0))) + } + } + + struct MissingUploadKeySource; + + impl UploadKeySource for MissingUploadKeySource { + fn private_key(&self) -> Result { + Err(KeyError::Unavailable) + } + } + struct TestJournal { root: PathBuf, } @@ -625,7 +696,7 @@ mod tests { let ledger = AuthorizedClientLedger::new(&journal.root); let private_key = recipient_private_key()?; let ingestor = UnusedIngest; - let deps = BlobDeps::new(&ledger, &private_key, [0_u8; 16], &ingestor); + let deps = blob_deps(ledger, private_key, ingestor); let gate = BlobAdmissionGate::default(); let emitter = Emitter::default(); @@ -639,6 +710,7 @@ mod tests { &mut malformed_reader, &mut malformed_sink, &deps, + [0_u8; 16], &gate, &emitter, ) @@ -655,6 +727,7 @@ mod tests { &mut nonmagic_reader, &mut nonmagic_sink, &deps, + [0_u8; 16], &gate, &emitter, ) @@ -665,12 +738,13 @@ mod tests { } #[tokio::test] - async fn unavailable_and_incomplete_ledgers_refuse_with_ready() -> Result<(), Box> { + async fn unavailable_absent_and_incomplete_ledgers_refuse_with_ready() + -> Result<(), Box> { let journal = TestJournal::create()?; let ledger = AuthorizedClientLedger::new(&journal.root); let private_key = recipient_private_key()?; let ingestor = UnusedIngest; - let deps = BlobDeps::new(&ledger, &private_key, [0_u8; 16], &ingestor); + let deps = blob_deps(ledger, private_key, ingestor); let gate = BlobAdmissionGate::default(); let emitter = Emitter::default(); @@ -682,12 +756,30 @@ mod tests { &mut unavailable_reader, &mut unavailable_sink, &deps, + [0_u8; 16], &gate, &emitter, ) .await?; assert_eq!(unavailable_sink.sent, [Bytes::from_static(b"SBR1\x01\x01")]); + journal.write_ledger("[]")?; + let mut absent_reader = BufferedWsReader::new(Frames::with_frames([Bytes::from( + valid_offer_header(0).to_vec(), + )])); + let mut absent_sink = Sink::default(); + receive_blob( + &mut absent_reader, + &mut absent_sink, + &deps, + [0_u8; 16], + &gate, + &emitter, + ) + .await?; + assert_eq!(absent_sink.sent, [Bytes::from_static(b"SBR1\x01\x01")]); + assert_eq!(absent_sink.closes, 1); + journal.write_ledger(&browser_entry(None))?; let mut incomplete_reader = BufferedWsReader::new(Frames::with_frames([Bytes::from( valid_offer_header(0).to_vec(), @@ -697,6 +789,7 @@ mod tests { &mut incomplete_reader, &mut incomplete_sink, &deps, + [0_u8; 16], &gate, &emitter, ) @@ -713,7 +806,7 @@ mod tests { let ledger = AuthorizedClientLedger::new(&journal.root); let private_key = recipient_private_key()?; let ingestor = UnusedIngest; - let deps = BlobDeps::new(&ledger, &private_key, [0_u8; 16], &ingestor); + let deps = blob_deps(ledger, private_key, ingestor); let gate = BlobAdmissionGate::new(2, 1); let sender = sender_key_for_test(); assert!(gate.try_acquire_sender(&sender)); @@ -723,7 +816,7 @@ mod tests { )])); let mut sink = Sink::default(); - receive_blob(&mut reader, &mut sink, &deps, &gate, &emitter).await?; + receive_blob(&mut reader, &mut sink, &deps, [0_u8; 16], &gate, &emitter).await?; assert!(sink.sent.is_empty()); assert_eq!(sink.closes, 1); @@ -754,7 +847,7 @@ mod tests { let ledger = AuthorizedClientLedger::new(&journal.root); let private_key = recipient_private_key()?; let ingestor = UnusedIngest; - let deps = BlobDeps::new(&ledger, &private_key, [0_u8; 16], &ingestor); + let deps = blob_deps(ledger, private_key, ingestor); let gate = BlobAdmissionGate::new(2, 1); let emitter = Emitter::default(); let mut reader = BufferedWsReader::new(Frames::with_frames([Bytes::from( @@ -762,7 +855,7 @@ mod tests { )])); let mut sink = Sink::default(); - receive_blob(&mut reader, &mut sink, &deps, &gate, &emitter).await?; + receive_blob(&mut reader, &mut sink, &deps, [0_u8; 16], &gate, &emitter).await?; assert_eq!(sink.sent, [Bytes::from_static(b"SBR1\x01\x00")]); assert_eq!(sink.closes, 1); @@ -774,7 +867,7 @@ mod tests { async fn authenticated_ingest_statuses_send_exact_acks_and_close() -> Result<(), Box> { let recipient_private_key = recipient_private_key()?; - let transfer = authenticated_transfer(&recipient_private_key)?; + let transfer = authenticated_transfer(&recipient_private_key, [0_u8; 16])?; for (ingest_status, acknowledgement_status) in [ (BlobIngestStatus::Ok, 0x00), @@ -790,7 +883,7 @@ mod tests { let ingestor = StatusIngest { status: ingest_status, }; - let deps = BlobDeps::new(&ledger, &recipient_private_key, [0_u8; 16], &ingestor); + let deps = blob_deps(ledger, recipient_private_key.clone(), ingestor); let gate = BlobAdmissionGate::new(2, 1); let emitter = Emitter::default(); let expected_acknowledgement = ack( @@ -805,7 +898,7 @@ mod tests { ])); let mut sink = Sink::default(); - receive_blob(&mut reader, &mut sink, &deps, &gate, &emitter).await?; + receive_blob(&mut reader, &mut sink, &deps, [0_u8; 16], &gate, &emitter).await?; assert_eq!( sink.sent, @@ -820,6 +913,131 @@ mod tests { Ok(()) } + #[tokio::test] + async fn missing_upload_key_closes_after_ready_without_ack() -> Result<(), Box> { + let recipient_private_key = recipient_private_key()?; + let transfer = authenticated_transfer(&recipient_private_key, [0_u8; 16])?; + let journal = TestJournal::create()?; + journal.write_ledger(&browser_entry_with_sender_key( + Some("observer-handle"), + &transfer.sender_public_key_hex, + ))?; + let deps = BlobDeps::new( + Arc::new(AuthorizedClientLedger::new(&journal.root)), + Arc::new(MissingUploadKeySource), + Arc::new(StatusIngest { + status: BlobIngestStatus::Ok, + }), + ); + let mut reader = BufferedWsReader::new(Frames::with_frames([ + Bytes::copy_from_slice(&transfer.offer.header), + Bytes::copy_from_slice(&transfer.encapsulated_key), + Bytes::copy_from_slice(&transfer.ciphertext), + ])); + let mut sink = Sink::default(); + let gate = BlobAdmissionGate::default(); + let emitter = Emitter::default(); + + receive_blob(&mut reader, &mut sink, &deps, [0_u8; 16], &gate, &emitter).await?; + + assert_eq!(sink.sent, [Bytes::from_static(b"SBR1\x01\x00")]); + assert_eq!(sink.closes, 1); + Ok(()) + } + + #[tokio::test] + async fn unexpected_ingest_status_closes_without_ack() -> Result<(), Box> { + let recipient_private_key = recipient_private_key()?; + let transfer = authenticated_transfer(&recipient_private_key, [0_u8; 16])?; + let journal = TestJournal::create()?; + journal.write_ledger(&browser_entry_with_sender_key( + Some("observer-handle"), + &transfer.sender_public_key_hex, + ))?; + let deps = blob_deps( + AuthorizedClientLedger::new(&journal.root), + recipient_private_key, + ConveyResponseIngest(json!({"status":"not-a-convey-status"})), + ); + let mut reader = BufferedWsReader::new(Frames::with_frames([ + Bytes::copy_from_slice(&transfer.offer.header), + Bytes::copy_from_slice(&transfer.encapsulated_key), + Bytes::copy_from_slice(&transfer.ciphertext), + ])); + let mut sink = Sink::default(); + let gate = BlobAdmissionGate::default(); + let emitter = Emitter::default(); + + receive_blob(&mut reader, &mut sink, &deps, [0_u8; 16], &gate, &emitter).await?; + + assert_eq!(sink.sent, [Bytes::from_static(b"SBR1\x01\x00")]); + assert_eq!(sink.closes, 1); + Ok(()) + } + + #[test] + fn convey_status_parser_rejects_unknown_missing_and_non_string_values() { + assert_eq!( + parse_convey_ingest_status(&json!({"status":"ok"})), + Ok(BlobIngestStatus::Ok) + ); + assert_eq!( + parse_convey_ingest_status(&json!({"status":"duplicate"})), + Ok(BlobIngestStatus::Duplicate) + ); + assert_eq!( + parse_convey_ingest_status(&json!({"status":"collision"})), + Ok(BlobIngestStatus::Collision) + ); + for response in [json!({}), json!({"status":"other"}), json!({"status":7})] { + assert_eq!( + parse_convey_ingest_status(&response), + Err(BlobIngestError::UnexpectedStatus) + ); + } + } + + #[tokio::test] + async fn supplied_instance_id_is_bound_into_authenticated_open() -> Result<(), Box> { + let recipient_private_key = recipient_private_key()?; + let instance_id = [0x71_u8; 16]; + let transfer = authenticated_transfer(&recipient_private_key, instance_id)?; + let journal = TestJournal::create()?; + journal.write_ledger(&browser_entry_with_sender_key( + Some("observer-handle"), + &transfer.sender_public_key_hex, + ))?; + let deps = blob_deps( + AuthorizedClientLedger::new(&journal.root), + recipient_private_key, + StatusIngest { + status: BlobIngestStatus::Ok, + }, + ); + let expected_acknowledgement = + ack(&transfer.offer.blob_id, 0x00, &transfer.acknowledgement_key)?; + let mut reader = BufferedWsReader::new(Frames::with_frames([ + Bytes::copy_from_slice(&transfer.offer.header), + Bytes::copy_from_slice(&transfer.encapsulated_key), + Bytes::copy_from_slice(&transfer.ciphertext), + ])); + let mut sink = Sink::default(); + let gate = BlobAdmissionGate::default(); + let emitter = Emitter::default(); + + receive_blob(&mut reader, &mut sink, &deps, instance_id, &gate, &emitter).await?; + + assert_eq!( + sink.sent, + [ + Bytes::from_static(b"SBR1\x01\x00"), + Bytes::copy_from_slice(&expected_acknowledgement), + ] + ); + assert_eq!(sink.closes, 1); + Ok(()) + } + fn recipient_private_key() -> Result> { P256Secret::from_slice(&[0x42; 32]).map_err(Into::into) } @@ -859,6 +1077,7 @@ mod tests { fn authenticated_transfer( recipient_private_key: &P256Secret, + instance_id: [u8; 16], ) -> Result> { let sender_private_key = P256Secret::from_slice(&[0x43; 32])?; let sender_public_key_der = sender_private_key @@ -873,7 +1092,7 @@ mod tests { let offer = parse_offer(&valid_offer_header(ciphertext_len))?; let mut info = Vec::with_capacity(b"spl-blob-v1".len() + 16 + 32); info.extend_from_slice(b"spl-blob-v1"); - info.extend_from_slice(&[0_u8; 16]); + info.extend_from_slice(&instance_id); info.extend_from_slice(&offer.sender_fingerprint); let recipient_public_key = hpke_public_key(recipient_private_key)?; diff --git a/core/crates/solstone-core-spl/src/lib.rs b/core/crates/solstone-core-spl/src/lib.rs index 2c73643ca..a15f54052 100644 --- a/core/crates/solstone-core-spl/src/lib.rs +++ b/core/crates/solstone-core-spl/src/lib.rs @@ -27,6 +27,7 @@ mod service; mod service_shutdown; mod service_transition; mod tunnel_route; +mod upload_key_source; mod ws_buffer; mod ws_sink; @@ -46,7 +47,8 @@ pub use blob_archive::{ pub use blob_content_type::blob_content_type; pub use blob_receive::{ BlobDeps, BlobError, BlobIngest, BlobIngestError, BlobIngestFuture, BlobIngestStatus, - BlobReceiveTiming, CallosumEmit, receive_blob, + BlobReceiveTiming, BrowserLedger, CallosumEmit, KeyError, LedgerError, LedgerRow, + UploadKeySource, parse_convey_ingest_status, receive_blob, }; pub use health::{ LINK_HEALTH_EVENT, OFFLINE_TUNNEL_REASONS, REASON_HOME_MISSING_MOBILE, @@ -85,5 +87,6 @@ pub use service_transition::{ PostureObservation, ServiceAction, ServiceLifecycle, TokenObservation, transition, }; pub use tunnel_route::{TunnelRoute, route_tunnel_prefix}; +pub use upload_key_source::JournalUploadKeySource; pub use ws_buffer::{BufferedWsReader, WsBufferError, WsByteSource, WsClosed}; pub use ws_sink::WsByteSink; diff --git a/core/crates/solstone-core-spl/src/upload_key_source.rs b/core/crates/solstone-core-spl/src/upload_key_source.rs new file mode 100644 index 000000000..838b963ab --- /dev/null +++ b/core/crates/solstone-core-spl/src/upload_key_source.rs @@ -0,0 +1,127 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +//! Load-only access to the home upload HPKE key. + +use std::{ + fs, + path::{Path, PathBuf}, +}; + +use p256::elliptic_curve::pkcs8::DecodePrivateKey; +use solstone_core_spl_hpke::P256Secret; + +use crate::{KeyError, UploadKeySource}; + +/// Load-only source for `/link/hpke/upload_private.pem`. +/// +/// This mirrors Python's `load_upload_key()`, deliberately not its distinct +/// pairing-only `load_or_generate_upload_key()` helper. It has no write or +/// generation operation, so a missing key cannot mint an incompatible one. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct JournalUploadKeySource { + path: PathBuf, +} + +impl JournalUploadKeySource { + /// Creates a source rooted at the journal directory without accessing it. + #[must_use] + pub fn new(journal_root: impl AsRef) -> Self { + Self { + path: journal_root + .as_ref() + .join("link") + .join("hpke") + .join("upload_private.pem"), + } + } +} + +impl UploadKeySource for JournalUploadKeySource { + fn private_key(&self) -> Result { + let pem = fs::read_to_string(&self.path).map_err(|_| KeyError::Unavailable)?; + P256Secret::from_pkcs8_pem(&pem).map_err(|_| KeyError::Invalid) + } +} + +#[cfg(test)] +mod tests { + use std::{ + error::Error, + fs, + path::PathBuf, + sync::atomic::{AtomicUsize, Ordering}, + time::{SystemTime, UNIX_EPOCH}, + }; + + use p256::elliptic_curve::pkcs8::EncodePrivateKey; + use solstone_core_spl_hpke::P256Secret; + + use super::JournalUploadKeySource; + use crate::{KeyError, UploadKeySource}; + + static NEXT_TEST_DIRECTORY: AtomicUsize = AtomicUsize::new(0); + + struct TestJournal { + root: PathBuf, + } + + impl TestJournal { + fn create() -> Result> { + let sequence = NEXT_TEST_DIRECTORY.fetch_add(1, Ordering::Relaxed); + let timestamp = SystemTime::now().duration_since(UNIX_EPOCH)?.as_nanos(); + let root = std::env::temp_dir().join(format!( + "solstone-spl-upload-key-source-{}-{timestamp}-{sequence}", + std::process::id() + )); + Ok(Self { root }) + } + + fn key_path(&self) -> PathBuf { + self.root + .join("link") + .join("hpke") + .join("upload_private.pem") + } + } + + impl Drop for TestJournal { + fn drop(&mut self) { + let _ = fs::remove_dir_all(&self.root); + } + } + + #[test] + fn missing_key_fails_without_creating_any_path() -> Result<(), Box> { + let journal = TestJournal::create()?; + let source = JournalUploadKeySource::new(&journal.root); + + assert_eq!(source.private_key(), Err(KeyError::Unavailable)); + assert!(!journal.root.exists()); + + Ok(()) + } + + #[test] + fn loads_only_a_p256_pkcs8_pem_key() -> Result<(), Box> { + let journal = TestJournal::create()?; + let path = journal.key_path(); + let parent = match path.parent() { + Some(parent) => parent, + None => return Err("test key path unexpectedly has no parent".into()), + }; + fs::create_dir_all(parent)?; + let expected = P256Secret::from_slice(&[0x42; 32])?; + fs::write(&path, expected.to_pkcs8_pem(Default::default())?.as_bytes())?; + + let loaded = JournalUploadKeySource::new(&journal.root).private_key()?; + assert_eq!(loaded.to_bytes(), expected.to_bytes()); + + fs::write(path, "not a private key")?; + assert_eq!( + JournalUploadKeySource::new(&journal.root).private_key(), + Err(KeyError::Invalid) + ); + Ok(()) + } +}