diff --git a/core/archipelago/src/media_registration.rs b/core/archipelago/src/media_registration.rs new file mode 100644 index 00000000..bac13fea --- /dev/null +++ b/core/archipelago/src/media_registration.rs @@ -0,0 +1,1080 @@ +//! Private, durable preparation for an explicitly authorized IndeeHub Cloud selection. +//! +//! This is NOT an RPC authorization boundary or a content-serving registration. +//! Call only on a blocking worker after authenticating the installed app, fetching +//! its intent through a trusted channel, and obtaining the operator's consent for +//! this exact selection and terms. Do not return the receipt to the app until the +//! caller has durably linked its snapshot to authenticated serving/entitlements. +//! Nothing here advertises content, changes a share policy, or creates an identity. +use anyhow::{Context, Result}; +use ed25519_dalek::{Signature, Signer}; +use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; +use std::ffi::CString; +use std::fs::File; +use std::io::{Read, Write}; +use std::os::fd::{AsRawFd, FromRawFd}; +use std::os::unix::ffi::OsStrExt; +use std::os::unix::fs::{MetadataExt, PermissionsExt}; +use std::path::{Component, Path, PathBuf}; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::time::{Duration, Instant}; + +const DOMAIN: &str = "archipelago.indeehub.media-registration.v1"; +const MAX_SAFE_INTEGER: u64 = 9_007_199_254_740_991; +const MAX_RECORD_BYTES: u64 = 64 * 1024; +const PRIVATE_DIRECTORY: &str = "media-registration"; + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase", deny_unknown_fields)] +pub struct Intent { + pub version: u8, + pub request_id: String, + pub nonce: String, + pub app_audience: String, + pub node_did: String, + pub producer: String, + pub project_id: String, + pub price_sats: u64, + pub viewing_seconds: u64, + pub created_at: u64, + pub expires_at: u64, +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase", deny_unknown_fields)] +pub struct Receipt { + pub version: u8, + pub request_id: String, + pub nonce: String, + pub app_audience: String, + pub node_did: String, + pub producer: String, + pub project_id: String, + pub price_sats: u64, + pub viewing_seconds: u64, + pub expires_at: u64, + pub content_id: String, + pub sha256: String, + /// Decimal string, preserving the full uint64 range in JavaScript. + pub size_bytes: String, + pub payment_methods: Vec, + pub issued_at: u64, + pub signature: String, +} + +impl Receipt { + /// Exactly the ordered UTF-8 JSON array consumed by IndeeHub v1. + pub fn preimage(&self) -> Result> { + Ok(serde_json::to_vec(&serde_json::json!([ + DOMAIN, + self.request_id, + self.nonce, + self.app_audience, + self.node_did, + self.producer, + self.project_id, + self.content_id, + self.sha256, + self.size_bytes, + self.price_sats, + self.viewing_seconds, + self.payment_methods, + self.issued_at, + self.expires_at + ]))?) + } +} + +/// Must come from authenticated installation state, never the request body. +#[derive(Clone, Debug)] +pub struct InstallationPin { + pub node_did: String, + pub app_audience: String, +} + +/// The caller must have obtained explicit operator consent for all these fields. +/// This struct carries assertions; constructing it does not authenticate them. +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct AuthorizedSelection { + pub relative_path: PathBuf, + /// Sorted, distinct, and restricted to methods this node can actually receive. + pub payment_methods: Vec, +} + +pub struct Limits<'a> { + pub max_bytes: u64, + pub cancelled: &'a AtomicBool, +} + +/// Private preparation only. Receipt delivery requires caller-owned durable +/// serving/entitlement linkage, described in this module's contract. +#[derive(Debug)] +pub struct PreparedRegistration { + pub receipt: Receipt, + pub snapshot_path: PathBuf, + /// Read-only descriptor for the verified immutable snapshot. + pub snapshot: File, +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct SourceStamp { + device: u64, + inode: u64, + size: u64, + modified_seconds: i64, + modified_nanos: i64, + changed_seconds: i64, + changed_nanos: i64, +} +impl SourceStamp { + fn read(file: &File) -> Result { + let m = file.metadata()?; + anyhow::ensure!( + m.is_file(), + "The selected Cloud item must be a regular file" + ); + Ok(Self { + device: m.dev(), + inode: m.ino(), + size: m.len(), + modified_seconds: m.mtime(), + modified_nanos: m.mtime_nsec(), + changed_seconds: m.ctime(), + changed_nanos: m.ctime_nsec(), + }) + } +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct Binding { + intent: Intent, + selection: AuthorizedSelection, + cloud_root: PathBuf, +} +#[derive(Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct Operation { + version: u8, + binding: Binding, + source: SourceStamp, + issued_at: u64, +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct SnapshotRecord { + sha256: String, + size: u64, +} + +fn ascii_id(value: &str) -> bool { + !value.is_empty() + && value.len() <= 128 + && value + .bytes() + .all(|v| v.is_ascii_alphanumeric() || v == b'_' || v == b'-') +} +fn lower_hex(value: &str, length: usize) -> bool { + value.len() == length + && value + .bytes() + .all(|v| v.is_ascii_digit() || (b'a'..=b'f').contains(&v)) +} +fn validate( + intent: &Intent, + selection: &AuthorizedSelection, + pin: &InstallationPin, + identity: &crate::identity::NodeIdentity, + now: u64, +) -> Result<()> { + let id = + uuid::Uuid::parse_str(&intent.request_id).context("Invalid registration request ID")?; + anyhow::ensure!( + id.get_version_num() == 4 + && id.get_variant() == uuid::Variant::RFC4122 + && id.to_string() == intent.request_id, + "Registration request ID must be a canonical v4 UUID" + ); + anyhow::ensure!( + intent.version == 1 + && lower_hex(&intent.nonce, 64) + && lower_hex(&intent.producer, 64) + && ascii_id(&intent.project_id) + && ascii_id(&pin.app_audience) + && intent.app_audience == pin.app_audience + && intent.node_did == pin.node_did + && pin.node_did == identity.did_key()?, + "Registration identity or installation binding mismatch" + ); + anyhow::ensure!( + intent.price_sats <= MAX_SAFE_INTEGER + && (1..=31_536_000).contains(&intent.viewing_seconds) + && intent.created_at <= MAX_SAFE_INTEGER + && intent.expires_at <= MAX_SAFE_INTEGER + && now <= MAX_SAFE_INTEGER + && intent.created_at <= now.saturating_add(30) + && intent.expires_at > intent.created_at + && intent.expires_at - intent.created_at <= 600, + "Invalid registration terms or timestamps" + ); + anyhow::ensure!( + !selection.relative_path.as_os_str().is_empty() + && selection + .relative_path + .components() + .all(|part| matches!(part, Component::Normal(_))) + && selection.relative_path.as_os_str().as_bytes().len() <= 4096, + "Cloud selection must be a relative path without traversal" + ); + let methods = &selection.payment_methods; + anyhow::ensure!( + !methods.is_empty() + && methods.len() <= 4 + && methods.iter().all(|m| matches!( + m.as_str(), + "cashu" | "fedimint" | "lightning" | "lightning-cashu" + )) + && methods.windows(2).all(|m| m[0] < m[1]), + "Invalid receiving methods" + ); + Ok(()) +} + +fn c_path(path: &Path) -> Result { + CString::new(path.as_os_str().as_bytes()).context("Invalid file path") +} +fn fd_result(fd: libc::c_int) -> Result { + if fd < 0 { + return Err(std::io::Error::last_os_error().into()); + } + // SAFETY: the successful syscall returned a newly owned descriptor. + Ok(unsafe { File::from_raw_fd(fd) }) +} +fn open_at(dir: &File, name: &str, flags: libc::c_int, mode: libc::mode_t) -> Result { + let name = CString::new(name)?; + // SAFETY: descriptor and NUL-terminated name remain alive through the call. + fd_result(unsafe { + libc::openat( + dir.as_raw_fd(), + name.as_ptr(), + flags | libc::O_CLOEXEC | libc::O_NOFOLLOW, + mode, + ) + }) +} +fn open_directory(path: &Path) -> Result { + let path = c_path(path)?; + // SAFETY: path is a valid NUL-terminated string. + fd_result(unsafe { + libc::open( + path.as_ptr(), + libc::O_RDONLY | libc::O_DIRECTORY | libc::O_CLOEXEC | libc::O_NOFOLLOW, + ) + }) +} +fn private_directory(parent: &File, name: &str) -> Result { + let c_name = CString::new(name)?; + // SAFETY: valid directory descriptor and string; no existing data is replaced. + let result = unsafe { libc::mkdirat(parent.as_raw_fd(), c_name.as_ptr(), 0o700) }; + if result < 0 { + let error = std::io::Error::last_os_error(); + if error.kind() != std::io::ErrorKind::AlreadyExists { + return Err(error.into()); + } + } + let dir = open_at(parent, name, libc::O_RDONLY | libc::O_DIRECTORY, 0)?; + let metadata = dir.metadata()?; + // SAFETY: geteuid has no pointer arguments or side effects. + anyhow::ensure!( + metadata.uid() == unsafe { libc::geteuid() } && metadata.mode() & 0o077 == 0, + "Registration storage must be private to the node service" + ); + parent.sync_all()?; + Ok(dir) +} + +#[cfg(target_os = "linux")] +fn open_cloud_file(root: &File, relative: &Path) -> Result { + #[repr(C)] + struct OpenHow { + flags: u64, + mode: u64, + resolve: u64, + } + let name = c_path(relative)?; + let how = OpenHow { + flags: (libc::O_PATH | libc::O_CLOEXEC | libc::O_NOFOLLOW) as u64, + mode: 0, + resolve: 0x08 | 0x04, + }; // RESOLVE_BENEATH | RESOLVE_NO_SYMLINKS + // SAFETY: openat2 receives a held directory descriptor, valid name and a + // correctly sized immutable open_how. Fail closed on unsupported kernels. + let path_fd = fd_result(unsafe { + libc::syscall( + libc::SYS_openat2, + root.as_raw_fd(), + name.as_ptr(), + &how as *const OpenHow, + std::mem::size_of::(), + ) as libc::c_int + }) + .context("Could not safely open the Cloud selection (Linux openat2 required)")?; + let expected = SourceStamp::read(&path_fd)?; + // Reopen only our already-held regular-file descriptor, not its mutable path. + // O_PATH lets us reject devices/FIFOs before invoking their open operations. + let descriptor = CString::new(format!("/proc/self/fd/{}", path_fd.as_raw_fd()))?; + let file = fd_result(unsafe { + libc::open( + descriptor.as_ptr(), + libc::O_RDONLY | libc::O_CLOEXEC | libc::O_NONBLOCK, + ) + })?; + anyhow::ensure!( + SourceStamp::read(&file)? == expected, + "Cloud selection changed while opening" + ); + Ok(file) +} +#[cfg(not(target_os = "linux"))] +fn open_cloud_file(_root: &File, _relative: &Path) -> Result { + anyhow::bail!("Safe Cloud registration currently requires Linux openat2") +} + +fn cancelled(limits: &Limits<'_>) -> Result<()> { + anyhow::ensure!( + !limits.cancelled.load(Ordering::Relaxed), + "Media registration cancelled" + ); + Ok(()) +} +fn lock_operation(dir: &File, limits: &Limits<'_>, deadline: Instant) -> Result<()> { + loop { + cancelled(limits)?; + anyhow::ensure!( + Instant::now() < deadline, + "Timed out waiting for this registration operation" + ); + // SAFETY: flock borrows the directory descriptor; closing it releases lock. + if unsafe { libc::flock(dir.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) } == 0 { + return Ok(()); + } + let error = std::io::Error::last_os_error(); + if error.kind() != std::io::ErrorKind::WouldBlock + && error.kind() != std::io::ErrorKind::Interrupted + { + return Err(error.into()); + } + std::thread::sleep(Duration::from_millis(20)); + } +} +fn read_record(dir: &File, name: &str) -> Result> { + let file = match open_at(dir, name, libc::O_RDONLY | libc::O_NONBLOCK, 0) { + Ok(file) => file, + Err(error) + if error + .downcast_ref::() + .is_some_and(|e| e.kind() == std::io::ErrorKind::NotFound) => + { + return Ok(None) + } + Err(error) => return Err(error), + }; + let m = file.metadata()?; + anyhow::ensure!( + m.is_file() && m.len() <= MAX_RECORD_BYTES && m.mode() & 0o077 == 0, + "Invalid private registration record" + ); + let mut bytes = Vec::new(); + file.take(MAX_RECORD_BYTES + 1).read_to_end(&mut bytes)?; + anyhow::ensure!( + bytes.len() as u64 <= MAX_RECORD_BYTES, + "Registration record is too large" + ); + Ok(Some(serde_json::from_slice(&bytes).context( + "Damaged registration record; preserve it for recovery", + )?)) +} +fn temporary(dir: &File) -> Result<(String, File)> { + let name = format!("pending-{}", uuid::Uuid::new_v4()); + Ok(( + name.clone(), + open_at( + dir, + &name, + libc::O_WRONLY | libc::O_CREAT | libc::O_EXCL, + 0o600, + )?, + )) +} +fn publish_file(dir: &File, temporary: &str, final_name: &str) -> Result<()> { + let from = CString::new(temporary)?; + let to = CString::new(final_name)?; + // linkat publishes without replacing an existing destination. Both names + // are generated locally inside the locked private operation directory. + if unsafe { + libc::linkat( + dir.as_raw_fd(), + from.as_ptr(), + dir.as_raw_fd(), + to.as_ptr(), + 0, + ) + } < 0 + { + return Err(std::io::Error::last_os_error().into()); + } + dir.sync_all()?; + if unsafe { libc::unlinkat(dir.as_raw_fd(), from.as_ptr(), 0) } < 0 { + return Err(std::io::Error::last_os_error().into()); + } + dir.sync_all()?; + Ok(()) +} +fn save_record(dir: &File, name: &str, value: &T) -> Result<()> { + let bytes = serde_json::to_vec(value)?; + anyhow::ensure!( + bytes.len() as u64 <= MAX_RECORD_BYTES, + "Registration record is too large" + ); + let (temporary, mut file) = temporary(dir)?; + file.write_all(&bytes)?; + file.sync_all()?; + publish_file(dir, &temporary, name) +} + +fn hash_file( + file: &mut File, + mut output: Option<&mut File>, + limits: &Limits<'_>, + progress: &mut impl FnMut(u64) -> Result<()>, +) -> Result<(String, u64)> { + let mut hash = Sha256::new(); + let mut size = 0u64; + let mut chunk = [0u8; 64 * 1024]; + loop { + cancelled(limits)?; + let count = file.read(&mut chunk)?; + if count == 0 { + break; + } + size = size + .checked_add(count as u64) + .context("Media size overflow")?; + anyhow::ensure!( + size <= limits.max_bytes, + "Media exceeds the approved snapshot size limit" + ); + hash.update(&chunk[..count]); + if let Some(destination) = output.as_mut() { + destination.write_all(&chunk[..count])?; + } + progress(size)?; + } + cancelled(limits)?; + Ok((hex::encode(hash.finalize()), size)) +} +fn receipt_for(operation: &Operation, hash: String, size: u64) -> Receipt { + let intent = &operation.binding.intent; + Receipt { + version: 1, + request_id: intent.request_id.clone(), + nonce: intent.nonce.clone(), + app_audience: intent.app_audience.clone(), + node_did: intent.node_did.clone(), + producer: intent.producer.clone(), + project_id: intent.project_id.clone(), + price_sats: intent.price_sats, + viewing_seconds: intent.viewing_seconds, + expires_at: intent.expires_at, + content_id: format!("registered_{}", intent.request_id), + sha256: hash, + size_bytes: size.to_string(), + payment_methods: operation.binding.selection.payment_methods.clone(), + issued_at: operation.issued_at, + signature: String::new(), + } +} + +/// Performs disk I/O and cross-process locking; run on a blocking worker. +/// `now` is caller-supplied trusted UTC seconds, never a request-body timestamp. +/// `progress` may return an error to cancel; it must not change authorized terms. +/// A completed exact retry may return the original receipt after expiry. An +/// incomplete operation cannot create a new receipt after its intent expires. +/// Failed/cancelled staging files are retained privately, never publicly served. +#[allow(clippy::too_many_arguments)] +pub fn prepare( + data_dir: &Path, + cloud_root: &Path, + identity: &crate::identity::NodeIdentity, + pin: &InstallationPin, + intent: &Intent, + selection: &AuthorizedSelection, + now: u64, + limits: &Limits<'_>, + mut progress: impl FnMut(u64) -> Result<()>, +) -> Result { + let started = Instant::now(); + validate(intent, selection, pin, identity, now)?; + cancelled(limits)?; + let data_dir = data_dir + .canonicalize() + .context("Node data directory unavailable")?; + let cloud_root = cloud_root + .canonicalize() + .context("Configured Cloud root unavailable")?; + let binding = Binding { + intent: intent.clone(), + selection: selection.clone(), + cloud_root: cloud_root.clone(), + }; + let data = open_directory(&data_dir)?; + let root = private_directory(&data, PRIVATE_DIRECTORY)?; + let operation_dir = private_directory(&root, &intent.request_id)?; + // Never wait indefinitely for another worker/process. Live intents also stop + // waiting at their expiry; a completed expired replay gets a bounded 30s wait. + let wait_seconds = if now < intent.expires_at { + (intent.expires_at - now).min(30) + } else { + 30 + }; + lock_operation( + &operation_dir, + limits, + started + Duration::from_secs(wait_seconds), + )?; + let operation: Operation = + if let Some(saved) = read_record::(&operation_dir, "operation.json")? { + anyhow::ensure!( + saved.version == 1 + && saved.binding == binding + && saved.issued_at >= intent.created_at.saturating_sub(30) + && saved.issued_at < intent.expires_at, + "Registration request was already bound to different terms" + ); + saved + } else { + anyhow::ensure!(now < intent.expires_at, "Registration intent expired"); + let cloud = open_directory(&cloud_root)?; + let source = open_cloud_file(&cloud, &selection.relative_path)?; + let stamp = SourceStamp::read(&source)?; + anyhow::ensure!( + stamp.size <= limits.max_bytes, + "Media exceeds the approved snapshot size limit" + ); + let operation = Operation { + version: 1, + binding, + source: stamp, + issued_at: now, + }; + save_record(&operation_dir, "operation.json", &operation)?; + operation + }; + let existing_receipt: Option = read_record(&operation_dir, "receipt.json")?; + let mut saved_snapshot: Option = read_record(&operation_dir, "snapshot.json")?; + let mut copied = None; + let mut snapshot = match open_at( + &operation_dir, + "media", + libc::O_RDONLY | libc::O_NONBLOCK, + 0, + ) { + Ok(file) => file, + Err(error) + if error + .downcast_ref::() + .is_some_and(|e| e.kind() == std::io::ErrorKind::NotFound) => + { + anyhow::ensure!( + existing_receipt.is_none(), + "Registered snapshot is missing; preserve the receipt for recovery" + ); + anyhow::ensure!( + now.saturating_add(started.elapsed().as_secs()) < intent.expires_at, + "Registration intent expired" + ); + let cloud = open_directory(&cloud_root)?; + let mut source = open_cloud_file(&cloud, &selection.relative_path)?; + anyhow::ensure!( + SourceStamp::read(&source)? == operation.source, + "Selected Cloud file changed; use a new reviewed intent" + ); + let (name, mut destination) = temporary(&operation_dir)?; + let (hash, size) = + hash_file(&mut source, Some(&mut destination), limits, &mut progress)?; + anyhow::ensure!( + SourceStamp::read(&source)? == operation.source && size == operation.source.size, + "Selected Cloud file changed while copying; no receipt issued" + ); + destination.set_permissions(std::fs::Permissions::from_mode(0o400))?; + destination.sync_all()?; + cancelled(limits)?; + let record = SnapshotRecord { + sha256: hash.clone(), + size, + }; + if let Some(saved) = &saved_snapshot { + anyhow::ensure!( + *saved == record, + "Snapshot recovery produced different bytes" + ); + } else { + save_record(&operation_dir, "snapshot.json", &record)?; + saved_snapshot = Some(record); + } + publish_file(&operation_dir, &name, "media")?; + copied = Some((hash, size)); + open_at( + &operation_dir, + "media", + libc::O_RDONLY | libc::O_NONBLOCK, + 0, + )? + } + Err(error) => return Err(error), + }; + let before = SourceStamp::read(&snapshot)?; + anyhow::ensure!( + snapshot.metadata()?.mode() & 0o7777 == 0o400, + "Registered snapshot permissions changed" + ); + let (hash, size) = match copied { + Some(result) => result, + None => hash_file(&mut snapshot, None, limits, &mut progress)?, + }; + anyhow::ensure!( + SourceStamp::read(&snapshot)? == before && size == operation.source.size, + "Registered snapshot changed or is incomplete" + ); + let saved_snapshot = saved_snapshot + .context("Snapshot commitment is missing; preserve registration data for recovery")?; + anyhow::ensure!( + saved_snapshot.sha256 == hash && saved_snapshot.size == size, + "Registered snapshot does not match its durable byte commitment" + ); + let mut receipt = receipt_for(&operation, hash, size); + if let Some(saved) = existing_receipt { + anyhow::ensure!( + lower_hex(&saved.signature, 128), + "Invalid stored registration signature encoding" + ); + let signature = Signature::from_slice(&hex::decode(&saved.signature)?)?; + identity + .signing_key() + .verifying_key() + .verify_strict(&saved.preimage()?, &signature)?; + receipt.signature = saved.signature.clone(); + anyhow::ensure!( + receipt == saved, + "Stored registration receipt does not match its immutable snapshot" + ); + } else { + anyhow::ensure!( + now.saturating_add(started.elapsed().as_secs()) < intent.expires_at, + "Registration intent expired before completion; no receipt issued" + ); + receipt.signature = + hex::encode(identity.signing_key().sign(&receipt.preimage()?).to_bytes()); + save_record(&operation_dir, "receipt.json", &receipt)?; + } + operation_dir.sync_all()?; + // Reopen from the held directory to return a descriptor positioned at zero. + let snapshot = open_at( + &operation_dir, + "media", + libc::O_RDONLY | libc::O_NONBLOCK, + 0, + )?; + Ok(PreparedRegistration { + receipt, + snapshot, + snapshot_path: data_dir + .join(PRIVATE_DIRECTORY) + .join(&intent.request_id) + .join("media"), + }) +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::Arc; + + struct Fixture { + root: tempfile::TempDir, + identity: crate::identity::NodeIdentity, + pin: InstallationPin, + intent: Intent, + selection: AuthorizedSelection, + cancelled: AtomicBool, + } + impl Fixture { + async fn new() -> Self { + let root = tempfile::tempdir().unwrap(); + std::fs::create_dir(root.path().join("cloud")).unwrap(); + std::fs::write(root.path().join("cloud/film.mp4"), vec![42u8; 150_000]).unwrap(); + let identity = + crate::identity::NodeIdentity::load_or_create(&root.path().join("identity")) + .await + .unwrap(); + let pin = InstallationPin { + node_did: identity.did_key().unwrap(), + app_audience: "fixture-indeehub".into(), + }; + let intent = Intent { + version: 1, + request_id: uuid::Uuid::new_v4().to_string(), + nonce: "ab".repeat(32), + app_audience: pin.app_audience.clone(), + node_did: pin.node_did.clone(), + producer: "cd".repeat(32), + project_id: "fixture-project".into(), + price_sats: 15, + viewing_seconds: 3600, + created_at: 1000, + expires_at: 1600, + }; + Self { + root, + identity, + pin, + intent, + selection: AuthorizedSelection { + relative_path: "film.mp4".into(), + payment_methods: vec!["cashu".into(), "lightning-cashu".into()], + }, + cancelled: AtomicBool::new(false), + } + } + fn run(&self, now: u64) -> Result { + self.with_progress(now, |_| Ok(())) + } + fn with_progress( + &self, + now: u64, + progress: impl FnMut(u64) -> Result<()>, + ) -> Result { + prepare( + self.root.path(), + &self.root.path().join("cloud"), + &self.identity, + &self.pin, + &self.intent, + &self.selection, + now, + &Limits { + max_bytes: 1_000_000, + cancelled: &self.cancelled, + }, + progress, + ) + } + fn operation_dir(&self) -> PathBuf { + self.root + .path() + .join(PRIVATE_DIRECTORY) + .join(&self.intent.request_id) + } + } + + #[tokio::test] + async fn independent_nodejs_golden_wire_and_signature_match_prepared_receipt() { + let vector: serde_json::Value = + serde_json::from_str(include_str!("media_registration/fixtures/v1.json")).unwrap(); + let mut fixture = Fixture::new().await; + // Fixed public test seed only, loaded through the existing identity API. + std::fs::write( + fixture.root.path().join("identity/node_key"), + hex::decode(vector["testSeedHex"].as_str().unwrap()).unwrap(), + ) + .unwrap(); + fixture.identity = + crate::identity::NodeIdentity::load_existing(&fixture.root.path().join("identity")) + .await + .unwrap(); + fixture.pin = InstallationPin { + node_did: vector["pin"]["nodeDid"].as_str().unwrap().into(), + app_audience: vector["pin"]["appAudience"].as_str().unwrap().into(), + }; + fixture.intent = serde_json::from_value(vector["intent"].clone()).unwrap(); + std::fs::write( + fixture.root.path().join("cloud/film.mp4"), + vector["mediaUtf8"].as_str().unwrap(), + ) + .unwrap(); + let expected: Receipt = serde_json::from_value(vector["receipt"].clone()).unwrap(); + assert_eq!( + fixture.identity.pubkey_hex(), + vector["pin"]["publicKey"].as_str().unwrap() + ); + assert_eq!( + expected.preimage().unwrap(), + vector["preimageUtf8"].as_str().unwrap().as_bytes() + ); + assert!(crate::identity::NodeIdentity::verify( + &fixture.identity.pubkey_hex(), + vector["preimageUtf8"].as_str().unwrap().as_bytes(), + &expected.signature + ) + .unwrap()); + assert_eq!(fixture.run(1000).unwrap().receipt, expected); + } + + #[test] + fn expired_lock_deadline_returns_without_waiting() { + let root = tempfile::tempdir().unwrap(); + let dir = open_directory(root.path()).unwrap(); + let cancelled = AtomicBool::new(false); + assert!(lock_operation( + &dir, + &Limits { + max_bytes: 1, + cancelled: &cancelled + }, + Instant::now() + ) + .is_err()); + cancelled.store(true, Ordering::Relaxed); + assert!(lock_operation( + &dir, + &Limits { + max_bytes: 1, + cancelled: &cancelled + }, + Instant::now() + Duration::from_secs(30) + ) + .is_err()); + } + + #[tokio::test] + async fn snapshot_and_signed_receipt_are_durable_private_and_exactly_recoverable() { + let fixture = Fixture::new().await; + let key_before = std::fs::read(fixture.root.path().join("identity/node_key")).unwrap(); + let mut prepared = fixture.run(1000).unwrap(); + let mut bytes = Vec::new(); + prepared.snapshot.read_to_end(&mut bytes).unwrap(); + assert_eq!(bytes, vec![42u8; 150_000]); + assert_eq!(prepared.receipt.sha256, hex::encode(Sha256::digest(&bytes))); + assert_eq!(prepared.receipt.size_bytes, "150000"); + assert_eq!(prepared.snapshot.metadata().unwrap().mode() & 0o7777, 0o400); + assert_eq!( + fixture.operation_dir().metadata().unwrap().mode() & 0o777, + 0o700 + ); + assert!(crate::identity::NodeIdentity::verify( + &fixture.identity.pubkey_hex(), + &prepared.receipt.preimage().unwrap(), + &prepared.receipt.signature + ) + .unwrap()); + let wire = serde_json::to_value(&prepared.receipt).unwrap(); + assert_eq!(wire.as_object().unwrap().len(), 16); + assert_eq!(wire["appAudience"], "fixture-indeehub"); + assert!(wire.get("createdAt").is_none()); + let expected = serde_json::json!([ + DOMAIN, + fixture.intent.request_id, + fixture.intent.nonce, + fixture.pin.app_audience, + fixture.pin.node_did, + fixture.intent.producer, + fixture.intent.project_id, + prepared.receipt.content_id, + prepared.receipt.sha256, + "150000", + 15, + 3600, + ["cashu", "lightning-cashu"], + 1000, + 1600 + ]); + assert_eq!( + prepared.receipt.preimage().unwrap(), + serde_json::to_vec(&expected).unwrap() + ); + // A finished retry uses the retained bytes, even after source removal and expiry. + std::fs::remove_file(fixture.root.path().join("cloud/film.mp4")).unwrap(); + assert_eq!(fixture.run(1700).unwrap().receipt, prepared.receipt); + assert_eq!( + std::fs::read(fixture.root.path().join("identity/node_key")).unwrap(), + key_before + ); + assert!(!fixture.root.path().join("content.json").exists()); + } + + #[tokio::test] + async fn changed_terms_identity_selection_or_methods_cannot_reuse_a_request() { + let mut fixture = Fixture::new().await; + let original = fixture.run(1000).unwrap().receipt; + fixture.intent.price_sats += 1; + assert!(fixture.run(1000).is_err()); + fixture.intent.price_sats -= 1; + fixture.intent.nonce = "ef".repeat(32); + assert!(fixture.run(1000).is_err()); + fixture.intent.nonce = "ab".repeat(32); + fixture.selection.relative_path = "other.mp4".into(); + assert!(fixture.run(1000).is_err()); + fixture.selection.relative_path = "film.mp4".into(); + fixture.selection.payment_methods = vec!["lightning".into()]; + assert!(fixture.run(1000).is_err()); + fixture.selection.payment_methods = vec!["cashu".into(), "lightning-cashu".into()]; + fixture.pin.app_audience = "different-installation".into(); + assert!(fixture.run(1000).is_err()); + fixture.pin.app_audience = "fixture-indeehub".into(); + assert_eq!(fixture.run(1000).unwrap().receipt, original); + } + + #[tokio::test] + async fn traversal_symlinks_directories_and_special_files_are_rejected() { + let mut fixture = Fixture::new().await; + let outside = fixture.root.path().join("outside"); + std::fs::create_dir(&outside).unwrap(); + std::fs::write(outside.join("secret"), b"not selected").unwrap(); + std::os::unix::fs::symlink(&outside, fixture.root.path().join("cloud/escape")).unwrap(); + std::os::unix::fs::symlink("film.mp4", fixture.root.path().join("cloud/link")).unwrap(); + let fifo = c_path(&fixture.root.path().join("cloud/fifo")).unwrap(); + assert_eq!(unsafe { libc::mkfifo(fifo.as_ptr(), 0o600) }, 0); + for path in [ + "../outside/secret", + "/etc/passwd", + "escape/secret", + "link", + ".", + "fifo", + ] { + fixture.intent.request_id = uuid::Uuid::new_v4().to_string(); + fixture.selection.relative_path = path.into(); + assert!(fixture.run(1000).is_err(), "accepted {path}"); + assert!(!fixture.operation_dir().join("receipt.json").exists()); + } + assert_eq!( + std::fs::read(outside.join("secret")).unwrap(), + b"not selected" + ); + } + + #[tokio::test] + async fn cancellation_preserves_bound_operation_and_retries_without_changing_receipt_terms() { + let fixture = Fixture::new().await; + assert!(fixture + .with_progress(1000, |_| { + fixture.cancelled.store(true, Ordering::Relaxed); + Ok(()) + }) + .is_err()); + assert!(fixture.operation_dir().join("operation.json").exists()); + assert!(!fixture.operation_dir().join("receipt.json").exists()); + fixture.cancelled.store(false, Ordering::Relaxed); + let receipt = fixture.run(1001).unwrap().receipt; + assert_eq!(receipt.issued_at, 1000); + assert_eq!(fixture.run(1002).unwrap().receipt, receipt); + } + + #[tokio::test] + async fn changed_cloud_bytes_during_copy_or_before_retry_never_get_a_receipt() { + let fixture = Fixture::new().await; + let path = fixture.root.path().join("cloud/film.mp4"); + let mut changed = false; + assert!(fixture + .with_progress(1000, |_| { + if !changed { + File::options().write(true).open(&path)?.set_len(10)?; + changed = true; + } + Ok(()) + }) + .is_err()); + assert!(!fixture.operation_dir().join("receipt.json").exists()); + assert!(fixture.run(1001).is_err()); + assert_eq!(path.metadata().unwrap().len(), 10); + } + + #[tokio::test] + async fn parallel_callers_receive_one_persisted_receipt() { + let fixture = Arc::new(Fixture::new().await); + let workers: Vec<_> = (0..4) + .map(|_| { + let fixture = Arc::clone(&fixture); + std::thread::spawn(move || fixture.run(1000).unwrap().receipt) + }) + .collect(); + let receipts: Vec<_> = workers + .into_iter() + .map(|worker| worker.join().unwrap()) + .collect(); + assert!(receipts.iter().all(|receipt| receipt == &receipts[0])); + assert_eq!(fixture.run(1000).unwrap().receipt, receipts[0]); + } + + #[tokio::test] + async fn durable_snapshot_recovers_before_receipt_but_expired_pending_work_cannot_sign() { + let fixture = Fixture::new().await; + let receipt = fixture.run(1000).unwrap().receipt; + // Simulate interruption after snapshot commit and before receipt persistence. + std::fs::remove_file(fixture.operation_dir().join("receipt.json")).unwrap(); + std::fs::remove_file(fixture.root.path().join("cloud/film.mp4")).unwrap(); + assert!(fixture.run(1600).is_err()); + assert_eq!(fixture.run(1001).unwrap().receipt, receipt); + } + + #[tokio::test] + async fn corrupt_snapshot_or_missing_commitment_never_reissues_a_receipt() { + let fixture = Fixture::new().await; + let prepared = fixture.run(1000).unwrap(); + std::fs::set_permissions( + &prepared.snapshot_path, + std::fs::Permissions::from_mode(0o600), + ) + .unwrap(); + std::fs::write(&prepared.snapshot_path, vec![21u8; 150_000]).unwrap(); + std::fs::set_permissions( + &prepared.snapshot_path, + std::fs::Permissions::from_mode(0o400), + ) + .unwrap(); + assert!(fixture.run(1000).is_err()); + std::fs::remove_file(fixture.operation_dir().join("receipt.json")).unwrap(); + assert!(fixture.run(1000).is_err()); + std::fs::remove_file(fixture.operation_dir().join("snapshot.json")).unwrap(); + assert!(fixture.run(1000).is_err()); + } + + #[tokio::test] + async fn invalid_intents_quota_and_damaged_records_fail_without_replacing_data() { + let mut fixture = Fixture::new().await; + fixture.intent.price_sats = MAX_SAFE_INTEGER + 1; + assert!(fixture.run(1000).is_err()); + fixture.intent.price_sats = 15; + assert!(fixture.run(1600).is_err()); + assert!(prepare( + fixture.root.path(), + &fixture.root.path().join("cloud"), + &fixture.identity, + &fixture.pin, + &fixture.intent, + &fixture.selection, + 1000, + &Limits { + max_bytes: 1, + cancelled: &fixture.cancelled + }, + |_| Ok(()) + ) + .is_err()); + assert!(!fixture.operation_dir().join("operation.json").exists()); + std::fs::write( + fixture.operation_dir().join("operation.json"), + b"damaged fixture", + ) + .unwrap(); + assert!(fixture.run(1000).is_err()); + assert_eq!( + std::fs::read(fixture.operation_dir().join("operation.json")).unwrap(), + b"damaged fixture" + ); + } +} diff --git a/core/archipelago/src/media_registration/fixtures/v1.json b/core/archipelago/src/media_registration/fixtures/v1.json new file mode 100644 index 00000000..811955b1 --- /dev/null +++ b/core/archipelago/src/media_registration/fixtures/v1.json @@ -0,0 +1,45 @@ +{ + "description": "Public deterministic test vector only. Seed 07 repeated 32 times is never a node or user key.", + "testSeedHex": "0707070707070707070707070707070707070707070707070707070707070707", + "mediaUtf8": "IndeeHub fixture media\n", + "pin": { + "publicKey": "ea4a6c63e29c520abef5507b132ec5f9954776aebebe7b92421eea691446d22c", + "nodeDid": "did:key:z6MkvDqGT54cXesYGvABpF1UapVNwjCqRcafi4Px6Thv5T3Z", + "appAudience": "fixture-indeehub" + }, + "intent": { + "version": 1, + "requestId": "00000000-0000-4000-8000-000000000001", + "nonce": "abababababababababababababababababababababababababababababababab", + "appAudience": "fixture-indeehub", + "nodeDid": "did:key:z6MkvDqGT54cXesYGvABpF1UapVNwjCqRcafi4Px6Thv5T3Z", + "producer": "cdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcd", + "projectId": "fixture-project", + "priceSats": 15, + "viewingSeconds": 3600, + "createdAt": 1000, + "expiresAt": 1600 + }, + "receipt": { + "version": 1, + "requestId": "00000000-0000-4000-8000-000000000001", + "nonce": "abababababababababababababababababababababababababababababababab", + "appAudience": "fixture-indeehub", + "nodeDid": "did:key:z6MkvDqGT54cXesYGvABpF1UapVNwjCqRcafi4Px6Thv5T3Z", + "producer": "cdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcd", + "projectId": "fixture-project", + "priceSats": 15, + "viewingSeconds": 3600, + "expiresAt": 1600, + "contentId": "registered_00000000-0000-4000-8000-000000000001", + "sha256": "bbe573fdac96c464b67257a97c8667713c124c80927c01bab75b7178f4c0a661", + "sizeBytes": "23", + "paymentMethods": [ + "cashu", + "lightning-cashu" + ], + "issuedAt": 1000, + "signature": "9cb11fca6a946b38049402e8833992130922140f6bb5f5cc3b0806bc9865657be084c3b6fcfa1e2aa31f78f81ba44df5f32d824caa1c9f926f8d90f6578a350e" + }, + "preimageUtf8": "[\"archipelago.indeehub.media-registration.v1\",\"00000000-0000-4000-8000-000000000001\",\"abababababababababababababababababababababababababababababababab\",\"fixture-indeehub\",\"did:key:z6MkvDqGT54cXesYGvABpF1UapVNwjCqRcafi4Px6Thv5T3Z\",\"cdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcd\",\"fixture-project\",\"registered_00000000-0000-4000-8000-000000000001\",\"bbe573fdac96c464b67257a97c8667713c124c80927c01bab75b7178f4c0a661\",\"23\",15,3600,[\"cashu\",\"lightning-cashu\"],1000,1600]" +} diff --git a/docs/indeehub-node-registration-followup.md b/docs/indeehub-node-registration-followup.md new file mode 100644 index 00000000..1514aceb --- /dev/null +++ b/docs/indeehub-node-registration-followup.md @@ -0,0 +1,153 @@ +# IndeeHub node media registration + +Status: node preparation primitive compiled and qualified locally; RPC/serving +integration, deployment and actual-node acceptance remain open. App-side v1 intent/receipt implementation is qualified in the +separate IndeeHub checkout at `c4b920f` (190 backend tests, production build). +This document records the integration contract, not a completed publishing flow. + +## Implemented boundary + +`core/archipelago/src/media_registration.rs` accepts an existing `NodeIdentity`, +an installer-owned node DID/app-instance audience pin, an authenticated v1 app +intent, the operator-approved relative Cloud selection and receiving methods, +trusted current UTC seconds, a byte limit, cancellation flag and progress callback. +It does not create/recover keys, select files for the operator, expose an RPC, +change existing content policies, publish a catalog entry or grant playback. + +Call `prepare` on a blocking worker: it performs filesystem I/O and waits for a +cross-process per-request filesystem lock. Lock acquisition polls cancellation +and stops at the earlier of 30 seconds or an unexpired intent's deadline. An +exact completed retry after expiry may still wait at most 30 seconds; caller +cancellation can end that wait earlier. Root has declared the module in `main.rs` for the queued combined isolated +backend qualification. RPC/UI and serving integration remain unwired. The public Rust +selection structure is a caller assertion, not an authorization token. + +The implementation validates canonical v4 request IDs, lowercase nonce/producer +keys, ASCII project/audience IDs, identity and installed audience equality, +JavaScript-safe numbers, ten-minute maximum intent lifetime, bounded viewing +windows and sorted distinct supported methods. It checks the existing node key's +DID against the installation pin, never a browser-supplied replacement key. + +Linux `openat2` with `RESOLVE_BENEATH|RESOLVE_NO_SYMLINKS` resolves the approved +relative path from a held Cloud-root descriptor. An `O_PATH` descriptor is checked +for regular-file type before reopening that same held file for reading; devices +and FIFOs are not opened for I/O. Unsupported kernels fail explicitly. Source +inode/device, size and modification/change times must remain stable through the +streaming copy. The original Cloud file is never altered by preparation. + +Each request gets its own private 0700 directory under +`/media-registration//`. Files are not placed under a +web-served Cloud or content directory. The durable records are: + +- `operation.json`: original intent, exact selection/methods/root and original + file identity/metadata, plus original receipt issue time. +- `snapshot.json`: SHA256 and uint64 byte size committed before snapshot publication. +- `media`: private 0400 snapshot retaining the approved original bytes. +- `receipt.json`: exact fixed-domain Ed25519 receipt after all preceding commits. + +Records and snapshots are fsynced; no-replace hard-link publication and directory +fsync establish durable names. Intermediate files are named `pending-` in +the private operation directory. Interrupted staging is retained for controlled +recovery/cleanup, never served. Future cleanup must distinguish these private +partials from a completed snapshot or outstanding receipt and must not remove +original Cloud data or published/rented bytes. + +An identical completed retry verifies the saved signature, snapshot byte commitment +and all original terms, then returns the same receipt even after its initial +expiry. A pending operation cannot issue its first receipt after expiry. Changed +terms, selection, audience, node identity or receiving methods reject reuse. +Pending snapshot recovery requires either the committed immutable snapshot or +unchanged source metadata/bytes. Damaged records or missing completed bytes fail +without recreating the operation or rewriting its commitment. + +The wire schema and fixed ordered JSON array match +`/home/archipelago/Projects/indeehub-followup/docs/archipelago-registration-wire-v1.md`. +`sizeBytes` is a decimal string; the receipt contains exactly the 16 v1 fields. +The signature uses the already-loaded node's existing Ed25519 identity. + +## Caller work required before enabling this feature + +1. Obtain the intent through the authenticated installed IndeeHub bridge and its + owner-authenticated API. A browser merely presenting matching-looking JSON does + not prove that the app persisted an intent or owns the project. Bind app origin, + installation identity, native dashboard session and CSRF protections. +2. Show project/producer identity, the actual Cloud selection, price, viewing time + and supported receiving methods for explicit operator consent. Pass only the + configured Cloud root and the exact selected relative path. Neither filesystem + root nor installation pins may come from arbitrary request fields. +3. Supply a configured per-file quota, cancellation/progress and actual UTC time; + preserve backpressure and track cumulative private storage use. File size is + streamed and checked rather than read into memory. Private staging retention + needs a separate reviewed quota/cleanup policy. +4. Persist the opaque `contentId` → immutable snapshot mapping and complete content + policy under the same original request ID. Authenticate all content access and + keep the pending item private throughout setup. Do not use the mutable Cloud + filename as the serving identity or temporarily make it free/public. +5. Only after durable serving and entitlement linkage succeeds may the caller + return the prepared receipt to the app for intent consumption. If this step + fails, retry the same prepared request/receipt; do not create another content + ID, silently change rental terms or return a success notification. +6. Complete purchase recovery and seller receipt settlement, timed FIPS playback, + full immutable-byte/expiry enforcement and real node UAT before enabling the + app's registration/publication flags. A registration receipt is neither payment + proof nor a playback entitlement. Nothing in this primitive advertises serving + availability or receives funds. + +The returned `snapshot` descriptor is positioned at zero and read-only. The +`snapshot_path` is private node storage for the caller's persistent mapping; +it must never be sent to the app or used as a browser-controlled path. Any future +serving reopen must preserve descriptor containment/integrity and authorization. + +## Local qualification + +Eleven isolated tests cover original-key signing and exact wire preimage, +private durable snapshot/retry after source removal and expiry, changed bindings, +traversal/symlink/directory/FIFO rejection, cancellation/retry, source modification +during copying, concurrent callers, recovery at the snapshot/receipt boundary, +corrupt snapshot/missing commitment, invalid terms/expiry/quota and damaged-record +preservation. Identity fixtures use the existing node identity test helper only. + +`media_registration/fixtures/v1.json` is a deterministic golden wire artifact +created with Node.js crypto from the public fixed test seed `07` repeated 32 times. +It contains independent canonical array bytes, expected public key/DID, signature, +receipt and media. A byte-identical copy lives in the IndeeHub app's test fixtures; +Rust preparation and signature verification and the app receipt parser both check +it. Fixture SHA256: +`8fbfcbc3beb0b4758fadf677c39c688d55a89ed200d8a7cd8741de0da569feb3`. +No live node identity or user key is included. A separate test checks expired lock +deadline and pre-cancelled acquisition without waiting. + +The combined `scripts/test-backend-isolated.sh` run passed 1,808 tests with +zero failures and five existing skips, including all eleven media-registration +tests and the independent Node.js golden wire/signature fixture. Compilation +took 9m07s and execution 15.01s. Source hashes matched the frozen build inputs. +Log: `/tmp/archy-resumed-combined-backend-tests.log`; source provenance: +`/tmp/archy-resumed-combined-source-provenance.json`. + +Later integration tests must also exercise caller origin/CSRF/consent, bridge interruption, +serving linkage failure/retry, receiving-method availability and real paid FIPS +playback. No live file, wallet, public relay, service or app was changed here. + +### Existing app parser golden check + +The lightweight Node check passed using the existing production-compiled IndeeHub +receipt parser, with no rebuild, Jest worker, database or network. It checks fixed +node DID, exact UTF-8 signature preimage, acceptance of the golden pinned-key +signature and rejection after changing `sizeBytes`. + +Log: `/tmp/indeehub-registration-golden-dist-check.log`. + +- Fixture SHA256: + `8fbfcbc3beb0b4758fadf677c39c688d55a89ed200d8a7cd8741de0da569feb3`. +- Existing `backend/dist/archipelago/media-registration.protocol.js` SHA256: + `8e4a9f17649381c0b6bd8b9e187bd266b94743f40c63f5bc384ade204114b98b`, + built at 2026-10-06T22:23:09.597Z by the successful resumed backend build. +- Its unchanged source protocol SHA256: + `ec5b614516dcbd9918ff8fdecfc980f4430ff20362aced4ab348197e284337d7`, + source modification time 2026-10-06T22:07:50.141Z. + +This verifies the actual compiled app parser against the shared artifact. The +subsequent combined isolated Rust run also passed its independent golden fixture +test, establishing agreement across both implementations. Node caller/serving +integration and live acceptance remain open. No fixture or Rust source changed +between the frozen compilation inputs and the successful run.