From 1ff0a344b43618e29315d9be8fd9eb27b2ca4940 Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Fri, 31 Jul 2026 16:31:17 -0600 Subject: [PATCH] Add U2 blob preparation units --- core/crates/solstone-core-spl/Cargo.toml | 25 + .../src/authenticated_blob.rs | 281 ++++++++++ .../src/authorized_client_ledger.rs | 210 +++++++ .../src/authorized_clients.rs | 400 ++++++++++++++ .../solstone-core-spl/src/blob_archive.rs | 516 ++++++++++++++++++ .../src/blob_content_type.rs | 49 ++ 6 files changed, 1481 insertions(+) create mode 100644 core/crates/solstone-core-spl/Cargo.toml create mode 100644 core/crates/solstone-core-spl/src/authenticated_blob.rs create mode 100644 core/crates/solstone-core-spl/src/authorized_client_ledger.rs create mode 100644 core/crates/solstone-core-spl/src/authorized_clients.rs create mode 100644 core/crates/solstone-core-spl/src/blob_archive.rs create mode 100644 core/crates/solstone-core-spl/src/blob_content_type.rs diff --git a/core/crates/solstone-core-spl/Cargo.toml b/core/crates/solstone-core-spl/Cargo.toml new file mode 100644 index 000000000..7f1c5f792 --- /dev/null +++ b/core/crates/solstone-core-spl/Cargo.toml @@ -0,0 +1,25 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +[package] +name = "solstone-core-spl" +version.workspace = true +edition.workspace = true +rust-version.workspace = true +license.workspace = true +publish = false + +[dependencies] +bytes = { workspace = true } +flate2 = "1.1.9" +hpke = { workspace = true } +p256 = { version = "0.14.0", default-features = false, features = ["pkcs8", "std"] } +regex = "1.12.3" +serde_json = { workspace = true } +solstone-core-spl-hpke = { workspace = true } +tar = "0.4.45" +thiserror = { workspace = true } +tokio = { workspace = true, features = ["io-util", "rt", "macros", "sync", "time"] } + +[lints] +workspace = true diff --git a/core/crates/solstone-core-spl/src/authenticated_blob.rs b/core/crates/solstone-core-spl/src/authenticated_blob.rs new file mode 100644 index 000000000..e43fd1d55 --- /dev/null +++ b/core/crates/solstone-core-spl/src/authenticated_blob.rs @@ -0,0 +1,281 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +//! Pure preparation of an authenticated SPL blob for later ingestion. + +use solstone_core_spl_hpke::{HpkeError, Offer, P256Secret, open_auth}; +use thiserror::Error; + +use crate::{BlobArchiveError, ValidatedBlobArchive, parse_blob_archive}; + +const BLOB_INFO_PREFIX: &[u8] = b"spl-blob-v1"; +const ACK_EXPORT_CONTEXT: &[u8] = b"spl-blob-ack-v1"; +const ACK_KEY_LEN: usize = 32; +const BLOB_INFO_LEN: usize = BLOB_INFO_PREFIX.len() + 16 + 32; + +/// A blob which has passed authenticated decryption and archive validation. +#[derive(Debug)] +pub struct PreparedAuthenticatedBlob { + /// The offer retained verbatim for acknowledgement construction. + pub offer: Offer, + /// Archive content which is safe for the caller to ingest. + pub archive: ValidatedBlobArchive, + /// The RFC 9180 exporter key for the later `SBA1` acknowledgement. + pub acknowledgement_key: [u8; ACK_KEY_LEN], +} + +/// Failures while preparing an authenticated blob without any side effects. +#[derive(Debug, Error)] +pub enum AuthenticatedBlobError { + /// The transfer did not contain exactly the length declared by its offer. + #[error("blob ciphertext length mismatch: expected {expected}, got {actual}")] + CiphertextLength { + /// The length from the validated offer header. + expected: u64, + /// The bytes supplied by the transport. + actual: usize, + }, + + /// The encrypted blob could not be authenticated and opened. + #[error("blob HPKE open failed: {0}")] + Hpke(#[from] HpkeError), + + /// The opened plaintext was not a safe SPL blob archive. + #[error("blob archive validation failed: {0}")] + Archive(#[from] BlobArchiveError), + + /// The HPKE exporter returned a length other than the protocol's 32 bytes. + #[error("blob acknowledgement exporter produced {actual} bytes, expected {ACK_KEY_LEN}")] + ExporterLength { + /// The number of bytes returned by the HPKE exporter. + actual: usize, + }, +} + +/// Authenticate, validate, and prepare a received browser blob. +/// +/// The caller supplies only already-read protocol data. This function neither +/// performs I/O nor decides authorization, admission, acknowledgement status, +/// or ingestion. It preserves the offer's exact 67-byte header as HPKE AAD. +pub fn prepare_authenticated_blob( + offer: Offer, + enc: &[u8], + ciphertext: &[u8], + recipient_private_key: &P256Secret, + sender_public_key_spki_der: &[u8], + instance_id: &[u8; 16], +) -> Result { + if !ciphertext_length_matches(ciphertext.len(), offer.ciphertext_len) { + return Err(AuthenticatedBlobError::CiphertextLength { + expected: offer.ciphertext_len, + actual: ciphertext.len(), + }); + } + + let info = blob_info(instance_id, &offer.sender_fingerprint); + let opened = open_auth( + enc, + recipient_private_key, + &info, + sender_public_key_spki_der, + ciphertext, + &offer.header, + )?; + let archive = parse_blob_archive(&opened.plaintext)?; + let acknowledgement_key = acknowledgement_key(&opened)?; + + Ok(PreparedAuthenticatedBlob { + offer, + archive, + acknowledgement_key, + }) +} + +fn ciphertext_length_matches(actual: usize, expected: u64) -> bool { + match u64::try_from(actual) { + Ok(actual) => actual == expected, + Err(_) => false, + } +} + +fn blob_info(instance_id: &[u8; 16], sender_fingerprint: &[u8; 32]) -> [u8; BLOB_INFO_LEN] { + let mut info = [0_u8; BLOB_INFO_LEN]; + let prefix_end = BLOB_INFO_PREFIX.len(); + let instance_end = prefix_end + instance_id.len(); + info[..prefix_end].copy_from_slice(BLOB_INFO_PREFIX); + info[prefix_end..instance_end].copy_from_slice(instance_id); + info[instance_end..].copy_from_slice(sender_fingerprint); + info +} + +fn acknowledgement_key( + opened: &solstone_core_spl_hpke::OpenedHpke, +) -> Result<[u8; ACK_KEY_LEN], AuthenticatedBlobError> { + let exported = opened.export(ACK_EXPORT_CONTEXT, ACK_KEY_LEN)?; + let actual = exported.len(); + exported + .try_into() + .map_err(|_| AuthenticatedBlobError::ExporterLength { actual }) +} + +#[cfg(test)] +mod tests { + use std::io::Cursor; + + use flate2::{Compression, write::GzEncoder}; + use hpke::{ + Deserializable, OpModeS, Serializable, + aead::AesGcm256, + kdf::HkdfSha256, + kem::{DhP256HkdfSha256, Kem as KemTrait}, + setup_sender, + }; + use p256::elliptic_curve::{pkcs8::EncodePublicKey, sec1::ToSec1Point}; + use solstone_core_spl_hpke::{OFFER_LEN, Offer, P256Secret, parse_offer}; + use tar::{Builder, Header}; + + use super::{AuthenticatedBlobError, blob_info, prepare_authenticated_blob}; + + type SuiteKem = DhP256HkdfSha256; + + #[test] + fn blob_info_is_the_exact_protocol_concatenation() { + let instance_id = [0x11; 16]; + let sender_fingerprint = [0x22; 32]; + let info = blob_info(&instance_id, &sender_fingerprint); + + assert_eq!(&info[..11], b"spl-blob-v1"); + assert_eq!(&info[11..27], &instance_id); + assert_eq!(&info[27..], &sender_fingerprint); + } + + #[test] + fn rejects_a_ciphertext_length_mismatch_before_hpke() -> Result<(), Box> + { + let offer = Offer { + header: [0_u8; OFFER_LEN], + sender_fingerprint: [0_u8; 32], + blob_id: [0_u8; 16], + ciphertext_len: 4, + }; + let recipient_private_key = P256Secret::from_slice(&[0x42; 32])?; + let result = prepare_authenticated_blob( + offer, + b"", + b"bad", + &recipient_private_key, + b"", + &[0_u8; 16], + ); + + assert!(matches!( + result, + Err(AuthenticatedBlobError::CiphertextLength { + expected: 4, + actual: 3 + }) + )); + Ok(()) + } + + #[test] + fn opens_an_authenticated_archive_and_derives_the_acknowledgement_key() + -> Result<(), Box> { + let recipient_private_key = P256Secret::from_slice(&[0x42; 32])?; + let sender_private_key = P256Secret::from_slice(&[0x43; 32])?; + let sender_public_key_spki_der = sender_private_key + .public_key() + .to_public_key_der()? + .as_bytes() + .to_vec(); + let archive = valid_archive()?; + let instance_id = [0x11; 16]; + let offer = valid_offer(archive.len() + 16)?; + let info = blob_info(&instance_id, &offer.sender_fingerprint); + + let recipient_public_key = hpke_public_key(&recipient_private_key)?; + let sender_private_key = hpke_private_key(&sender_private_key)?; + let sender_public_key = hpke_public_key_from_private_key(&sender_private_key)?; + let (encapsulated_key, mut sender_context) = setup_sender::( + &OpModeS::Auth((sender_private_key, sender_public_key)), + &recipient_public_key, + &info, + )?; + let ciphertext = sender_context.seal(&archive, &offer.header)?; + let mut expected_acknowledgement_key = [0_u8; 32]; + sender_context.export(b"spl-blob-ack-v1", &mut expected_acknowledgement_key)?; + + let prepared = prepare_authenticated_blob( + offer, + &encapsulated_key.to_bytes(), + &ciphertext, + &recipient_private_key, + &sender_public_key_spki_der, + &instance_id, + )?; + + assert_eq!(prepared.offer, offer); + assert_eq!(prepared.archive.metadata.day, "20260731"); + assert_eq!(prepared.archive.metadata.segment, "120001_3"); + assert_eq!(prepared.archive.entries.len(), 1); + assert_eq!(prepared.archive.entries[0].name, "entry.bin"); + assert_eq!(prepared.archive.entries[0].bytes, b"body"); + assert_eq!(prepared.acknowledgement_key, expected_acknowledgement_key); + Ok(()) + } + + fn valid_archive() -> Result, Box> { + let encoder = GzEncoder::new(Vec::new(), Compression::default()); + let mut builder = Builder::new(encoder); + append_regular( + &mut builder, + "blob.json", + br#"{"v":1,"day":"20260731","segment":"120001_3","host":"home.example","meta":{}}"#, + )?; + append_regular(&mut builder, "entry.bin", b"body")?; + let encoder = builder.into_inner()?; + Ok(encoder.finish()?) + } + + fn append_regular( + builder: &mut Builder>>, + name: &str, + bytes: &[u8], + ) -> Result<(), std::io::Error> { + let mut header = Header::new_gnu(); + header.set_size(bytes.len() as u64); + header.set_mode(0o600); + header.set_cksum(); + builder.append_data(&mut header, name, Cursor::new(bytes)) + } + + fn valid_offer(ciphertext_len: usize) -> Result> { + let mut header = [0_u8; OFFER_LEN]; + header[..11].copy_from_slice(&[ + b'S', b'B', b'O', b'1', 0x01, 0x00, 0x10, 0x00, 0x01, 0x00, 0x02, + ]); + header[11..43].copy_from_slice(&[0x22; 32]); + header[43..59].copy_from_slice(&[0x33; 16]); + header[59..].copy_from_slice(&u64::try_from(ciphertext_len)?.to_be_bytes()); + Ok(parse_offer(&header)?) + } + + fn hpke_private_key( + private_key: &P256Secret, + ) -> Result<::PrivateKey, hpke::HpkeError> { + ::PrivateKey::from_bytes(private_key.to_bytes().as_slice()) + } + + fn hpke_public_key( + private_key: &P256Secret, + ) -> Result<::PublicKey, hpke::HpkeError> { + let point = private_key.public_key().to_sec1_point(false); + ::PublicKey::from_bytes(point.as_bytes()) + } + + fn hpke_public_key_from_private_key( + private_key: &::PrivateKey, + ) -> Result<::PublicKey, hpke::HpkeError> { + Ok(::sk_to_pk(private_key)) + } +} diff --git a/core/crates/solstone-core-spl/src/authorized_client_ledger.rs b/core/crates/solstone-core-spl/src/authorized_client_ledger.rs new file mode 100644 index 000000000..0b70673c6 --- /dev/null +++ b/core/crates/solstone-core-spl/src/authorized_client_ledger.rs @@ -0,0 +1,210 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +//! Fresh, fail-closed browser authorization lookups. + +use std::fs; +use std::path::{Path, PathBuf}; + +use crate::{AuthorizedClients, BrowserUploadAuthorization, parse_authorized_clients}; + +/// An owned result from a browser authorization lookup. +/// +/// This owns the authorization material so a lookup can read a fresh file for +/// each call rather than retaining a last-good ledger in memory. +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum BrowserLedgerLookup { + /// The ledger could not be statted, read, or parsed as a JSON list. + LedgerUnavailable, + /// No browser authorization matched the sender fingerprint. + NotAuthorized, + /// A browser entry matched but lacks required upload material. + Incomplete, + /// The matching browser entry has the upload material required by SPL. + Authorized { + /// Hex-encoded sender SPKI. + pubkey_spki: String, + /// Observer-side upload handle. + observer_handle: String, + }, +} + +/// Read-only, fail-closed access to `/link/authorized_clients.json`. +/// +/// Delegated U2/C1s contract: every [`Self::lookup_browser`] call stats the +/// ledger and reads it afresh, so a Python writer's completed record is visible +/// in the same process without a restart or sleep. A missing, unreadable, +/// malformed, or non-list file has no retained authorization. This follows the +/// mtime reload and fail-closed behavior in +/// `solstone/think/link/auth.py:87-105,248-293`. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct AuthorizedClientLedger { + path: PathBuf, +} + +impl AuthorizedClientLedger { + /// Creates a reader rooted at a journal directory. + #[must_use] + pub fn new(journal_root: impl AsRef) -> Self { + Self { + path: journal_root + .as_ref() + .join("link") + .join("authorized_clients.json"), + } + } + + /// 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(), + }, + } + } + + 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) + .and_then(|metadata| metadata.modified()) + .is_err() + { + return AuthorizedClients::unavailable(); + } + + match fs::read(&self.path) { + Ok(bytes) => parse_authorized_clients(&bytes), + Err(_) => AuthorizedClients::unavailable(), + } + } +} + +#[cfg(test)] +mod tests { + use std::error::Error; + use std::fs; + use std::path::PathBuf; + use std::sync::atomic::{AtomicUsize, Ordering}; + use std::time::{SystemTime, UNIX_EPOCH}; + + use super::{AuthorizedClientLedger, BrowserLedgerLookup}; + + 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-authorized-client-ledger-{}-{timestamp}-{sequence}", + std::process::id() + )); + fs::create_dir_all(root.join("link"))?; + Ok(Self { root }) + } + + fn ledger_path(&self) -> PathBuf { + self.root.join("link").join("authorized_clients.json") + } + } + + impl Drop for TestJournal { + fn drop(&mut self) { + let _ = fs::remove_dir_all(&self.root); + } + } + + #[test] + fn existing_reader_observes_completed_python_write_without_sleep() -> Result<(), Box> + { + let journal = TestJournal::create()?; + let fingerprint = [0xa5; 32]; + let ledger = AuthorizedClientLedger::new(&journal.root); + + fs::write(journal.ledger_path(), browser_entry("browser", "null"))?; + assert_eq!( + ledger.lookup_browser(&fingerprint), + BrowserLedgerLookup::Incomplete + ); + + fs::write( + journal.ledger_path(), + browser_entry("browser", "\"observer-handle\""), + )?; + assert_eq!( + ledger.lookup_browser(&fingerprint), + BrowserLedgerLookup::Authorized { + pubkey_spki: "30aa".to_owned(), + observer_handle: "observer-handle".to_owned(), + } + ); + + Ok(()) + } + + #[test] + fn a_non_browser_kind_is_not_authorized() -> Result<(), Box> { + let journal = TestJournal::create()?; + let ledger = AuthorizedClientLedger::new(&journal.root); + + fs::write( + journal.ledger_path(), + browser_entry("cert", "\"observer-handle\""), + )?; + + assert_eq!( + ledger.lookup_browser(&[0xa5; 32]), + BrowserLedgerLookup::NotAuthorized + ); + Ok(()) + } + + #[test] + fn disappearance_and_invalid_contents_do_not_retain_authorization() -> Result<(), Box> + { + let journal = TestJournal::create()?; + let ledger = AuthorizedClientLedger::new(&journal.root); + let path = journal.ledger_path(); + + fs::write(&path, browser_entry("browser", "\"observer-handle\""))?; + assert!(matches!( + ledger.lookup_browser(&[0xa5; 32]), + BrowserLedgerLookup::Authorized { .. } + )); + + fs::remove_file(&path)?; + assert_eq!( + ledger.lookup_browser(&[0xa5; 32]), + BrowserLedgerLookup::LedgerUnavailable + ); + + fs::write(&path, "{}")?; + assert_eq!( + ledger.lookup_browser(&[0xa5; 32]), + BrowserLedgerLookup::LedgerUnavailable + ); + + Ok(()) + } + + fn browser_entry(kind: &str, observer_handle: &str) -> String { + let fingerprint = "a5".repeat(32); + format!( + "[{{\"fingerprint\":\"sha256:{fingerprint}\",\"kind\":\"{kind}\",\"pubkey_spki\":\"30aa\",\"observer_handle\":{observer_handle}}}]" + ) + } +} diff --git a/core/crates/solstone-core-spl/src/authorized_clients.rs b/core/crates/solstone-core-spl/src/authorized_clients.rs new file mode 100644 index 000000000..8c463b339 --- /dev/null +++ b/core/crates/solstone-core-spl/src/authorized_clients.rs @@ -0,0 +1,400 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +//! Pure parsing and lookup for the paired-client authorization ledger. + +use std::collections::BTreeMap; + +use serde_json::Value; + +/// A parsed entry from `link/authorized_clients.json`. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct ClientEntry { + /// The persisted client fingerprint, such as `sha256:abcd`. + pub fingerprint: String, + /// The home-assigned device label. + pub device_label: String, + /// The persisted pairing timestamp. + pub paired_at: String, + /// The home instance that paired the client. + pub instance_id: String, + /// The optional pairing provenance role. + pub role: String, + /// The last-seen timestamp, when present as a string. + pub last_seen_at: Option, + /// The local network display name, when present as a string. + pub network: Option, + /// The client-provided display label. + pub client_label: String, + /// The client kind, defaulting to `cert`. + pub kind: String, + /// The browser sender public key SPKI encoding, when present as a string. + pub pubkey_spki: Option, + /// The observer upload handle, when present as a string. + pub observer_handle: Option, +} + +/// Whether source bytes represented a usable ledger list. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum LedgerStatus { + /// The source was a JSON list, including an empty list. + Available, + /// The source was missing to the caller or did not decode to a JSON list. + Unavailable, +} + +/// Parsed paired-client entries plus source availability. +/// +/// An unavailable ledger intentionally has no entries. This preserves Python's +/// fail-closed behavior while allowing a caller to distinguish it from an +/// available ledger that simply has no matching client. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct AuthorizedClients { + status: LedgerStatus, + entries: BTreeMap, +} + +/// Browser upload authorization resolved from a raw sender fingerprint. +#[derive(Debug, Eq, PartialEq)] +pub enum BrowserUploadAuthorization<'a> { + /// The caller had no usable authorization ledger. + LedgerUnavailable, + /// The fingerprint has no browser entry in the ledger. + NotAuthorized, + /// A browser entry matched but lacks a non-empty upload key or observer handle. + Incomplete, + /// A browser entry contains the materials needed for an upload request. + Authorized { + /// Hex-encoded browser sender SPKI. + pubkey_spki: &'a str, + /// The observer upload handle. + observer_handle: &'a str, + }, +} + +impl AuthorizedClients { + /// Represents a missing ledger without introducing file I/O into this module. + #[must_use] + pub fn unavailable() -> Self { + Self { + status: LedgerStatus::Unavailable, + entries: BTreeMap::new(), + } + } + + /// Returns whether the source was a JSON ledger list. + #[must_use] + pub fn status(&self) -> LedgerStatus { + self.status + } + + /// Returns the parsed entries, which are empty for an unavailable ledger. + #[must_use] + pub fn entries(&self) -> &BTreeMap { + &self.entries + } + + /// Resolves browser upload authorization for a raw SHA-256 sender fingerprint. + #[must_use] + pub fn browser_upload_authorization( + &self, + sender_fingerprint: &[u8; 32], + ) -> BrowserUploadAuthorization<'_> { + if self.status == LedgerStatus::Unavailable { + return BrowserUploadAuthorization::LedgerUnavailable; + } + + let fingerprint = fingerprint_key(sender_fingerprint); + let Some(entry) = self.entries.get(&fingerprint) else { + return BrowserUploadAuthorization::NotAuthorized; + }; + if entry.kind != "browser" { + return BrowserUploadAuthorization::NotAuthorized; + } + let (Some(pubkey_spki), Some(observer_handle)) = ( + entry + .pubkey_spki + .as_deref() + .filter(|value| !value.is_empty()), + entry + .observer_handle + .as_deref() + .filter(|value| !value.is_empty()), + ) else { + return BrowserUploadAuthorization::Incomplete; + }; + + BrowserUploadAuthorization::Authorized { + pubkey_spki, + observer_handle, + } + } +} + +/// Parses JSON ledger bytes using the Python reader's fail-closed rules. +#[must_use] +pub fn parse_authorized_clients(input: &[u8]) -> AuthorizedClients { + let value: Value = match serde_json::from_slice(input) { + Ok(value) => value, + Err(_) => return AuthorizedClients::unavailable(), + }; + let Some(items) = value.as_array() else { + return AuthorizedClients::unavailable(); + }; + + let mut entries = BTreeMap::new(); + for item in items { + let Some(object) = item.as_object() else { + continue; + }; + let Some(fingerprint) = object.get("fingerprint").and_then(Value::as_str) else { + continue; + }; + + let entry = ClientEntry { + fingerprint: fingerprint.to_owned(), + device_label: python_str(object.get("device_label")), + paired_at: python_str(object.get("paired_at")), + instance_id: python_str(object.get("instance_id")), + role: string_or_empty(object.get("role")), + last_seen_at: string_or_none(object.get("last_seen_at")), + network: string_or_none(object.get("network")), + client_label: string_or_empty(object.get("client_label")), + kind: nonempty_string_or_cert(object.get("kind")), + pubkey_spki: string_or_none(object.get("pubkey_spki")), + observer_handle: string_or_none(object.get("observer_handle")), + }; + entries.insert(fingerprint.to_owned(), entry); + } + + AuthorizedClients { + status: LedgerStatus::Available, + entries, + } +} + +fn string_or_empty(value: Option<&Value>) -> String { + value + .and_then(Value::as_str) + .map_or_else(String::new, ToOwned::to_owned) +} + +fn string_or_none(value: Option<&Value>) -> Option { + value.and_then(Value::as_str).map(ToOwned::to_owned) +} + +fn nonempty_string_or_cert(value: Option<&Value>) -> String { + value + .and_then(Value::as_str) + .filter(|value| !value.is_empty()) + .map_or_else(|| "cert".to_owned(), ToOwned::to_owned) +} + +/// Reproduces Python's `str(item.get(field, ""))` for JSON values. +fn python_str(value: Option<&Value>) -> String { + value.map_or_else(String::new, python_repr_or_string) +} + +fn python_repr_or_string(value: &Value) -> String { + match value { + Value::String(value) => value.clone(), + _ => python_repr(value), + } +} + +fn python_repr(value: &Value) -> String { + match value { + Value::Null => "None".to_owned(), + Value::Bool(true) => "True".to_owned(), + Value::Bool(false) => "False".to_owned(), + Value::Number(number) => number.to_string(), + Value::String(value) => python_quoted_string(value), + Value::Array(values) => { + let values = values + .iter() + .map(python_repr) + .collect::>() + .join(", "); + format!("[{values}]") + } + Value::Object(values) => { + let values = values + .iter() + .map(|(key, value)| { + format!("{}: {}", python_quoted_string(key), python_repr(value)) + }) + .collect::>() + .join(", "); + format!("{{{values}}}") + } + } +} + +fn python_quoted_string(value: &str) -> String { + let quote = if value.contains('\'') && !value.contains('"') { + '"' + } else { + '\'' + }; + let mut quoted = String::with_capacity(value.len() + 2); + quoted.push(quote); + for character in value.chars() { + match character { + '\\' => quoted.push_str("\\\\"), + '\n' => quoted.push_str("\\n"), + '\r' => quoted.push_str("\\r"), + '\t' => quoted.push_str("\\t"), + '\u{08}' => quoted.push_str("\\x08"), + '\u{0c}' => quoted.push_str("\\x0c"), + character if character == quote => { + quoted.push('\\'); + quoted.push(character); + } + character if character.is_control() => { + let code_point = character as u32; + if code_point <= 0xff { + format_hex_escape(&mut quoted, code_point as u8); + } else { + quoted.push_str(&format!("\\u{code_point:04x}")); + } + } + character => quoted.push(character), + } + } + quoted.push(quote); + quoted +} + +fn format_hex_escape(output: &mut String, value: u8) { + const HEX: &[u8; 16] = b"0123456789abcdef"; + + output.push_str("\\x"); + output.push(char::from(HEX[usize::from(value >> 4)])); + output.push(char::from(HEX[usize::from(value & 0x0f)])); +} + +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 super::{ + AuthorizedClients, BrowserUploadAuthorization, LedgerStatus, parse_authorized_clients, + }; + + #[test] + fn malformed_and_non_list_ledgers_fail_closed_with_no_entries() { + for source in [b"{".as_slice(), b"{}".as_slice(), b"42".as_slice()] { + let ledger = parse_authorized_clients(source); + assert_eq!(ledger.status(), LedgerStatus::Unavailable); + assert!(ledger.entries().is_empty()); + } + } + + #[test] + fn non_objects_and_non_string_fingerprints_are_skipped() { + let ledger = parse_authorized_clients( + br#"["entry", 3, null, {}, {"fingerprint": 3}, {"fingerprint":"ok"}]"#, + ); + + assert_eq!(ledger.status(), LedgerStatus::Available); + assert_eq!(ledger.entries().len(), 1); + assert!(ledger.entries().contains_key("ok")); + } + + #[test] + fn fields_follow_python_coercions_and_defaults() { + let ledger = parse_authorized_clients( + br#"[{"fingerprint":"client","device_label":false,"paired_at":null,"instance_id":["x",true],"role":false,"last_seen_at":1,"network":false,"client_label":[],"kind":"","pubkey_spki":3,"observer_handle":null}]"#, + ); + let entry = ledger.entries().get("client"); + + assert!(entry.is_some()); + let entry = match entry { + Some(entry) => entry, + None => return, + }; + assert_eq!(entry.device_label, "False"); + assert_eq!(entry.paired_at, "None"); + assert_eq!(entry.instance_id, "['x', True]"); + assert_eq!(entry.role, ""); + assert_eq!(entry.last_seen_at, None); + assert_eq!(entry.network, None); + assert_eq!(entry.client_label, ""); + assert_eq!(entry.kind, "cert"); + assert_eq!(entry.pubkey_spki, None); + assert_eq!(entry.observer_handle, None); + } + + #[test] + fn duplicate_fingerprints_keep_the_later_entry() { + let ledger = parse_authorized_clients( + br#"[{"fingerprint":"same","device_label":"first"},{"fingerprint":"same","device_label":"later"}]"#, + ); + let entry = ledger.entries().get("same"); + + assert_eq!(ledger.entries().len(), 1); + assert_eq!( + entry.map(|entry| entry.device_label.as_str()), + Some("later") + ); + } + + #[test] + fn raw_sender_fingerprint_resolves_browser_upload_authorization() { + let sender_fingerprint = [0xab; 32]; + let fingerprint = "ab".repeat(32); + let source = format!( + "[{{\"fingerprint\":\"sha256:{fingerprint}\",\"kind\":\"browser\",\"pubkey_spki\":\"30aa\",\"observer_handle\":\"observer-a\"}}]" + ); + let ledger = parse_authorized_clients(source.as_bytes()); + + assert_eq!( + ledger.browser_upload_authorization(&sender_fingerprint), + BrowserUploadAuthorization::Authorized { + pubkey_spki: "30aa", + observer_handle: "observer-a", + } + ); + } + + #[test] + fn incomplete_browser_entries_are_rejected_separately() { + let sender_fingerprint = [0xcd; 32]; + let fingerprint = "cd".repeat(32); + let source = format!( + "[{{\"fingerprint\":\"sha256:{fingerprint}\",\"kind\":\"browser\",\"pubkey_spki\":\"\",\"observer_handle\":\"observer-a\"}}]" + ); + let ledger = parse_authorized_clients(source.as_bytes()); + + assert_eq!( + ledger.browser_upload_authorization(&sender_fingerprint), + BrowserUploadAuthorization::Incomplete + ); + } + + #[test] + fn unavailable_ledger_is_not_conflated_with_no_matching_client() { + let unavailable = AuthorizedClients::unavailable(); + let available = parse_authorized_clients(br#"[]"#); + let sender_fingerprint = [0; 32]; + + assert_eq!( + unavailable.browser_upload_authorization(&sender_fingerprint), + BrowserUploadAuthorization::LedgerUnavailable + ); + assert_eq!( + available.browser_upload_authorization(&sender_fingerprint), + BrowserUploadAuthorization::NotAuthorized + ); + } +} diff --git a/core/crates/solstone-core-spl/src/blob_archive.rs b/core/crates/solstone-core-spl/src/blob_archive.rs new file mode 100644 index 000000000..971a55e81 --- /dev/null +++ b/core/crates/solstone-core-spl/src/blob_archive.rs @@ -0,0 +1,516 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +//! Validation for the self-contained archive carried by an SPL blob offer. + +use std::io::{Cursor, Read}; + +use flate2::read::GzDecoder; +use regex::Regex; +use serde_json::{Map, Value}; +use tar::{Archive, EntryType}; +use thiserror::Error; + +const MAX_ENTRIES: usize = 64; +const MAX_ENTRY_BYTES: u64 = 16 * 1024 * 1024; +const MAX_TOTAL_BYTES: u64 = 64 * 1024 * 1024; +const METADATA_NAME: &str = "blob.json"; + +/// A regular payload file that passed archive validation. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct BlobArchiveEntry { + /// The single-component name stored in the archive. + pub name: String, + /// The complete in-memory content of the validated file. + pub bytes: Vec, +} + +/// Validated metadata from `blob.json`. +#[derive(Clone, Debug, PartialEq)] +pub struct BlobArchiveMetadata { + /// Archive format version, always `1` for a parsed archive. + pub version: u64, + /// Eight-digit archive day. + pub day: String, + /// Segment identifier in `HHMMSS_N` form. + pub segment: String, + /// The archive's nonempty source host. + pub host: String, + /// Caller-defined metadata, constrained to a JSON object. + pub meta: Map, +} + +/// A fully validated gzip-compressed tar archive. +#[derive(Clone, Debug, PartialEq)] +pub struct ValidatedBlobArchive { + /// Parsed metadata formerly held in `blob.json`. + pub metadata: BlobArchiveMetadata, + /// Every regular file other than `blob.json`. + pub entries: Vec, +} + +/// Reasons an archive cannot safely be made available to an SPL receiver. +#[derive(Debug, Error)] +pub enum BlobArchiveError { + /// The gzip stream or tar reader returned an I/O error. + #[error("archive stream error: {0}")] + Stream(#[from] std::io::Error), + + /// More than the permitted number of tar entries was present. + #[error("archive has more than {MAX_ENTRIES} entries")] + TooManyEntries, + + /// A tar entry was not a regular file. + #[error("archive entry {name:?} has unsupported type {entry_type:?}")] + UnsupportedEntryType { + /// The path reported by the tar header, when it was representable. + name: String, + /// The tar entry type reported by the header. + entry_type: EntryType, + }, + + /// A tar name was not valid UTF-8. + #[error("archive entry name is not valid UTF-8")] + NonUtf8EntryName, + + /// A tar name violated the single-file-name contract. + #[error("archive entry name {name:?} is not allowed")] + InvalidEntryName { + /// The rejected archive name. + name: String, + }, + + /// One file exceeded the per-file size limit. + #[error("archive entry {name:?} exceeds {MAX_ENTRY_BYTES} bytes")] + EntryTooLarge { + /// The oversized entry's name. + name: String, + /// The size declared by the tar header. + size: u64, + }, + + /// The declared regular-file sizes exceeded the total size limit. + #[error("archive regular files exceed {MAX_TOTAL_BYTES} bytes")] + TotalTooLarge, + + /// The archive included more than one metadata file. + #[error("archive has more than one {METADATA_NAME}")] + DuplicateMetadata, + + /// The archive did not include its required metadata file. + #[error("archive is missing {METADATA_NAME}")] + MissingMetadata, + + /// The archive has no regular payload file besides its metadata. + #[error("archive has no payload files")] + MissingPayload, + + /// `blob.json` was not valid JSON. + #[error("{METADATA_NAME} is not valid JSON: {0}")] + InvalidMetadataJson(serde_json::Error), + + /// `blob.json` was not a JSON object. + #[error("{METADATA_NAME} must be a JSON object")] + MetadataNotObject, + + /// The metadata version was absent or not numeric version one. + #[error("{METADATA_NAME}.v must equal 1")] + InvalidVersion, + + /// The metadata day was absent or did not use eight ASCII digits. + #[error("{METADATA_NAME}.day must contain exactly eight digits")] + InvalidDay, + + /// The metadata segment was absent or did not use `HHMMSS_N` form. + #[error("{METADATA_NAME}.segment must use HHMMSS_N form")] + InvalidSegment, + + /// The metadata host was absent, not a string, or empty. + #[error("{METADATA_NAME}.host must be a nonempty string")] + InvalidHost, + + /// The metadata meta field was absent or not a JSON object. + #[error("{METADATA_NAME}.meta must be a JSON object")] + InvalidMeta, +} + +/// Parse and validate an SPL blob archive without writing any file to disk. +pub fn parse_blob_archive(input: &[u8]) -> Result { + let decoder = GzDecoder::new(Cursor::new(input)); + let mut archive = Archive::new(decoder); + let mut entry_count = 0_usize; + let mut total_bytes = 0_u64; + let mut metadata_bytes = None; + let mut payload_entries = Vec::new(); + + for entry_result in archive.entries()? { + entry_count = entry_count + .checked_add(1) + .ok_or(BlobArchiveError::TooManyEntries)?; + if entry_count > MAX_ENTRIES { + return Err(BlobArchiveError::TooManyEntries); + } + + let mut entry = entry_result?; + let path = entry.path()?; + let name = path + .to_str() + .ok_or(BlobArchiveError::NonUtf8EntryName)? + .to_owned(); + validate_entry_name(&name)?; + + let entry_type = entry.header().entry_type(); + if !entry_type.is_file() { + return Err(BlobArchiveError::UnsupportedEntryType { name, entry_type }); + } + + let size = entry.size(); + if size > MAX_ENTRY_BYTES { + return Err(BlobArchiveError::EntryTooLarge { name, size }); + } + total_bytes = total_bytes + .checked_add(size) + .filter(|total| *total <= MAX_TOTAL_BYTES) + .ok_or(BlobArchiveError::TotalTooLarge)?; + + let mut bytes = Vec::with_capacity(size as usize); + entry.read_to_end(&mut bytes)?; + if name == METADATA_NAME { + if metadata_bytes.replace(bytes).is_some() { + return Err(BlobArchiveError::DuplicateMetadata); + } + } else { + payload_entries.push(BlobArchiveEntry { name, bytes }); + } + } + + let metadata_bytes = metadata_bytes.ok_or(BlobArchiveError::MissingMetadata)?; + if payload_entries.is_empty() { + return Err(BlobArchiveError::MissingPayload); + } + + let metadata = parse_metadata(&metadata_bytes)?; + Ok(ValidatedBlobArchive { + metadata, + entries: payload_entries, + }) +} + +fn validate_entry_name(name: &str) -> Result<(), BlobArchiveError> { + if name.is_empty() + || matches!(name, "." | "..") + || name.starts_with('/') + || name.contains('\\') + || name.contains('/') + { + return Err(BlobArchiveError::InvalidEntryName { + name: name.to_owned(), + }); + } + Ok(()) +} + +fn parse_metadata(bytes: &[u8]) -> Result { + let value: Value = + serde_json::from_slice(bytes).map_err(BlobArchiveError::InvalidMetadataJson)?; + let object = value + .as_object() + .ok_or(BlobArchiveError::MetadataNotObject)?; + + let version = object + .get("v") + .and_then(Value::as_f64) + .filter(|version| *version == 1.0) + .map(|_| 1) + .ok_or(BlobArchiveError::InvalidVersion)?; + let day = object + .get("day") + .and_then(Value::as_str) + .filter(|day| matches_day(day)) + .ok_or(BlobArchiveError::InvalidDay)? + .to_owned(); + let segment = object + .get("segment") + .and_then(Value::as_str) + .filter(|segment| matches_segment(segment)) + .ok_or(BlobArchiveError::InvalidSegment)? + .to_owned(); + let host = object + .get("host") + .and_then(Value::as_str) + .filter(|host| !host.is_empty()) + .ok_or(BlobArchiveError::InvalidHost)? + .to_owned(); + let meta = object + .get("meta") + .and_then(Value::as_object) + .ok_or(BlobArchiveError::InvalidMeta)? + .clone(); + + Ok(BlobArchiveMetadata { + version, + day, + segment, + host, + meta, + }) +} + +fn matches_day(value: &str) -> bool { + matches_pattern(r"^\d{8}$", value) +} + +fn matches_segment(value: &str) -> bool { + matches_pattern(r"^\d{6}_\d+$", value) +} + +fn matches_pattern(pattern: &str, value: &str) -> bool { + Regex::new(pattern).is_ok_and(|regex| regex.is_match(value)) +} + +#[cfg(test)] +mod tests { + use std::io::Cursor; + + use flate2::{Compression, write::GzEncoder}; + use tar::{Builder, EntryType, Header}; + + use super::{BlobArchiveError, parse_blob_archive, validate_entry_name}; + + const METADATA: &str = + r#"{"v":1,"day":"20260731","segment":"120001_3","host":"home.example","meta":{"a":true}}"#; + + fn regular_archive(entries: &[(&str, &[u8])]) -> Result, std::io::Error> { + let encoder = GzEncoder::new(Vec::new(), Compression::default()); + let mut builder = Builder::new(encoder); + for (name, bytes) in entries { + let mut header = Header::new_gnu(); + header.set_size(bytes.len() as u64); + header.set_mode(0o600); + header.set_cksum(); + builder.append_data(&mut header, *name, Cursor::new(*bytes))?; + } + let encoder = builder.into_inner()?; + encoder.finish() + } + + fn valid_archive() -> Result, std::io::Error> { + regular_archive(&[("blob.json", METADATA.as_bytes()), ("entry.bin", b"body")]) + } + + fn archive_with_metadata(metadata: &[u8]) -> Result, std::io::Error> { + regular_archive(&[("blob.json", metadata), ("entry.bin", b"body")]) + } + + #[test] + fn returns_metadata_and_payload_entries() -> Result<(), Box> { + let parsed = parse_blob_archive(&valid_archive()?)?; + + assert_eq!(parsed.metadata.version, 1); + assert_eq!(parsed.metadata.day, "20260731"); + assert_eq!(parsed.entries.len(), 1); + assert_eq!(parsed.entries[0].name, "entry.bin"); + assert_eq!(parsed.entries[0].bytes, b"body"); + Ok(()) + } + + #[test] + fn accepts_float_one_for_metadata_version() -> Result<(), Box> { + let metadata = + br#"{"v":1.0,"day":"20260731","segment":"120001_3","host":"home.example","meta":{}}"#; + let parsed = parse_blob_archive(&archive_with_metadata(metadata)?)?; + + assert_eq!(parsed.metadata.version, 1); + Ok(()) + } + + #[test] + fn accepts_unicode_decimal_digits_in_metadata() -> Result<(), Box> { + let metadata = + r#"{"v":1,"day":"١٢٣٤٥٦٧٨","segment":"١٢٣٤٥٦_٧","host":"home.example","meta":{}}"#; + let parsed = parse_blob_archive(&archive_with_metadata(metadata.as_bytes())?)?; + + assert_eq!(parsed.metadata.day, "١٢٣٤٥٦٧٨"); + assert_eq!(parsed.metadata.segment, "١٢٣٤٥٦_٧"); + Ok(()) + } + + #[test] + fn rejects_more_than_64_entries() -> Result<(), Box> { + let mut names = Vec::new(); + names.push(("blob.json".to_owned(), METADATA.as_bytes().to_vec())); + for number in 0..64 { + names.push((format!("{number}.bin"), Vec::new())); + } + let references: Vec<(&str, &[u8])> = names + .iter() + .map(|(name, bytes)| (name.as_str(), bytes.as_slice())) + .collect(); + let archive = regular_archive(&references)?; + + assert!(matches!( + parse_blob_archive(&archive), + Err(BlobArchiveError::TooManyEntries) + )); + Ok(()) + } + + #[test] + fn rejects_non_regular_entries() -> Result<(), Box> { + let encoder = GzEncoder::new(Vec::new(), Compression::default()); + let mut builder = Builder::new(encoder); + let mut header = Header::new_gnu(); + header.set_entry_type(EntryType::Directory); + header.set_size(0); + header.set_mode(0o700); + header.set_cksum(); + builder.append_data(&mut header, "directory", Cursor::new(Vec::::new()))?; + let encoder = builder.into_inner()?; + let archive = encoder.finish()?; + + assert!(matches!( + parse_blob_archive(&archive), + Err(BlobArchiveError::UnsupportedEntryType { .. }) + )); + Ok(()) + } + + #[test] + fn rejects_disallowed_names() { + for name in ["", ".", "..", "/absolute", "back\\slash", "nested/name"] { + assert!(matches!( + validate_entry_name(name), + Err(BlobArchiveError::InvalidEntryName { .. }) + )); + } + } + + #[test] + fn rejects_oversized_regular_file() -> Result<(), Box> { + let large = vec![0_u8; (16 * 1024 * 1024) + 1]; + let archive = regular_archive(&[("blob.json", METADATA.as_bytes()), ("large", &large)])?; + + assert!(matches!( + parse_blob_archive(&archive), + Err(BlobArchiveError::EntryTooLarge { .. }) + )); + Ok(()) + } + + #[test] + fn rejects_total_oversized_regular_files() -> Result<(), Box> { + let first = vec![0_u8; 16 * 1024 * 1024]; + let second = vec![0_u8; 16 * 1024 * 1024]; + let third = vec![0_u8; 16 * 1024 * 1024]; + let fourth = vec![0_u8; 16 * 1024 * 1024]; + let extra = vec![0_u8; 1]; + let archive = regular_archive(&[ + ("blob.json", METADATA.as_bytes()), + ("one", &first), + ("two", &second), + ("three", &third), + ("four", &fourth), + ("five", &extra), + ])?; + + assert!(matches!( + parse_blob_archive(&archive), + Err(BlobArchiveError::TotalTooLarge) + )); + Ok(()) + } + + #[test] + fn requires_exactly_one_metadata_file() -> Result<(), Box> { + let missing = regular_archive(&[("entry.bin", b"body")])?; + assert!(matches!( + parse_blob_archive(&missing), + Err(BlobArchiveError::MissingMetadata) + )); + + let duplicate = regular_archive(&[ + ("blob.json", METADATA.as_bytes()), + ("blob.json", METADATA.as_bytes()), + ("entry.bin", b"body"), + ])?; + assert!(matches!( + parse_blob_archive(&duplicate), + Err(BlobArchiveError::DuplicateMetadata) + )); + Ok(()) + } + + #[test] + fn requires_payload_file() -> Result<(), Box> { + let archive = regular_archive(&[("blob.json", METADATA.as_bytes())])?; + + assert!(matches!( + parse_blob_archive(&archive), + Err(BlobArchiveError::MissingPayload) + )); + Ok(()) + } + + #[test] + fn validates_metadata_shape() -> Result<(), Box> { + let cases = [ + (b"not json".as_slice(), BlobArchiveErrorKind::InvalidJson), + (b"[]".as_slice(), BlobArchiveErrorKind::NotObject), + ( + br#"{"v":2,"day":"20260731","segment":"120001_3","host":"home","meta":{}}"# + .as_slice(), + BlobArchiveErrorKind::Version, + ), + ( + br#"{"v":1,"day":"bad","segment":"120001_3","host":"home","meta":{}}"#.as_slice(), + BlobArchiveErrorKind::Day, + ), + ( + br#"{"v":1,"day":"20260731","segment":"no","host":"home","meta":{}}"#.as_slice(), + BlobArchiveErrorKind::Segment, + ), + ( + br#"{"v":1,"day":"20260731","segment":"120001_3","host":"","meta":{}}"#.as_slice(), + BlobArchiveErrorKind::Host, + ), + ( + br#"{"v":1,"day":"20260731","segment":"120001_3","host":"home","meta":[]}"# + .as_slice(), + BlobArchiveErrorKind::Meta, + ), + ]; + for (metadata, kind) in cases { + let archive = regular_archive(&[("blob.json", metadata), ("entry.bin", b"body")])?; + let result = parse_blob_archive(&archive); + assert!(kind.matches(result)); + } + Ok(()) + } + + enum BlobArchiveErrorKind { + InvalidJson, + NotObject, + Version, + Day, + Segment, + Host, + Meta, + } + + impl BlobArchiveErrorKind { + fn matches(&self, result: Result) -> bool { + matches!( + (self, result), + ( + Self::InvalidJson, + Err(BlobArchiveError::InvalidMetadataJson(_)) + ) | (Self::NotObject, Err(BlobArchiveError::MetadataNotObject)) + | (Self::Version, Err(BlobArchiveError::InvalidVersion)) + | (Self::Day, Err(BlobArchiveError::InvalidDay)) + | (Self::Segment, Err(BlobArchiveError::InvalidSegment)) + | (Self::Host, Err(BlobArchiveError::InvalidHost)) + | (Self::Meta, Err(BlobArchiveError::InvalidMeta)) + ) + } + } +} diff --git a/core/crates/solstone-core-spl/src/blob_content_type.rs b/core/crates/solstone-core-spl/src/blob_content_type.rs new file mode 100644 index 000000000..4a9a144fa --- /dev/null +++ b/core/crates/solstone-core-spl/src/blob_content_type.rs @@ -0,0 +1,49 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +//! Content-type selection for decoded SPL blob entries. + +/// Returns the media type assigned by the legacy SPL blob receiver. +/// +/// Suffix matching is deliberately case-sensitive to preserve the existing +/// service behaviour. +pub fn blob_content_type(name: &str) -> &'static str { + if name.ends_with(".jsonl") { + "application/jsonl" + } else if name.ends_with(".json") { + "application/json" + } else { + "application/octet-stream" + } +} + +#[cfg(test)] +mod tests { + use super::blob_content_type; + + #[test] + fn selects_jsonl_before_json() { + assert_eq!(blob_content_type("journal.jsonl"), "application/jsonl"); + } + + #[test] + fn selects_json_for_the_json_suffix() { + assert_eq!(blob_content_type("blob.json"), "application/json"); + } + + #[test] + fn preserves_case_sensitive_suffix_matching() { + assert_eq!(blob_content_type("blob.JSONL"), "application/octet-stream"); + assert_eq!(blob_content_type("blob.Json"), "application/octet-stream"); + } + + #[test] + fn assigns_octet_stream_to_arbitrary_names() { + assert_eq!( + blob_content_type("archive.tar.gz"), + "application/octet-stream" + ); + assert_eq!(blob_content_type("json"), "application/octet-stream"); + assert_eq!(blob_content_type(""), "application/octet-stream"); + } +} -- 2.51.2