diff --git a/.github/workflows/release-ffi.yml b/.github/workflows/release-ffi.yml index 4758ffb..d5f2b4d 100644 --- a/.github/workflows/release-ffi.yml +++ b/.github/workflows/release-ffi.yml @@ -46,7 +46,7 @@ jobs: # Dry-run needs no registry auth -- it only verifies metadata, # compresses the package, and checks the result, never uploads. # This is also the step that proves macula-rust-ffi's `macula-rust - # = { path = "..", version = "0.4" }` dependency actually resolves + # = { path = "..", version = "0.5" }` dependency actually resolves # against the REAL published macula-rust on crates.io, not just the # local workspace path -- `cargo publish` downloads and rebuilds # against the registry version during verification, confirmed diff --git a/CHANGELOG.md b/CHANGELOG.md index 04f985c..01bd1b5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,7 +13,28 @@ usually touches both, but their version numbers don't move in lockstep. ## macula-rust -### [0.4.0] - Unreleased +### [0.5.0] - Unreleased + +#### Added + +- Node-served content (macula 12's D27), ported from macula-go v0.12.0's + `pool/content.go`: `Pool::share_content` keeps the content, serves it on the + node's own `~/content_v1` server stream and announces it in the DHT, + renewed at half its hour; `unshare_content` withdraws it with a tombstone; + `get_content` finds the announcements bound to their announcer's own content + procedure in the realm, dials each sharer's station, and checks the block, + the manifest and every chunk against the content id, within + `ContentOptions` (256 MiB, 16,384 chunks, 4 streams at a time, 15 s each by + default). No realm key on either side. `content_procedure_bound` is public. + New `PoolError`s: `NotShared`, `ContentUnavailable`, `ContentMismatch`, + `ContentTooLarge`, `ContentReply`. +- `manifest`: macula 12's content manifests, byte for byte with macula's own + (`tests/vectors/manifest/erlang_manifests.json`): 256 KiB chunks, SHA-384, + 50-byte content ids, the odd-leaf Merkle fold, and the wire form, read with + macula_manifest's checks (sha384 only, whole chunks, fields in range). +- Example `content`. + +### [0.4.0] - 2026-09-26 The macula 12 wire. **Breaking throughout**: a 0.3 node cannot reach a macula 12 station, and nothing of the 0.3 API carries over. See the README's @@ -340,7 +361,20 @@ same day, not a separate feature set. Independently versioned from the core crate since day one (this crate started at 0.1.0 the same day the core crate did, but the two have moved at different paces ever since). -### [ffi-0.4.0] - Unreleased +### [ffi-0.5.0] - Unreleased + +#### Added + +- `FfiPool::share_content`, `unshare_content` and `get_content`, with + `FfiContentOptions` (zero for macula's defaults). New `FfiError`s: + `NotShared` and `ContentUnavailable`, which names why each sharer failed; a + content id of another length than 50 bytes is `WrongByteLength`. + +#### Changed + +- Builds on `macula-rust` 0.5. + +### [ffi-0.4.0] - 2026-09-26 **Breaking throughout**: rewritten on `macula-rust` 0.4's pool, the macula 12 wire. diff --git a/Cargo.lock b/Cargo.lock index 27d8431..b7501f7 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1225,7 +1225,7 @@ dependencies = [ [[package]] name = "macula-rust" -version = "0.4.0" +version = "0.5.0" dependencies = [ "apple-native-keyring-store", "aws-lc-rs", @@ -1247,7 +1247,7 @@ dependencies = [ [[package]] name = "macula-rust-ffi" -version = "0.4.0" +version = "0.5.0" dependencies = [ "async-trait", "hex", diff --git a/Cargo.toml b/Cargo.toml index 2e5fb44..ad7b8cf 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -3,7 +3,7 @@ members = ["macula-rust-ffi"] [package] name = "macula-rust" -version = "0.4.0" +version = "0.5.0" edition = "2021" rust-version = "1.89" authors = ["Macula "] diff --git a/README.md b/README.md index dfc866d..e8c9780 100644 --- a/README.md +++ b/README.md @@ -23,11 +23,11 @@ > ML-DSA-87 identities (in pq_hybrid, the fleet's profile, the ML-DSA-87 + > RSA-PSS-4096 composite), ML-KEM hybrid key exchange, and signed requests. > Calls and streams by direct dial, serving (under an org or in a node's own -> namespace), publish/subscribe and the DHT are tested against in-process -> macula 12 stations on every `cargo test`, and live against the fleet. Not -> here yet: UCAN-gated calls and node-served content; see [Not yet -> implemented](#not-yet-implemented). Releases before 0.4.0 speak the retired -> 10.x wire and cannot reach the current fleet. +> namespace), publish/subscribe, node-served content and the DHT are tested +> against in-process macula 12 stations on every `cargo test`, and calls, +> pubsub and the DHT live against the fleet. Not here yet: UCAN-gated calls; +> see [Not yet implemented](#not-yet-implemented). Releases before 0.4.0 speak +> the retired 10.x wire and cannot reach the current fleet. ## What is this? @@ -48,7 +48,7 @@ Swift. The core crate has no FFI dependency and no FFI-shaped types. ```toml [dependencies] -macula-rust = "0.4" +macula-rust = "0.5" tokio = { version = "1", features = ["full"] } ``` @@ -109,8 +109,8 @@ let served = pool .await?; ``` -Runnable versions are in [`examples/`](examples): `quickstart`, `serve` and -`publish_subscribe`, each reading the environment described at the top of +Runnable versions are in [`examples/`](examples): `quickstart`, `serve`, +`publish_subscribe` and `content`, each reading the environment described at the top of [`examples/common/mod.rs`](examples/common/mod.rs). ### Coming from 0.3 and earlier @@ -129,8 +129,10 @@ compatibility layer. `direct_dial::call` is simply `Pool::call`; `resolve` is `Pool::providers`; `serve_one_call` is `Pool::serve` with a handler; `Trust::WebPki` is gone: every station is pinned by its node_id. -- `ucan`, `cert_chain` and the content-transfer modules are gone until - macula 12's own arrive (see [Not yet implemented](#not-yet-implemented)). +- The 10.x content transfer (`content`, `manifest`, `put_direct`/`get_direct`) + is replaced by node-served content: `Pool::share_content`/`get_content`, with + macula 12's SHA-384 `manifest`. `ucan` and `cert_chain` are gone until macula + 12's own arrive (see [Not yet implemented](#not-yet-implemented)). - Serving an org procedure needs the realm's org directory and the org's delegation to your node in the DHT: a realm admits orgs through a human. @@ -145,6 +147,7 @@ compatibility layer. | A node's own namespace (`record::own_procedure`) | ✅ | ✅ | `~/`: served and called with no org and no realm key | | Streams (`open_stream`, `Offer::stream`) | ✅ | ✅ | Server, client and bidi; a QUIC stream per session, released on every path | | Publish/subscribe | ✅ | ✅ | Signed publications, delivered once across links | +| Node-served content (`share_content`, `unshare_content`, `get_content`) | ✅ | ✅ | Shared on the node's own `~/content_v1` and announced; a fetch checks the block, the manifest and every chunk against the content id, bounded (`ContentOptions`), with no realm key; manifests match macula's byte for byte (`manifest`) | | DHT (`find_record`, `find_records`, `find_records_by_type`, `put_record`) | ✅ | — | Records verified before they are handed on | | Mobile bindings (Kotlin, Swift) | ✅ | ✅ | `macula-rust-ffi`, below | @@ -165,8 +168,10 @@ integers within ±2^63, text or integer map keys, no duplicates). `macula-rust-ffi` wraps the pool with [UniFFI](https://mozilla.github.io/uniffi-rs/) proc macros: `FfiNodeKey`, `FfiPool`, `FfiSubscription`, `FfiStream`, and two handlers the app implements, `FfiCallHandler` and `FfiStreamHandler` -(`suspend fun` in Kotlin, `async throws` in Swift). Every 32-byte id crosses as -bytes and is checked. +(`suspend fun` in Kotlin, `async throws` in Swift). `FfiPool` also shares and +fetches content (`shareContent`, `getContent`, `FfiContentOptions`). Every id +crosses as bytes and is checked for its length (32 bytes, or 50 for a content +id). ```bash cargo build -p macula-rust-ffi --release @@ -203,7 +208,6 @@ crate 1.89. - **UCAN-gated calls and serving.** macula 12 uses post-quantum UCANs; calls carry no token yet, and a gated procedure cannot be served. -- **Node-served content** (macula 12's D27): planned for 0.5.0. - **Station discovery beyond the seeds.** macula's discovery call is not served by the fleet today (macula-io/macula#31); give the pool its seeds. diff --git a/examples/content.rs b/examples/content.rs new file mode 100644 index 0000000..b527788 --- /dev/null +++ b/examples/content.rs @@ -0,0 +1,32 @@ +//! Shares content from one node and fetches it from another by its content +//! id. The sharer keeps the content and serves it from its own namespace; +//! the fetcher checks everything it receives against the content id, so +//! neither needs a realm key. +//! +//! Run: `cargo run --example content`, with the environment +//! examples/common/mod.rs reads. The fetcher's key is `fetcher.key`. + +mod common; + +use macula_rust::pool::ContentOptions; + +#[tokio::main] +async fn main() -> Result<(), Box> { + let sharer = common::connect(None).await; + let data: Vec = (0..600_000usize).map(|i| (i % 251) as u8).collect(); + let mcid = sharer + .share_content(&common::realm(), &data, "example.bin") + .await?; + println!("shared {} bytes as {}", data.len(), common::hex(&mcid)); + + let fetcher = common::connect(Some("fetcher.key")).await; + let got = fetcher + .get_content(&common::realm(), &mcid, ContentOptions::default()) + .await?; + println!("fetched {} bytes, the same: {}", got.len(), got == data); + + sharer.unshare_content(&common::realm(), &mcid).await?; + fetcher.close().await; + sharer.close().await; + Ok(()) +} diff --git a/macula-rust-ffi/Cargo.toml b/macula-rust-ffi/Cargo.toml index 1e35d5d..0e5dea2 100644 --- a/macula-rust-ffi/Cargo.toml +++ b/macula-rust-ffi/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "macula-rust-ffi" -version = "0.4.0" +version = "0.5.0" edition = "2021" rust-version = "1.91" authors = ["Macula "] @@ -27,7 +27,7 @@ name = "uniffi-bindgen" path = "uniffi-bindgen.rs" [dependencies] -macula-rust = { path = "..", version = "0.4" } +macula-rust = { path = "..", version = "0.5" } uniffi = { version = "0.32", features = ["cli", "tokio"] } tokio = { version = "1", features = ["full"] } thiserror = "2" diff --git a/macula-rust-ffi/src/content.rs b/macula-rust-ffi/src/content.rs new file mode 100644 index 0000000..ec5bfaf --- /dev/null +++ b/macula-rust-ffi/src/content.rs @@ -0,0 +1,95 @@ +//! Node-served content (D27): a node shares content it keeps, served from +//! its own `~/content_v1` and announced in the DHT; another node +//! fetches it by its 50-byte content id, checking everything it receives +//! against that id. No realm key is needed on either side. + +use macula_rust::manifest::Mcid; +use macula_rust::pool::ContentOptions; + +use crate::pool::FfiPool; +use crate::{millis, to_32, FfiError}; + +/// A fetch's bounds. Zero is macula's default: 256 MiB, 16,384 chunks, 4 +/// chunk streams at a time, 15 seconds per stream. +#[derive(uniffi::Record, Debug, Clone, Copy, PartialEq, Eq, Default)] +pub struct FfiContentOptions { + #[uniffi(default = 0)] + pub max_bytes: u64, + #[uniffi(default = 0)] + pub max_chunks: u64, + #[uniffi(default = 0)] + pub parallel: u32, + #[uniffi(default = 0)] + pub chunk_timeout_ms: u64, +} + +impl From for ContentOptions { + fn from(o: FfiContentOptions) -> Self { + let d = ContentOptions::default(); + let or = |v: u64, fallback: u64| if v == 0 { fallback } else { v }; + ContentOptions { + max_bytes: or(o.max_bytes, d.max_bytes), + max_chunks: or(o.max_chunks, d.max_chunks), + parallel: or(u64::from(o.parallel), d.parallel as u64) as usize, + chunk_timeout: if o.chunk_timeout_ms == 0 { + d.chunk_timeout + } else { + millis(o.chunk_timeout_ms) + }, + } + } +} + +/// A content id: 50 bytes. +fn to_mcid(bytes: Vec) -> Result { + let actual = bytes.len() as u32; + bytes.try_into().map_err(|_| FfiError::WrongByteLength { + expected: 50, + actual, + }) +} + +#[uniffi::export(async_runtime = "tokio")] +impl FfiPool { + /// Keeps `data`, serves it and announces it in `realm`, until + /// [`unshare_content`](Self::unshare_content). Returns its 50-byte + /// content id: a raw block's for up to 256 KiB, a manifest's, named + /// `name`, above that. + pub async fn share_content( + &self, + realm: Vec, + data: Vec, + name: String, + ) -> Result, FfiError> { + Ok(self + .0 + .share_content(&to_32(realm)?, &data, &name) + .await? + .to_vec()) + } + + /// Stops sharing `mcid` in `realm` and withdraws its announcement. + pub async fn unshare_content(&self, realm: Vec, mcid: Vec) -> Result<(), FfiError> { + Ok(self + .0 + .unshare_content(&to_32(realm)?, &to_mcid(mcid)?) + .await?) + } + + /// Fetches the content `mcid` names in `realm` from a node that shares + /// it, checked against `mcid` throughout. [`FfiError::NotShared`] when + /// nobody announces it; [`FfiError::ContentUnavailable`] when every + /// sharer failed, naming why. + pub async fn get_content( + &self, + realm: Vec, + mcid: Vec, + options: FfiContentOptions, + ) -> Result, FfiError> { + let mcid = to_mcid(mcid)?; + Ok(self + .0 + .get_content(&to_32(realm)?, &mcid, options.into()) + .await?) + } +} diff --git a/macula-rust-ffi/src/lib.rs b/macula-rust-ffi/src/lib.rs index e664fbe..f12afd6 100644 --- a/macula-rust-ffi/src/lib.rs +++ b/macula-rust-ffi/src/lib.rs @@ -10,7 +10,8 @@ //! calls to a provider at its own station, serving a procedure with a //! handler the foreign side implements ([`FfiCallHandler`]), pubsub //! ([`FfiSubscription`]), streaming sessions on either side ([`FfiStream`], -//! [`FfiStreamHandler`]), and DHT records. +//! [`FfiStreamHandler`]), node-served content (`share_content`, +//! `get_content`, [`FfiContentOptions`]), and DHT records. //! //! [`FfiValue`] mirrors every variant [`macula_rust::cbor::Value`] has, //! narrowed only where the FFI boundary forces it: `Int` is `i64`, and an @@ -27,12 +28,14 @@ //! --language kotlin --out-dir bindings/kotlin //! ``` +mod content; mod node_key; mod pool; mod pubsub; mod serve; mod stream; +pub use content::FfiContentOptions; pub use node_key::{FfiNodeKey, FfiProfile}; pub use pool::{ own_procedure, FfiLinkStatus, FfiPool, FfiPoolOptions, FfiProvider, FfiRealmKey, FfiRecord, @@ -105,6 +108,14 @@ pub enum FfiError { /// An operation on a closed pool, subscription, stream or serving. #[error("closed")] Closed, + /// Content no node announces in the realm. + #[error("the content is not shared")] + NotShared, + /// Content every announcing sharer failed to give, and why each failed: + /// content that does not match its id, is over the fetch's bounds, or a + /// sharer that could not be reached. + #[error("{message}")] + ContentUnavailable { message: String }, /// Anything else the core crate reports, as its text. #[error("{message}")] Other { message: String }, @@ -148,6 +159,10 @@ impl From for FfiError { PoolError::Link(link) => link.into(), PoolError::NoRealmKey => FfiError::NoRealmKey, PoolError::Closed => FfiError::Closed, + PoolError::NotShared => FfiError::NotShared, + e @ PoolError::ContentUnavailable(_) => FfiError::ContentUnavailable { + message: e.to_string(), + }, e @ PoolError::NoProvider(_) => FfiError::NoProvider { message: e.to_string(), }, diff --git a/macula-rust-ffi/tests/pool_ffi.rs b/macula-rust-ffi/tests/pool_ffi.rs index de7c43e..bac14db 100644 --- a/macula-rust-ffi/tests/pool_ffi.rs +++ b/macula-rust-ffi/tests/pool_ffi.rs @@ -12,9 +12,9 @@ use std::time::Duration; use common::lab::{Lab, LabStation}; use macula_rust::profile::Profile; use macula_rust_ffi::{ - own_procedure, FfiCallHandler, FfiError, FfiNodeKey, FfiPool, FfiPoolOptions, FfiProfile, - FfiRealmKey, FfiRequest, FfiSeed, FfiStream, FfiStreamEvent, FfiStreamHandler, FfiStreamMode, - FfiValue, + own_procedure, FfiCallHandler, FfiContentOptions, FfiError, FfiNodeKey, FfiPool, + FfiPoolOptions, FfiProfile, FfiRealmKey, FfiRequest, FfiSeed, FfiStream, FfiStreamEvent, + FfiStreamHandler, FfiStreamMode, FfiValue, }; const ORG_PROCEDURE: &str = "mcl-echo/echo"; @@ -362,3 +362,69 @@ async fn records_go_through_the_pool() { "{bad_id:?}" ); } + +#[tokio::test(flavor = "multi_thread")] +async fn content_is_shared_by_one_node_and_fetched_by_another() { + let lab = Lab::start(Profile::PqPure); + let (sharing, fetching) = ( + lab.station("ffi content sharing"), + lab.station("ffi content fetching"), + ); + lab.share(&[&sharing, &fetching]); + let realm = vec![0x0c; 32]; + let sharer = pool( + &FfiNodeKey::generate(FfiProfile::PqPure).unwrap(), + &sharing, + options(), + ) + .await; + let fetcher = pool( + &FfiNodeKey::generate(FfiProfile::PqPure).unwrap(), + &fetching, + options(), + ) + .await; + let data: Vec = (0..600_000usize).map(|i| (i % 251) as u8).collect(); + let mcid = sharer + .share_content(realm.clone(), data.clone(), "blob.bin".into()) + .await + .unwrap(); + assert_eq!(mcid.len(), 50); + let got = fetcher + .get_content(realm.clone(), mcid.clone(), FfiContentOptions::default()) + .await + .unwrap(); + assert!(got == data, "{} bytes back", got.len()); + let small = FfiContentOptions { + max_bytes: 1_000, + ..FfiContentOptions::default() + }; + let refused = fetcher + .get_content(realm.clone(), mcid.clone(), small) + .await; + assert!( + matches!(refused, Err(FfiError::ContentUnavailable { .. })), + "{refused:?}" + ); + sharer + .unshare_content(realm.clone(), mcid.clone()) + .await + .unwrap(); + let gone = fetcher + .get_content(realm.clone(), mcid.clone(), FfiContentOptions::default()) + .await; + assert!(matches!(gone, Err(FfiError::NotShared)), "{gone:?}"); + let bad = fetcher + .get_content(realm, vec![2; 49], FfiContentOptions::default()) + .await; + assert!( + matches!( + bad, + Err(FfiError::WrongByteLength { + expected: 50, + actual: 49 + }) + ), + "{bad:?}" + ); +} diff --git a/src/lib.rs b/src/lib.rs index 4b458fc..62c49ba 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -10,6 +10,7 @@ pub mod cbor; pub mod frame; pub mod handshake; pub mod keystore; +pub mod manifest; pub mod node_key; pub mod petname; pub mod pool; diff --git a/src/manifest.rs b/src/manifest.rs new file mode 100644 index 0000000..09af1c2 --- /dev/null +++ b/src/manifest.rs @@ -0,0 +1,417 @@ +//! Content manifests as macula 12's macula_manifest builds them, byte for +//! byte: fixed-size chunks (256 KiB by default), SHA-384 hashes, a 50-byte +//! content id `<<2, Codec, SHA-384>>` (tag 2 names SHA-384, D24; codec 0x55 a +//! raw block, 0x56 a manifest), and a Merkle fold that pairs an odd last hash +//! with itself. +//! +//! A manifest's name has two encodings that must not be confused: its +//! content id hashes the name as CBOR text, while the wire form a manifest +//! travels in carries it as a byte string. + +use std::fmt; + +use sha2::{Digest, Sha384}; + +use crate::cbor::{self, Value}; + +/// 256 KiB, macula_manifest's default chunk size. +pub const DEFAULT_CHUNK_SIZE: u64 = 262_144; + +/// A SHA-384 digest's length. +pub const HASH_SIZE: usize = 48; + +/// A content id: `<>`. +pub type Mcid = [u8; 50]; + +/// A SHA-384 digest. +pub type Hash = [u8; HASH_SIZE]; + +/// The one hash algorithm a manifest names. +pub const SHA384: &str = "sha384"; + +const VERSION: u32 = 1; +const TAG_SHA384: u8 = 2; +const CODEC_RAW: u8 = 0x55; +const CODEC_MANIFEST: u8 = 0x56; + +/// One chunk of a manifest. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ChunkInfo { + pub index: u64, + pub offset: u64, + pub size: u64, + pub hash: Hash, +} + +/// A chunked content's manifest. `name` is bytes, as the wire carries it, and +/// `hash_algorithm` the wire's own text: a manifest read from a peer may name +/// anything, and [`verify_mcid`] refuses what is not a UTF-8 name and sha384. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct Manifest { + pub mcid: Mcid, + pub version: u32, + pub name: Vec, + pub size: u64, + /// Unix seconds. + pub created: u64, + pub chunk_size: u64, + pub chunk_count: u64, + pub hash_algorithm: String, + pub root_hash: Hash, + pub chunks: Vec, +} + +/// Why a manifest was refused. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum ManifestError { + /// A chunk size of zero. + ChunkSizeZero, + /// Content that does not match the manifest, and how. + ContentMismatch(String), + /// A manifest that does not describe the content id it was asked for by. + McidMismatch, + /// A manifest whose chunks do not describe its content whole, and where. + NotWhole(String), + /// A manifest whose chunk hashes do not make its root hash. + ChunkHashes, + /// A wire form that is not a manifest, and why. + Malformed(String), +} + +impl fmt::Display for ManifestError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + ManifestError::ChunkSizeZero => f.write_str("a chunk size of zero"), + ManifestError::ContentMismatch(why) => { + write!(f, "the content does not match the manifest: {why}") + } + ManifestError::McidMismatch => { + f.write_str("the manifest does not describe the content id") + } + ManifestError::NotWhole(why) => write!( + f, + "the manifest's chunks do not describe its content whole: {why}" + ), + ManifestError::ChunkHashes => { + f.write_str("the manifest's chunk hashes do not make its root hash") + } + ManifestError::Malformed(why) => write!(f, "not a manifest: {why}"), + } + } +} + +impl std::error::Error for ManifestError {} + +fn sha384(data: &[u8]) -> Hash { + Sha384::digest(data).into() +} + +fn make_mcid(codec: u8, hash: &Hash) -> Mcid { + let mut out = [0u8; 50]; + out[0] = TAG_SHA384; + out[1] = codec; + out[2..].copy_from_slice(hash); + out +} + +/// Splits `data` into chunks of `chunk_size` bytes and builds its manifest, +/// created now; returns it with the chunks in order. +pub fn create( + data: &[u8], + name: &str, + chunk_size: u64, +) -> Result<(Manifest, Vec>), ManifestError> { + let now = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.as_secs()) + .unwrap_or(0); + create_at(data, name, chunk_size, now) +} + +/// [`create`] with the manifest's creation time given, in unix seconds. +pub fn create_at( + data: &[u8], + name: &str, + chunk_size: u64, + created: u64, +) -> Result<(Manifest, Vec>), ManifestError> { + if chunk_size == 0 { + return Err(ManifestError::ChunkSizeZero); + } + let chunks: Vec> = data + .chunks(chunk_size as usize) + .map(<[u8]>::to_vec) + .collect(); + let infos = chunk_infos(&chunks); + let root_hash = root_hash_for(&infos); + let mut m = Manifest { + mcid: [0; 50], + version: VERSION, + name: name.as_bytes().to_vec(), + size: data.len() as u64, + created, + chunk_size, + chunk_count: infos.len() as u64, + hash_algorithm: SHA384.into(), + root_hash, + chunks: infos, + }; + m.mcid = mcid_for(&m); + Ok((m, chunks)) +} + +/// The content id chunk `index` is stored and fetched under: the block id of +/// its bytes, which a sharer derives from the manifest alone. +pub fn chunk_mcid(m: &Manifest, index: usize) -> Option { + m.chunks.get(index).map(|c| make_mcid(CODEC_RAW, &c.hash)) +} + +/// The content id of a single block: `<<2, 0x55, SHA-384(data)>>`. +pub fn block_mcid(data: &[u8]) -> Mcid { + make_mcid(CODEC_RAW, &sha384(data)) +} + +/// Whether a content id names a manifest (chunked content) rather than a +/// single block, read from its codec byte. +pub fn mcid_is_chunked(mcid: &Mcid) -> bool { + mcid[1] == CODEC_MANIFEST +} + +/// The content id the manifest's canonical fields describe: its name, size, +/// chunk size and count, hash algorithm and root hash. Its creation time and +/// chunk list are not part of it. +pub fn mcid_for(m: &Manifest) -> Mcid { + let canonical = Value::Map(vec![ + ( + Value::text("name"), + Value::text(String::from_utf8_lossy(&m.name)), + ), + (Value::text("size"), Value::Int(i128::from(m.size))), + ( + Value::text("chunk_size"), + Value::Int(i128::from(m.chunk_size)), + ), + ( + Value::text("chunk_count"), + Value::Int(i128::from(m.chunk_count)), + ), + ( + Value::text("hash_algorithm"), + Value::text(m.hash_algorithm.clone()), + ), + (Value::text("root_hash"), Value::Bytes(m.root_hash.to_vec())), + ]); + let encoded = cbor::encode(&canonical).expect("a manifest's canonical fields always encode"); + make_mcid(CODEC_MANIFEST, &sha384(&encoded)) +} + +/// Whether `m` describes `mcid`, as macula_manifest's verify_mcid/2 checks: +/// a UTF-8 name, sha384, and the content id its canonical fields recompute +/// to. The manifest's own `mcid` field is not consulted. +pub fn verify_mcid(m: &Manifest, mcid: &Mcid) -> Result<(), ManifestError> { + if std::str::from_utf8(&m.name).is_err() || m.hash_algorithm != SHA384 || mcid_for(m) != *mcid { + return Err(ManifestError::McidMismatch); + } + Ok(()) +} + +/// Checks reassembled `data` against `m`: its size, then a root hash over +/// `data` cut the same way. +pub fn verify(m: &Manifest, data: &[u8]) -> Result<(), ManifestError> { + if m.chunk_size == 0 { + return Err(ManifestError::ChunkSizeZero); + } + if data.len() as u64 != m.size { + return Err(ManifestError::ContentMismatch(format!( + "{} bytes, the manifest's {}", + data.len(), + m.size + ))); + } + let chunks: Vec> = data + .chunks(m.chunk_size as usize) + .map(<[u8]>::to_vec) + .collect(); + if root_hash_for(&chunk_infos(&chunks)) != m.root_hash { + return Err(ManifestError::ContentMismatch("another root hash".into())); + } + Ok(()) +} + +/// Checks that `m`'s chunks describe its content whole, cut as [`create`] +/// cuts it: a positive chunk size, ceil(size / chunk size) chunks, which is +/// its chunk count, chunk i at offset i × chunk size and chunk size long but +/// for the last, which holds what is left, between 1 and chunk size bytes. +pub fn check_whole(m: &Manifest) -> Result<(), ManifestError> { + let wanted = if m.chunk_size == 0 { + None + } else { + Some(m.size.div_ceil(m.chunk_size)) + }; + if wanted != Some(m.chunk_count) || m.chunk_count != m.chunks.len() as u64 { + return Err(ManifestError::NotWhole(format!( + "chunk size {}, size {}, {} chunks counted, {} listed", + m.chunk_size, + m.size, + m.chunk_count, + m.chunks.len() + ))); + } + for (i, c) in m.chunks.iter().enumerate() { + let offset = i as u64 * m.chunk_size; + let size = m.chunk_size.min(m.size - offset); + if c.index != i as u64 || c.offset != offset || c.size != size { + return Err(ManifestError::NotWhole(format!("chunk {i}"))); + } + } + Ok(()) +} + +/// Checks that `m`'s chunk hashes make its root hash. The root hash is part +/// of the content id and the chunk hashes are not, so after [`verify_mcid`] +/// this is what ties each chunk, fetched by its hash, to the content id. +pub fn check_chunk_hashes(m: &Manifest) -> Result<(), ManifestError> { + if root_hash_for(&m.chunks) != m.root_hash { + return Err(ManifestError::ChunkHashes); + } + Ok(()) +} + +fn chunk_infos(chunks: &[Vec]) -> Vec { + let mut offset = 0u64; + chunks + .iter() + .enumerate() + .map(|(i, chunk)| { + let info = ChunkInfo { + index: i as u64, + offset, + size: chunk.len() as u64, + hash: sha384(chunk), + }; + offset += chunk.len() as u64; + info + }) + .collect() +} + +/// The Merkle root: pairs from the front, hash(left || right), an odd last +/// hash paired with itself, until one is left; SHA-384 of nothing for no +/// chunks. +fn root_hash_for(infos: &[ChunkInfo]) -> Hash { + let mut level: Vec = infos.iter().map(|c| c.hash).collect(); + if level.is_empty() { + return sha384(&[]); + } + while level.len() > 1 { + level = level + .chunks(2) + .map(|pair| { + let right = pair.get(1).unwrap_or(&pair[0]); + sha384(&[pair[0].as_slice(), right.as_slice()].concat()) + }) + .collect(); + } + level[0] +} + +/// The manifest as it travels: the name as a byte string. +pub fn to_wire(m: &Manifest) -> Value { + let chunks = m + .chunks + .iter() + .map(|c| { + Value::Map(vec![ + (Value::text("index"), Value::Int(i128::from(c.index))), + (Value::text("offset"), Value::Int(i128::from(c.offset))), + (Value::text("size"), Value::Int(i128::from(c.size))), + (Value::text("hash"), Value::Bytes(c.hash.to_vec())), + ]) + }) + .collect(); + Value::Map(vec![ + (Value::text("mcid"), Value::Bytes(m.mcid.to_vec())), + (Value::text("version"), Value::Int(i128::from(m.version))), + (Value::text("name"), Value::Bytes(m.name.clone())), + (Value::text("size"), Value::Int(i128::from(m.size))), + (Value::text("created"), Value::Int(i128::from(m.created))), + ( + Value::text("chunk_size"), + Value::Int(i128::from(m.chunk_size)), + ), + ( + Value::text("chunk_count"), + Value::Int(i128::from(m.chunk_count)), + ), + ( + Value::text("hash_algorithm"), + Value::text(m.hash_algorithm.clone()), + ), + (Value::text("root_hash"), Value::Bytes(m.root_hash.to_vec())), + (Value::text("chunks"), Value::List(chunks)), + ]) +} + +/// A manifest read from its wire form, as macula_manifest's from_wire/1 +/// reads it. One that names another hash algorithm than sha384 (as text or +/// bytes), whose chunks do not describe its content whole, or holding a +/// number outside its field, is refused. Nothing is allocated from the size +/// or count it claims: only the chunks it lists are read. +pub fn from_wire(v: &Value) -> Result { + let chunks = match v.get("chunks") { + Some(Value::List(items)) => items + .iter() + .map(chunk_from_wire) + .collect::, _>>()?, + _ => return Err(malformed("chunks")), + }; + let m = Manifest { + mcid: bytes_exact(v, "mcid")?, + version: u32::try_from(uint(v, "version")?).map_err(|_| malformed("version"))?, + name: match v.get("name") { + Some(Value::Bytes(b)) => b.clone(), + _ => return Err(malformed("name")), + }, + size: uint(v, "size")?, + created: uint(v, "created")?, + chunk_size: uint(v, "chunk_size")?, + chunk_count: uint(v, "chunk_count")?, + hash_algorithm: match v.get("hash_algorithm") { + Some(Value::Text(t)) if t == SHA384 => SHA384.into(), + Some(Value::Bytes(b)) if b == SHA384.as_bytes() => SHA384.into(), + _ => return Err(malformed("hash_algorithm")), + }, + root_hash: bytes_exact(v, "root_hash")?, + chunks, + }; + check_whole(&m)?; + Ok(m) +} + +fn chunk_from_wire(v: &Value) -> Result { + Ok(ChunkInfo { + index: uint(v, "index")?, + offset: uint(v, "offset")?, + size: uint(v, "size")?, + hash: bytes_exact(v, "hash")?, + }) +} + +fn malformed(field: &str) -> ManifestError { + ManifestError::Malformed(format!("field {field:?} is missing or of the wrong type")) +} + +/// A field that is an integer between 0 and 2^63 - 1. +fn uint(v: &Value, field: &str) -> Result { + match v.get(field) { + Some(Value::Int(n)) if (0..=i128::from(i64::MAX)).contains(n) => Ok(*n as u64), + _ => Err(malformed(field)), + } +} + +fn bytes_exact(v: &Value, field: &str) -> Result<[u8; N], ManifestError> { + match v.get(field) { + Some(Value::Bytes(b)) => b.as_slice().try_into().map_err(|_| malformed(field)), + _ => Err(malformed(field)), + } +} diff --git a/src/pool.rs b/src/pool.rs index f39f524..cd1784f 100644 --- a/src/pool.rs +++ b/src/pool.rs @@ -17,11 +17,13 @@ //! (macula-io/macula#31). mod call; +mod content; mod member; mod pubsub; mod serve; pub use call::{Call, Provider, StreamCall}; +pub use content::{content_procedure_bound, ContentOptions, CONTENT_PROCEDURE}; pub use pubsub::Subscription; pub use serve::{Offer, Served}; @@ -86,6 +88,19 @@ pub enum PoolError { NotServed(Vec), /// A link's own failure, or a provider's or station's answer. Link(LinkError), + /// Content no node announces in the realm, or a sharer that does not + /// hold it. + NotShared, + /// Content every announcing sharer failed to give: each sharer's node + /// and why. + ContentUnavailable(Vec<([u8; 32], PoolError)>), + /// A block, manifest or whole that does not match the content id it was + /// asked for by, and which. + ContentMismatch(String), + /// Content over the fetch's bounds, and how large. + ContentTooLarge(String), + /// An answer a sharer never gives, and what it was. + ContentReply(String), } impl fmt::Display for PoolError { @@ -102,6 +117,13 @@ impl fmt::Display for PoolError { } Ok(()) } + PoolError::ContentUnavailable(tried) => { + f.write_str("no sharer gave the content:")?; + for (node, e) in tried { + write!(f, " [{}: {e}]", short(node))?; + } + Ok(()) + } other => write!(f, "{other:?}"), } } @@ -209,6 +231,7 @@ pub(crate) struct PoolInner { dedup: Arc, state: Mutex, ticks: tokio::task::JoinHandle<()>, + content: content::Sharer, } struct State { @@ -261,6 +284,7 @@ impl Pool { closed: false, }), ticks, + content: content::Sharer::default(), opts, }); let pool = Pool { inner }; diff --git a/src/pool/call.rs b/src/pool/call.rs index aee751b..7ed313f 100644 --- a/src/pool/call.rs +++ b/src/pool/call.rs @@ -375,7 +375,7 @@ impl PoolInner { .find(|l| l.station_node_id() == *station) } - async fn link_to( + pub(super) async fn link_to( self: &Arc, station: &[u8; 32], deadline: Instant, @@ -447,7 +447,7 @@ impl PoolInner { /// Runs `ask` on the links in selection order and returns the first /// answer, moving on only when a link could not carry the request: a /// station's own answer, not_found included, is final. - async fn first_answer<'a, T, F, Fut>(&self, ask: F) -> Result + pub(super) async fn first_answer<'a, T, F, Fut>(&self, ask: F) -> Result where F: Fn(Link) -> Fut, Fut: Future> + 'a, diff --git a/src/pool/content.rs b/src/pool/content.rs new file mode 100644 index 0000000..ba66c59 --- /dev/null +++ b/src/pool/content.rs @@ -0,0 +1,662 @@ +//! Node-served content, as macula 12 has it (D27). A station keeps no +//! content: the node that shares content keeps it, serves it on a +//! server_stream procedure of its own, `~/content_v1`, and announces +//! it in the DHT under the content id's key, naming the realm, the station it +//! is reachable through and that procedure. A fetcher finds the +//! announcements, dials each sharer's station from the station's own endpoint +//! record, and asks for the root and then each chunk, one stream each; +//! everything it receives is checked against the content id it asked for, so +//! it trusts no realm key and no sharer. + +use std::collections::HashMap; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::{Arc, Mutex, MutexGuard, Weak}; +use std::time::Duration; + +use tokio::sync::watch; +use tokio::time::Instant; + +use crate::cbor::Value; +use crate::frame::StreamMode; +use crate::manifest::{self, block_mcid, chunk_mcid, Manifest, Mcid, DEFAULT_CHUNK_SIZE}; +use crate::record::{ + self, new_content_announcement, ContentAnnouncementOptions, Reason, Record, RecordType, + TombstoneOptions, +}; +use crate::station_link::{ + self, stream_handler, Link, LinkError, Stream, StreamEvent, DEFAULT_CALL_TIMEOUT, +}; + +use super::{shuffle, Offer, Pool, PoolError, PoolInner, Served}; + +/// The content procedure's name in a sharer's own namespace. +pub const CONTENT_PROCEDURE: &str = "content_v1"; + +/// How long an announcement is signed for; it is renewed at half that. +const ANNOUNCEMENT_TTL: Duration = Duration::from_secs(60 * 60); + +/// One DATA body holds at most one chunk. +const MAX_BLOCK_BYTES: u64 = DEFAULT_CHUNK_SIZE; + +/// A fetch's bounds, macula_content_fetch's defaults: 256 MiB, 16,384 +/// chunks, 4 chunk streams at a time, 15 seconds per stream. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct ContentOptions { + pub max_bytes: u64, + pub max_chunks: u64, + pub parallel: usize, + pub chunk_timeout: Duration, +} + +impl Default for ContentOptions { + fn default() -> Self { + ContentOptions { + max_bytes: 256 << 20, + max_chunks: 16_384, + parallel: 4, + chunk_timeout: Duration::from_secs(15), + } + } +} + +/// What a pool shares: per realm, the content it keeps (roots and chunks, by +/// content id), the served content procedure, and each root's latest signed +/// announcement. +#[derive(Default)] +pub(super) struct Sharer { + realms: Mutex>, + /// Serves one realm's content procedure at a time. + serving: tokio::sync::Mutex<()>, +} + +struct SharedRealm { + _served: Served, + roots: HashMap, + chunks: HashMap>, + announcements: HashMap, +} + +#[derive(Clone)] +enum Root { + Block(Vec), + Chunked(Arc), +} + +struct Announced { + latest: Arc>, + stop: watch::Sender, +} + +impl Sharer { + fn lock(&self) -> MutexGuard<'_, HashMap<[u8; 32], SharedRealm>> { + self.realms.lock().unwrap_or_else(|p| p.into_inner()) + } +} + +/// Whether `procedure` is `node`'s content procedure: `~/content_v1`, +/// or `/content_v1_` under an org that is not empty, `_`, nor a +/// `~` namespace, each node as 64 lowercase hex, so an announcement cannot +/// point a fetch at another node's procedure. +pub fn content_procedure_bound(procedure: &str, node: &[u8; 32]) -> bool { + let node_hex: String = node.iter().map(|b| format!("{b:02x}")).collect(); + if procedure == record::own_procedure(node, CONTENT_PROCEDURE) { + return true; + } + match procedure.split_once('/') { + Some((org, name)) => { + !org.is_empty() + && org != "_" + && !org.starts_with(record::OWN_NAMESPACE_PREFIX) + && name == format!("{CONTENT_PROCEDURE}_{node_hex}") + } + None => false, + } +} + +impl Pool { + /// Keeps `data`, serves it on this node's content procedure in `realm` + /// and announces it, renewing the announcement until unshared or the + /// pool is dropped. Data of at most one chunk is one raw block; larger + /// data is a manifest over 256 KiB chunks, named `name`. Returns the + /// content id. + pub async fn share_content( + &self, + realm: &[u8; 32], + data: &[u8], + name: &str, + ) -> Result { + let inner = &self.inner; + let (mcid, root, chunks) = if data.len() as u64 <= MAX_BLOCK_BYTES { + (block_mcid(data), Root::Block(data.to_vec()), HashMap::new()) + } else { + let (m, parts) = manifest::create(data, name, DEFAULT_CHUNK_SIZE) + .map_err(|e| PoolError::ContentMismatch(e.to_string()))?; + let chunks: HashMap> = parts + .into_iter() + .enumerate() + .filter_map(|(i, part)| chunk_mcid(&m, i).map(|c| (c, part))) + .collect(); + (m.mcid, Root::Chunked(Arc::new(m)), chunks) + }; + inner.shared_realm(realm).await?; + let already = { + let mut realms = inner.content.lock(); + let shared = realms.get_mut(realm).ok_or(PoolError::Closed)?; + shared.roots.insert(mcid, root.clone()); + shared.chunks.extend(chunks); + shared.announcements.contains_key(&mcid) + }; + if already { + return Ok(mcid); + } + let latest = match inner.announce_content(realm, &mcid, &root).await { + Ok(latest) => latest, + Err(e) => { + inner.forget_shared(realm, &mcid); + return Err(e); + } + }; + let latest = Arc::new(Mutex::new(latest)); + let (stop, stopped) = watch::channel(false); + if let Some(shared) = inner.content.lock().get_mut(realm) { + shared.announcements.insert( + mcid, + Announced { + latest: latest.clone(), + stop, + }, + ); + } + tokio::spawn(renew_announcement( + Arc::downgrade(inner), + *realm, + mcid, + root, + latest, + stopped, + )); + Ok(mcid) + } + + /// Stops sharing `mcid` in `realm`: it is no longer served, and its + /// announcement is withdrawn with a tombstone this node signs. Content + /// not shared is nothing to do. + pub async fn unshare_content(&self, realm: &[u8; 32], mcid: &Mcid) -> Result<(), PoolError> { + let Some(announced) = self.inner.forget_shared(realm, mcid) else { + return Ok(()); + }; + announced.stop.send_replace(true); + let latest = announced + .latest + .lock() + .unwrap_or_else(|p| p.into_inner()) + .clone(); + let tombstone = + record::new_tombstone(&latest, Reason::Shutdown, &TombstoneOptions::default()) + .map_err(LinkError::from)?; + let signed = + record::sign(&tombstone, &self.inner.opts.identity).map_err(LinkError::from)?; + let wire = record::encode(&signed).map_err(LinkError::from)?; + self.put_record(&wire).await + } + + /// Fetches the content `mcid` names in `realm` from a node that shares + /// it, as macula_content_fetch does: the announcements under the content + /// id's key that name it, `realm`, a serving station and a content + /// procedure bound to their announcer; the sharers tried one at a time in + /// a random order; a block that must hash to `mcid`, or a manifest that + /// must match `mcid` before its sizes are read and fit `opts` before any + /// chunk is asked for, each chunk on its own stream, `opts.parallel` at a + /// time, checked against its own content id, and the whole against the + /// manifest. No realm key is needed. Content nobody announces is + /// [`PoolError::NotShared`]; when every sharer fails, + /// [`PoolError::ContentUnavailable`] names each failure. + pub async fn get_content( + &self, + realm: &[u8; 32], + mcid: &Mcid, + opts: ContentOptions, + ) -> Result, PoolError> { + let sharers = self.inner.content_sharers(realm, mcid).await?; + if sharers.is_empty() { + return Err(PoolError::NotShared); + } + let mut tried = Vec::new(); + for sharer in sharers { + match self.inner.fetch_from(realm, &sharer, mcid, &opts).await { + Ok(data) => return Ok(data), + Err(e) => tried.push((sharer.node, e)), + } + } + Err(PoolError::ContentUnavailable(tried)) + } +} + +/// An announcing node, where it is served, and on what. +#[derive(Debug, Clone)] +struct Sharing { + node: [u8; 32], + station: [u8; 32], + procedure: String, +} + +impl PoolInner { + /// `realm`'s share, serving the content procedure on first use. + async fn shared_realm(self: &Arc, realm: &[u8; 32]) -> Result<(), PoolError> { + let _serving = self.content.serving.lock().await; + if self.content.lock().contains_key(realm) { + return Ok(()); + } + let pool = Arc::downgrade(self); + let r = *realm; + let answer = stream_handler(move |s| { + let pool = pool.clone(); + async move { + match pool.upgrade() { + Some(inner) => inner.answer_fetch(&r, &s).await, + None => { + s.abort("not_shared", "the node no longer shares content") + .await + } + } + .map_err(|e| e.to_string()) + } + }); + let procedure = record::own_procedure(&self.self_id, CONTENT_PROCEDURE); + let served = Pool { + inner: self.clone(), + } + .serve(Offer::stream( + *realm, + &procedure, + StreamMode::ServerStream, + answer, + )) + .await?; + self.content.lock().insert( + *realm, + SharedRealm { + _served: served, + roots: HashMap::new(), + chunks: HashMap::new(), + announcements: HashMap::new(), + }, + ); + Ok(()) + } + + /// Answers one fetch, as macula_content_serve does: one DATA body, then + /// the end, or not_shared or malformed. + async fn answer_fetch(&self, realm: &[u8; 32], s: &Stream) -> Result<(), LinkError> { + let args = &s.request().payload; + let mcid: Option = match wire_bytes(args, "mcid") { + Some(b) if b.len() == 50 && b[0] == 2 => b.as_slice().try_into().ok(), + _ => None, + }; + let want = wire_text(args, "want"); + let (Some(mcid), Some(want @ ("root" | "block"))) = (mcid, want.as_deref()) else { + return s + .abort( + "malformed", + "a fetch names one content id and wants root or block", + ) + .await; + }; + let body = self + .content + .lock() + .get(realm) + .and_then(|shared| shared.body(want, &mcid)); + match body { + None => { + s.abort("not_shared", "this node does not share that content") + .await + } + Some(body) => { + s.send_value(body).await?; + s.close().await + } + } + } + + /// Drops `mcid`'s root, its chunks and its announcement from `realm`, and + /// returns the announcement, if there was one. + fn forget_shared(&self, realm: &[u8; 32], mcid: &Mcid) -> Option { + let mut realms = self.content.lock(); + let shared = realms.get_mut(realm)?; + if let Some(Root::Chunked(m)) = shared.roots.remove(mcid) { + for i in 0..m.chunks.len() { + if let Some(c) = chunk_mcid(&m, i) { + shared.chunks.remove(&c); + } + } + } + shared.announcements.remove(mcid) + } + + /// Signs an announcement of `mcid` in `realm`, naming a station this + /// node is linked to and its content procedure, and puts it in the DHT. + async fn announce_content( + self: &Arc, + realm: &[u8; 32], + mcid: &Mcid, + root: &Root, + ) -> Result { + let station = self + .links() + .first() + .map(Link::station_node_id) + .ok_or(PoolError::NoLink(Vec::new()))?; + let mut opts = ContentAnnouncementOptions { + realm_id: *realm, + serving_station: station, + procedure: record::own_procedure(&self.self_id, CONTENT_PROCEDURE), + ttl_ms: ANNOUNCEMENT_TTL.as_millis() as u64, + ..ContentAnnouncementOptions::default() + }; + match root { + Root::Chunked(m) => { + opts.name = String::from_utf8_lossy(&m.name).into_owned(); + opts.size = Some(m.size); + opts.chunk_count = Some(m.chunk_count); + } + Root::Block(b) => opts.size = Some(b.len() as u64), + } + let unsigned = + new_content_announcement(&self.self_id, mcid, &opts).map_err(LinkError::from)?; + let signed = record::sign(&unsigned, &self.opts.identity).map_err(LinkError::from)?; + let wire = record::encode(&signed).map_err(LinkError::from)?; + Pool { + inner: self.clone(), + } + .put_record(&wire) + .await?; + Ok(signed) + } + + /// Every node whose verified announcement of `mcid` in `realm` names a + /// serving station and a content procedure bound to it, shuffled. + async fn content_sharers( + &self, + realm: &[u8; 32], + mcid: &Mcid, + ) -> Result, PoolError> { + let key = record::content_key(mcid).map_err(LinkError::from)?; + let found = match self + .first_answer(|l| async move { l.find_records(&key).await }) + .await + { + Ok((found, _)) => found, + Err(PoolError::Link(LinkError::RecordNotFound)) => Vec::new(), + Err(e) => return Err(e), + }; + let mut out: Vec = found + .iter() + .filter(|v| v.record().record_type == RecordType::CONTENT_ANNOUNCEMENT) + .filter_map(|v| record::read_content_announcement(v.record()).ok()) + .filter(|a| { + a.mcid.as_slice() == mcid.as_slice() + && a.realm_id == *realm + && a.serving_station != [0; 32] + && content_procedure_bound(&a.procedure, &a.announcer_node) + }) + .map(|a| Sharing { + node: a.announcer_node, + station: a.serving_station, + procedure: a.procedure, + }) + .collect(); + shuffle(&mut out); + Ok(out) + } + + /// Fetches `mcid` from one sharer. + async fn fetch_from( + self: &Arc, + realm: &[u8; 32], + s: &Sharing, + mcid: &Mcid, + opts: &ContentOptions, + ) -> Result, PoolError> { + let link = self + .link_to(&s.station, Instant::now() + DEFAULT_CALL_TIMEOUT) + .await?; + let (kind, body) = fetch_one(&link, realm, s, mcid, "root", opts.chunk_timeout).await?; + if kind == "block" { + let bytes = wire_bytes(&body, "bytes").unwrap_or_default(); + if block_mcid(&bytes) != *mcid { + return Err(PoolError::ContentMismatch("the block".into())); + } + return Ok(bytes); + } + let m = manifest::from_wire(body.get("manifest").unwrap_or(&Value::Null)) + .map_err(|e| PoolError::ContentReply(format!("the manifest: {e}")))?; + manifest::verify_mcid(&m, mcid) + .map_err(|e| PoolError::ContentMismatch(format!("the manifest: {e}")))?; + if m.size > opts.max_bytes || m.chunk_count > opts.max_chunks { + return Err(PoolError::ContentTooLarge(format!( + "{} bytes in {} chunks", + m.size, m.chunk_count + ))); + } + manifest::check_whole(&m) + .map_err(|e| PoolError::ContentMismatch(format!("the manifest: {e}")))?; + fetch_chunks(link, *realm, s.clone(), Arc::new(m), opts).await + } +} + +impl SharedRealm { + /// What a fetch of `want` for `mcid` is answered with; a raw root is not + /// served as a chunk. + fn body(&self, want: &str, mcid: &Mcid) -> Option { + let block = |b: &[u8]| { + Value::Map(vec![ + (Value::text("kind"), Value::text("block")), + (Value::text("mcid"), Value::Bytes(mcid.to_vec())), + (Value::text("bytes"), Value::Bytes(b.to_vec())), + ]) + }; + if want == "block" { + return self.chunks.get(mcid).map(|b| block(b)); + } + match self.roots.get(mcid)? { + Root::Block(b) => Some(block(b)), + Root::Chunked(m) => Some(Value::Map(vec![ + (Value::text("kind"), Value::text("manifest")), + (Value::text("mcid"), Value::Bytes(mcid.to_vec())), + (Value::text("manifest"), manifest::to_wire(m)), + ])), + } + } +} + +/// Signs `mcid`'s announcement again at half its lifetime, naming the +/// station the pool is linked to then, until it is unshared or the pool is +/// dropped; a failed renewal is tried again at the next half. +async fn renew_announcement( + pool: Weak, + realm: [u8; 32], + mcid: Mcid, + root: Root, + latest: Arc>, + mut stopped: watch::Receiver, +) { + loop { + tokio::select! { + _ = stopped.wait_for(|s| *s) => return, + _ = tokio::time::sleep(ANNOUNCEMENT_TTL / 2) => {} + } + let Some(inner) = pool.upgrade() else { return }; + if inner.lock().closed { + return; + } + let renewed = tokio::time::timeout( + DEFAULT_CALL_TIMEOUT, + inner.announce_content(&realm, &mcid, &root), + ) + .await; + if let Ok(Ok(record)) = renewed { + *latest.lock().unwrap_or_else(|p| p.into_inner()) = record; + } + } +} + +/// Fetches every chunk of `m`, `opts.parallel` at a time, each checked +/// against its own content id, then the whole against `m`. The first failure +/// ends the fetch: no chunk is asked for after it. +async fn fetch_chunks( + link: Link, + realm: [u8; 32], + s: Sharing, + m: Arc, + opts: &ContentOptions, +) -> Result, PoolError> { + let count = m.chunks.len(); + let parts: Arc>>>> = Arc::new(Mutex::new(vec![None; count])); + let next = Arc::new(AtomicUsize::new(0)); + let (failed_tx, failed) = watch::channel::>(None); + let failed_tx = Arc::new(failed_tx); + let mut workers = tokio::task::JoinSet::new(); + for _ in 0..opts.parallel.max(1).min(count) { + let (link, s, m, parts, next, failed_tx) = ( + link.clone(), + s.clone(), + m.clone(), + parts.clone(), + next.clone(), + failed_tx.clone(), + ); + let timeout = opts.chunk_timeout; + workers.spawn(async move { + loop { + if failed_tx.borrow().is_some() { + return; + } + let i = next.fetch_add(1, Ordering::SeqCst); + let Some(want) = chunk_mcid(&m, i) else { + return; + }; + let fetched = fetch_one(&link, &realm, &s, &want, "block", timeout).await; + let outcome = fetched.and_then(|(_, body)| { + let bytes = wire_bytes(&body, "bytes").unwrap_or_default(); + if block_mcid(&bytes) == want { + Ok(bytes) + } else { + Err(PoolError::ContentMismatch(format!("chunk {i}"))) + } + }); + match outcome { + Ok(bytes) => parts.lock().unwrap_or_else(|p| p.into_inner())[i] = Some(bytes), + Err(e) => { + failed_tx.send_if_modified(|f| { + let first = f.is_none(); + if first { + *f = Some(e); + } + first + }); + return; + } + } + } + }); + } + while workers.join_next().await.is_some() {} + if let Some(e) = failed.borrow().clone() { + return Err(e); + } + let parts = std::mem::take(&mut *parts.lock().unwrap_or_else(|p| p.into_inner())); + let mut whole = Vec::with_capacity(m.size as usize); + for part in parts { + whole.extend(part.ok_or_else(|| PoolError::ContentReply("a chunk never arrived".into()))?); + } + manifest::verify(&m, &whole) + .map_err(|e| PoolError::ContentMismatch(format!("the whole: {e}")))?; + Ok(whole) +} + +/// Asks a sharer for `want` of `mcid` on a stream of its own and reads its +/// one DATA body; the stream is released on every path. +async fn fetch_one( + link: &Link, + realm: &[u8; 32], + s: &Sharing, + mcid: &Mcid, + want: &str, + timeout: Duration, +) -> Result<(String, Value), PoolError> { + let asked = async { + let stream = link + .open_stream(station_link::StreamCall { + realm: *realm, + procedure: s.procedure.clone(), + target: s.node, + mode: StreamMode::ServerStream, + payload: Value::Map(vec![ + (Value::text("mcid"), Value::Bytes(mcid.to_vec())), + (Value::text("want"), Value::text(want)), + ]), + deadline: timeout, + ..station_link::StreamCall::default() + }) + .await?; + let event = stream.recv().await; + let _ = stream.close().await; + read_body(event, mcid, want) + }; + tokio::time::timeout(timeout, asked) + .await + .unwrap_or(Err(PoolError::Link(LinkError::CallTimeout))) +} + +fn read_body( + event: Result, + mcid: &Mcid, + want: &str, +) -> Result<(String, Value), PoolError> { + let body = match event { + Err(LinkError::Stream { code, .. }) if code == "not_shared" => { + return Err(PoolError::NotShared) + } + Err(LinkError::EndOfStream) => { + return Err(PoolError::ContentReply( + "the stream ended with no body".into(), + )) + } + Err(e) => return Err(e.into()), + Ok(StreamEvent::Data { body, .. }) => body, + Ok(_) => return Err(PoolError::ContentReply("a frame that is not DATA".into())), + }; + let kind = wire_text(&body, "kind").unwrap_or_default(); + if wire_bytes(&body, "mcid").as_deref() != Some(mcid.as_slice()) { + return Err(PoolError::ContentReply( + "a body for another content id".into(), + )); + } + let bytes = wire_bytes(&body, "bytes"); + match (kind.as_str(), bytes) { + ("block", Some(b)) if b.len() as u64 > MAX_BLOCK_BYTES => Err(PoolError::ContentTooLarge( + format!("a block of {} bytes", b.len()), + )), + ("block", Some(_)) => Ok((kind, body)), + ("manifest", _) if want == "root" => Ok((kind, body)), + _ => Err(PoolError::ContentReply(format!("kind {kind:?}"))), + } +} + +/// Field `name` of `v` as text, however it arrives. +fn wire_text(v: &Value, name: &str) -> Option { + match v.get(name) { + Some(Value::Text(t)) => Some(t.clone()), + Some(Value::Bytes(b)) => String::from_utf8(b.clone()).ok(), + _ => None, + } +} + +/// Field `name` of `v` as bytes. +fn wire_bytes(v: &Value, name: &str) -> Option> { + match v.get(name) { + Some(Value::Bytes(b)) => Some(b.clone()), + _ => None, + } +} diff --git a/tests/manifest.rs b/tests/manifest.rs new file mode 100644 index 0000000..b747066 --- /dev/null +++ b/tests/manifest.rs @@ -0,0 +1,437 @@ +//! Content manifests as macula 12's macula_manifest builds them, byte for +//! byte (tests/vectors/manifest/erlang_manifests.json holds manifests macula +//! built): fixed-size chunks, SHA-384, a 50-byte content id, a Merkle fold +//! that pairs an odd last hash with itself, and the wire form a manifest +//! travels in. Ported from macula-go v0.12.0's manifest tests. + +use macula_rust::cbor::{self, Value}; +use macula_rust::manifest::{ + block_mcid, check_chunk_hashes, check_whole, chunk_mcid, create_at, from_wire, mcid_for, + mcid_is_chunked, to_wire, verify, verify_mcid, ChunkInfo, Manifest, ManifestError, + DEFAULT_CHUNK_SIZE, +}; +use sha2::{Digest, Sha384}; + +const EMPTY_ROOT_HASH: &str = + "38b060a751ac96384cd9327eb1b1e36a21fdb71114be07434c0cc7bf63f6e1da274edebfe76f65fbd51ad2f14898b95b"; + +fn sha384(data: &[u8]) -> [u8; 48] { + Sha384::digest(data).into() +} + +fn made(data: &[u8], name: &str, chunk_size: u64) -> (Manifest, Vec>) { + create_at(data, name, chunk_size, 1000).unwrap() +} + +fn three_chunks() -> Manifest { + made( + &vec![1; 2 * DEFAULT_CHUNK_SIZE as usize + 3], + "unnamed", + DEFAULT_CHUNK_SIZE, + ) + .0 +} + +fn edited(m: &Manifest, change: impl FnOnce(&mut Manifest)) -> Manifest { + let mut c = m.clone(); + change(&mut c); + c +} + +fn with_own_mcid(mut m: Manifest) -> Manifest { + m.mcid = mcid_for(&m); + m +} + +fn set_field(v: &Value, field: &str, value: Option) -> Value { + let Value::Map(entries) = v else { + panic!("a map") + }; + let mut out: Vec<(Value, Value)> = entries + .iter() + .filter(|(k, _)| *k != Value::text(field)) + .cloned() + .collect(); + if let Some(value) = value { + out.push((Value::text(field), value)); + } + Value::Map(out) +} + +fn set_chunk_field(v: &Value, i: usize, field: &str, value: Value) -> Value { + let Some(Value::List(chunks)) = v.get("chunks") else { + panic!("chunks") + }; + let mut changed = chunks.clone(); + changed[i] = set_field(&changed[i], field, Some(value)); + set_field(v, "chunks", Some(Value::List(changed))) +} + +#[test] +fn a_single_block_s_mcid_is_the_block_s_sha384() { + let data = b"small blob, single block"; + let got = block_mcid(data); + let mut want = vec![2, 0x55]; + want.extend(sha384(data)); + assert_eq!(got.to_vec(), want); + assert!(!mcid_is_chunked(&got)); +} + +#[test] +fn chunks_are_cut_at_the_chunk_size_with_no_empty_trailing_chunk() { + let (m, chunks) = made(&[0x42; 20], "x", 10); + assert_eq!((chunks.len(), m.chunk_count), (2, 2)); + assert!(chunks.iter().all(|c| c.len() == 10)); + let (m, chunks) = made(&[7; 25], "x", 10); + assert_eq!(chunks.len(), 3); + assert_eq!(chunks[2].len(), 5); + assert_eq!(m.size, 25); + assert!( + create_at(b"x", "x", 0, 1000).is_err(), + "a chunk size of zero" + ); +} + +#[test] +fn a_manifest_s_mcid_is_chunked_and_a_block_s_is_not() { + let (m, _) = made(&[1; 100], "x", 10); + assert!(mcid_is_chunked(&m.mcid)); + assert!(!mcid_is_chunked(&block_mcid(&[1; 100]))); +} + +#[test] +fn a_chunk_s_mcid_is_the_block_mcid_of_the_chunk() { + let (m, chunks) = made(&[9; 25], "x", 10); + for (i, chunk) in chunks.iter().enumerate() { + assert_eq!(chunk_mcid(&m, i), Some(block_mcid(chunk))); + } + assert_eq!(chunk_mcid(&m, chunks.len()), None); +} + +#[test] +fn an_odd_last_hash_is_paired_with_itself() { + let (m, _) = made(&[3; 25], "x", 10); + let fold = |l: &[u8; 48], r: &[u8; 48]| sha384(&[l.as_slice(), r.as_slice()].concat()); + let (h0, h1, h2) = (m.chunks[0].hash, m.chunks[1].hash, m.chunks[2].hash); + assert_eq!(m.root_hash, fold(&fold(&h0, &h1), &fold(&h2, &h2))); +} + +#[test] +fn verify_accepts_the_content_and_refuses_tampering() { + let data = vec![5u8; 1000]; + let (m, _) = made(&data, "x", 300); + verify(&m, &data).unwrap(); + let mut tampered = data.clone(); + tampered[500] ^= 0xFF; + assert!(verify(&m, &tampered).is_err()); + assert!(verify(&m, &data[..999]).is_err()); + let zero = edited(&m, |c| c.chunk_size = 0); + assert!(verify(&zero, &data).is_err(), "a chunk size of zero"); +} + +#[test] +fn the_wire_form_round_trips_with_the_name_as_bytes() { + let data: Vec = [0x11, 0x22].repeat(500); + let (m, _) = create_at(&data, "my-file.bin", 300, 1_700_000_000).unwrap(); + let wire = to_wire(&m); + assert_eq!( + wire.get("name"), + Some(&Value::Bytes(b"my-file.bin".to_vec())) + ); + assert_eq!(from_wire(&wire).unwrap(), m); +} + +#[test] +fn verify_mcid_checks_the_fields_macula_checks() { + let m = three_chunks(); + let not_utf8 = with_own_mcid(edited(&m, |c| c.name = vec![0xff, 0xfe])); + let unknown = with_own_mcid(edited(&m, |c| c.hash_algorithm = "sha256".into())); + let cases: Vec<(&str, Manifest, [u8; 50], bool)> = vec![ + ("as created", m.clone(), m.mcid, true), + ( + "another created time, version, own mcid and chunk list", + edited(&m, |c| { + c.created += 1; + c.version = 9; + c.mcid = [0; 50]; + c.chunks.clear(); + }), + m.mcid, + true, + ), + ( + "another name", + edited(&m, |c| c.name = b"other".to_vec()), + m.mcid, + false, + ), + ("another size", edited(&m, |c| c.size += 1), m.mcid, false), + ( + "another chunk count", + edited(&m, |c| c.chunk_count += 1), + m.mcid, + false, + ), + ( + "another root hash", + edited(&m, |c| c.root_hash[0] ^= 1), + m.mcid, + false, + ), + ( + "a name that isn't UTF-8", + not_utf8.clone(), + not_utf8.mcid, + false, + ), + ( + "another hash algorithm", + unknown.clone(), + unknown.mcid, + false, + ), + ]; + for (name, manifest, mcid, ok) in cases { + let got = verify_mcid(&manifest, &mcid); + assert_eq!(got.is_ok(), ok, "{name}: {got:?}"); + if !ok { + assert_eq!(got, Err(ManifestError::McidMismatch), "{name}"); + } + } +} + +#[test] +fn check_whole_refuses_chunks_that_do_not_describe_the_content_whole() { + let m = three_chunks(); + let (empty, _) = made(&[], "unnamed", DEFAULT_CHUNK_SIZE); + let cases: Vec<(&str, Manifest, bool)> = vec![ + ("as created", m.clone(), true), + ("empty content, as created", empty, true), + ( + "a chunk cut short before the last, with the count still right", + edited(&m, |c| { + c.chunks[0].size -= 1; + c.chunks[1].offset -= 1; + c.chunks[2].offset -= 1; + c.chunks[2].size += 1; + }), + false, + ), + ( + "a chunk counted that isn't listed", + edited(&m, |c| c.chunk_count += 1), + false, + ), + ( + "chunks out of index order", + edited(&m, |c| { + c.chunks[0].index = 1; + c.chunks[1].index = 0; + }), + false, + ), + ( + "a gap between chunks", + edited(&m, |c| c.chunks[1].offset += 1), + false, + ), + ( + "an empty chunk", + edited(&m, |c| c.chunks[2].size = 0), + false, + ), + ( + "a chunk larger than the chunk size", + edited(&m, |c| c.chunks[2].size = c.chunk_size + 1), + false, + ), + ( + "a size the chunks don't add up to", + edited(&m, |c| c.size += 1), + false, + ), + ( + "a chunk size of zero", + edited(&m, |c| c.chunk_size = 0), + false, + ), + ]; + for (name, manifest, whole) in cases { + let got = check_whole(&manifest); + assert_eq!(got.is_ok(), whole, "{name}: {got:?}"); + if !whole { + assert!( + matches!(got, Err(ManifestError::NotWhole(_))), + "{name}: {got:?}" + ); + } + } +} + +#[test] +fn chunk_hashes_must_make_the_root_hash() { + let m = three_chunks(); + check_chunk_hashes(&m).unwrap(); + let swapped = edited(&m, |c| c.chunks.swap(0, 2)); + assert_eq!( + check_chunk_hashes(&swapped), + Err(ManifestError::ChunkHashes) + ); +} + +#[test] +fn empty_content_has_one_whole_form() { + let (empty, _) = made(&[], "unnamed", DEFAULT_CHUNK_SIZE); + let back = from_wire(&to_wire(&empty)).unwrap(); + assert_eq!((back.size, back.chunk_count, back.chunks.len()), (0, 0, 0)); + assert!(back.chunk_size > 0); + assert_eq!(hex::encode(back.root_hash), EMPTY_ROOT_HASH); + verify_mcid(&back, &empty.mcid).unwrap(); + verify(&back, &[]).unwrap(); + let one_empty_chunk = edited(&empty, |c| { + c.chunk_count = 1; + c.chunks = vec![ChunkInfo { + index: 0, + offset: 0, + size: 0, + hash: c.root_hash, + }]; + }); + assert!(matches!( + from_wire(&to_wire(&one_empty_chunk)), + Err(ManifestError::NotWhole(_)) + )); +} + +#[test] +fn only_sha384_is_read_as_the_hash_algorithm() { + let (m, _) = made(b"some content", "unnamed", DEFAULT_CHUNK_SIZE); + let wire = to_wire(&m); + let cases: Vec<(&str, Option, bool)> = vec![ + ("missing", None, false), + ("sha384 as text", Some(Value::text("sha384")), true), + ( + "sha384 as bytes", + Some(Value::Bytes(b"sha384".to_vec())), + true, + ), + ("blake3 as text", Some(Value::text("blake3")), false), + ("sha256 as text", Some(Value::text("sha256")), false), + ( + "sha256 as bytes", + Some(Value::Bytes(b"sha256".to_vec())), + false, + ), + ("an unknown name", Some(Value::text("md5")), false), + ]; + for (name, value, ok) in cases { + let got = from_wire(&set_field(&wire, "hash_algorithm", value)); + assert_eq!(got.is_ok(), ok, "{name}: {got:?}"); + if let Ok(read) = got { + assert_eq!(read.hash_algorithm, "sha384"); + } + } +} + +#[test] +fn a_number_too_large_for_its_field_is_refused() { + let m = three_chunks(); + let wire = to_wire(&m); + from_wire(&wire).unwrap(); + let past32 = |n: u64| Value::Int(i128::from(n) + (1 << 32)); + let last = m.chunks.last().unwrap().clone(); + let i = last.index as usize; + let cases = [ + ( + "version", + set_field(&wire, "version", Some(past32(u64::from(m.version)))), + ), + ( + "a chunk index", + set_chunk_field(&wire, i, "index", past32(last.index)), + ), + ( + "a chunk offset", + set_chunk_field(&wire, i, "offset", past32(last.offset)), + ), + ( + "a chunk size", + set_chunk_field(&wire, i, "size", past32(last.size)), + ), + ( + "chunk_size", + set_field(&wire, "chunk_size", Some(past32(m.chunk_size))), + ), + ( + "chunk_count", + set_field(&wire, "chunk_count", Some(past32(m.chunk_count))), + ), + ]; + for (name, v) in cases { + assert!(from_wire(&v).is_err(), "{name} plus 2^32 was accepted"); + } + let negative = set_field(&wire, "size", Some(Value::Int(-1))); + assert!(from_wire(&negative).is_err(), "a negative size"); + // created is read by nothing that checks it, so its own range is its only + // guard. + for created in [Value::Int(-1), Value::Int(i128::from(i64::MAX) + 1)] { + let v = set_field(&wire, "created", Some(created.clone())); + assert!(from_wire(&v).is_err(), "created {created:?}"); + } +} + +#[test] +fn manifests_match_macula_s_byte_for_byte() { + let raw = std::fs::read(concat!( + env!("CARGO_MANIFEST_DIR"), + "/tests/vectors/manifest/erlang_manifests.json" + )) + .unwrap(); + let fixture: serde_json::Value = serde_json::from_slice(&raw).unwrap(); + let text = |v: &serde_json::Value| v.as_str().unwrap().to_string(); + let pattern = |n: usize| (0..n).map(|i| (i % 251) as u8).collect::>(); + let manifests = fixture["manifests"].as_array().unwrap(); + assert_eq!(manifests.len(), 5); + for want in manifests { + let name = String::from_utf8(hex::decode(text(&want["name"])).unwrap()).unwrap(); + let size = want["size"].as_u64().unwrap() as usize; + let chunk_size = want["chunk_size"].as_u64().unwrap(); + let (m, chunks) = create_at(&pattern(size), &name, chunk_size, 1_789_000_000).unwrap(); + assert_eq!(hex::encode(m.mcid), text(&want["mcid_hex"]), "{name}: mcid"); + assert_eq!( + hex::encode(m.root_hash), + text(&want["root_hex"]), + "{name}: root" + ); + let hashes = want["chunk_hashes"].as_array().unwrap(); + let mcids = want["chunk_mcids"].as_array().unwrap(); + assert_eq!( + (chunks.len(), m.chunks.len()), + (hashes.len(), hashes.len()), + "{name}" + ); + for (i, c) in m.chunks.iter().enumerate() { + assert_eq!( + hex::encode(c.hash), + text(&hashes[i]), + "{name}: chunk {i} hash" + ); + assert_eq!( + hex::encode(chunk_mcid(&m, i).unwrap()), + text(&mcids[i]), + "{name}: chunk {i} mcid" + ); + } + let wire = hex::decode(text(&want["wire_hex"])).unwrap(); + assert_eq!( + hex::encode(cbor::encode(&to_wire(&m)).unwrap()), + hex::encode(&wire), + "{name}: wire" + ); + let read = from_wire(&cbor::decode(&wire).unwrap()).unwrap(); + verify_mcid(&read, &read.mcid).unwrap(); + assert_eq!(read.mcid, m.mcid, "{name}: read back"); + } + let block = block_mcid(&hex::decode(text(&fixture["block"]["data_hex"])).unwrap()); + assert_eq!(hex::encode(block), text(&fixture["block"]["mcid_hex"])); +} diff --git a/tests/pool_content.rs b/tests/pool_content.rs new file mode 100644 index 0000000..79099f5 --- /dev/null +++ b/tests/pool_content.rs @@ -0,0 +1,409 @@ +//! Node-served content through the pool (D27): a node keeps what it shares, +//! serves it on its own `~/content_v1` server stream and announces it +//! in the DHT; a fetcher finds the announcements, dials each sharer's station +//! and checks everything it receives against the content id it asked for. No +//! realm key is needed on either side. Ported from macula-go v0.12.0's +//! content tests. + +mod common; + +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::Arc; +use std::time::{Duration, Instant}; + +use common::lab::{Lab, LabStation}; +use macula_rust::cbor::Value; +use macula_rust::frame::StreamMode; +use macula_rust::manifest::{ + self, block_mcid, chunk_mcid, mcid_is_chunked, Mcid, DEFAULT_CHUNK_SIZE, +}; +use macula_rust::node_key::{NodeKey, PUZZLE_DIFFICULTY}; +use macula_rust::pool::{ + content_procedure_bound, ContentOptions, Offer, Opts, Pool, PoolError, CONTENT_PROCEDURE, +}; +use macula_rust::profile::Profile; +use macula_rust::record::{self, new_content_announcement, ContentAnnouncementOptions}; +use macula_rust::station_link::stream_handler; + +const REALM: [u8; 32] = [0x0c; 32]; + +fn key() -> Arc { + Arc::new(NodeKey::generate_identity(Profile::PqPure, PUZZLE_DIFFICULTY).unwrap()) +} + +/// A pool as `key` on one station, trusting no realm. +async fn connect(key: &Arc, s: &LabStation) -> Pool { + let mut opts = Opts::new(key.clone()); + opts.respawn_delay = Duration::from_millis(100); + opts.connect_timeout = Duration::from_secs(15); + Pool::connect(vec![s.seed()], opts).await.unwrap() +} + +/// A sharer on one station and a fetcher on another, the two sharing a DHT. +async fn content_pools(lab: &Lab, name: &str) -> (Pool, Pool, LabStation, LabStation) { + let sharing = lab.station(&format!("{name} sharing")); + let fetching = lab.station(&format!("{name} fetching")); + lab.share(&[&sharing, &fetching]); + let sharer = connect(&key(), &sharing).await; + let fetcher = connect(&key(), &fetching).await; + (sharer, fetcher, sharing, fetching) +} + +fn pattern(n: usize) -> Vec { + (0..n).map(|i| (i % 251) as u8).collect() +} + +fn block_body(mcid: &[u8], bytes: &[u8]) -> Value { + Value::Map(vec![ + (Value::text("kind"), Value::text("block")), + (Value::text("mcid"), Value::Bytes(mcid.to_vec())), + (Value::text("bytes"), Value::Bytes(bytes.to_vec())), + ]) +} + +fn manifest_body(mcid: &Mcid, m: &manifest::Manifest) -> Value { + Value::Map(vec![ + (Value::text("kind"), Value::text("manifest")), + (Value::text("mcid"), Value::Bytes(mcid.to_vec())), + (Value::text("manifest"), manifest::to_wire(m)), + ]) +} + +/// Announces `mcid` as shared by `key`'s node on `procedure`, served through +/// `station`. +async fn announce(pool: &Pool, key: &NodeKey, mcid: &Mcid, station: &LabStation, procedure: &str) { + let unsigned = new_content_announcement( + &key.node_id().unwrap(), + mcid, + &ContentAnnouncementOptions { + realm_id: REALM, + serving_station: station.node_id, + procedure: procedure.to_string(), + ..ContentAnnouncementOptions::default() + }, + ) + .unwrap(); + let wire = record::encode(&record::sign(&unsigned, key).unwrap()).unwrap(); + pool.put_record(&wire).await.unwrap(); +} + +/// A node on `station` that serves, on its own content procedure, whatever +/// `answer` gives for a fetch's want and content id, and announces `mcid` as +/// shared there. Returns its pool, kept alive by the caller. +async fn liar_for( + station: &LabStation, + mcid: &Mcid, + answer: impl Fn(String, Vec) -> Value + Send + Sync + 'static, +) -> Pool { + let k = key(); + let liar = connect(&k, station).await; + let own = record::own_procedure(&liar.node_id(), CONTENT_PROCEDURE); + let answer = Arc::new(answer); + liar.serve(Offer::stream( + REALM, + &own, + StreamMode::ServerStream, + stream_handler(move |s| { + let answer = answer.clone(); + async move { + let payload = s.request().payload.clone(); + let want = match payload.get("want") { + Some(Value::Text(t)) => t.clone(), + _ => String::new(), + }; + let asked = match payload.get("mcid") { + Some(Value::Bytes(b)) => b.clone(), + _ => Vec::new(), + }; + s.send_value(answer(want, asked)) + .await + .map_err(|e| e.to_string())?; + s.close().await.map_err(|e| e.to_string()) + } + }), + )) + .await + .unwrap(); + announce(&liar, &k, mcid, station, &own).await; + liar +} + +fn chunk_of(m: &manifest::Manifest, chunks: &[Vec], asked: &[u8]) -> Vec { + (0..chunks.len()) + .find(|i| chunk_mcid(m, *i).is_some_and(|c| c.as_slice() == asked)) + .map(|i| chunks[i].clone()) + .unwrap_or_default() +} + +async fn within(limit: Duration, what: &str, mut ok: impl FnMut() -> bool) { + let deadline = Instant::now() + limit; + while Instant::now() < deadline { + if ok() { + return; + } + tokio::time::sleep(Duration::from_millis(20)).await; + } + panic!("not within {limit:?}: {what}"); +} + +#[tokio::test(flavor = "multi_thread")] +async fn content_is_fetched_from_the_node_that_shares_it() { + let lab = Lab::start(Profile::PqPure); + let (sharer, fetcher, sharing, _) = content_pools(&lab, "content").await; + for (name, data) in [ + ("a raw block", pattern(10_000)), + ("a chunked blob", pattern(600_000)), + ] { + let mcid = sharer + .share_content(&REALM, &data, "blob.bin") + .await + .unwrap(); + assert_eq!( + mcid_is_chunked(&mcid), + data.len() as u64 > DEFAULT_CHUNK_SIZE, + "{name}: chunked" + ); + let got = fetcher + .get_content(&REALM, &mcid, ContentOptions::default()) + .await + .unwrap(); + assert!(got == data, "{name}: {} bytes back", got.len()); + assert!( + lab.connected(&sharing, &fetcher.node_id()), + "{name}: the fetcher dialed the sharer's station" + ); + sharer.unshare_content(&REALM, &mcid).await.unwrap(); + let after = fetcher + .get_content(&REALM, &mcid, ContentOptions::default()) + .await; + assert!( + matches!(after, Err(PoolError::NotShared)), + "{name}: after unsharing {after:?}" + ); + } + let never = fetcher + .get_content( + &[0x0d; 32], + &block_mcid(b"never shared"), + ContentOptions::default(), + ) + .await; + assert!(matches!(never, Err(PoolError::NotShared)), "{never:?}"); +} + +#[tokio::test(flavor = "multi_thread")] +async fn a_fetch_refuses_content_over_its_bounds() { + let lab = Lab::start(Profile::PqPure); + let (sharer, fetcher, _, _) = content_pools(&lab, "bounds").await; + let mcid = sharer + .share_content(&REALM, &pattern(600_000), "big.bin") + .await + .unwrap(); + let bytes = ContentOptions { + max_bytes: 500_000, + ..ContentOptions::default() + }; + let refused = fetcher.get_content(&REALM, &mcid, bytes).await; + assert!( + matches!(&refused, Err(PoolError::ContentUnavailable(tried)) if matches!(tried[0].1, PoolError::ContentTooLarge(_))), + "{refused:?}" + ); + let chunks = ContentOptions { + max_chunks: 2, + ..ContentOptions::default() + }; + let refused = fetcher.get_content(&REALM, &mcid, chunks).await; + assert!( + matches!(&refused, Err(PoolError::ContentUnavailable(tried)) if matches!(tried[0].1, PoolError::ContentTooLarge(_))), + "{refused:?}" + ); +} + +#[tokio::test(flavor = "multi_thread")] +async fn nothing_a_sharer_says_is_trusted() { + let lab = Lab::start(Profile::PqPure); + let (_, fetcher, sharing, _) = content_pools(&lab, "liar").await; + let asked = b"the content asked for"; + let mcid = block_mcid(asked); + let wanted = mcid; + let liar_key = key(); + let liar = connect(&liar_key, &sharing).await; + let own = record::own_procedure(&liar.node_id(), CONTENT_PROCEDURE); + liar.serve(Offer::stream( + REALM, + &own, + StreamMode::ServerStream, + stream_handler(move |s| async move { + s.send_value(block_body(&wanted, b"something else")) + .await + .map_err(|e| e.to_string())?; + s.close().await.map_err(|e| e.to_string()) + }), + )) + .await + .unwrap(); + // An announcement pointing at another node's procedure is never tried. + announce( + &liar, + &liar_key, + &mcid, + &sharing, + &record::own_procedure(&fetcher.node_id(), CONTENT_PROCEDURE), + ) + .await; + let got = fetcher + .get_content(&REALM, &mcid, ContentOptions::default()) + .await; + assert!(matches!(got, Err(PoolError::NotShared)), "{got:?}"); + // A sharer answering with other bytes is refused. + announce(&liar, &liar_key, &mcid, &sharing, &own).await; + let got = fetcher + .get_content(&REALM, &mcid, ContentOptions::default()) + .await; + assert!( + matches!(&got, Err(PoolError::ContentUnavailable(tried)) if matches!(tried[0].1, PoolError::ContentMismatch(_))), + "{got:?}" + ); +} + +#[test] +fn a_content_procedure_is_bound_to_its_announcer() { + let mut node = [0u8; 32]; + node[0] = 0xab; + let hex_node = format!("ab{}", "00".repeat(31)); + let cases = [ + (format!("~{hex_node}/content_v1"), true), + (format!("acme/content_v1_{hex_node}"), true), + (format!("~{}/content_v1", hex_node.to_uppercase()), false), + (format!("~{hex_node}/content_v2"), false), + (format!("acme/content_v1_{}", "cd".repeat(32)), false), + (format!("/content_v1_{hex_node}"), false), + (format!("_/content_v1_{hex_node}"), false), + (format!("~x/content_v1_{hex_node}"), false), + (format!("acme/sub/content_v1_{hex_node}"), false), + ]; + for (procedure, want) in cases { + assert_eq!( + content_procedure_bound(&procedure, &node), + want, + "{procedure}" + ); + } +} + +#[tokio::test(flavor = "multi_thread")] +async fn an_announcement_names_where_the_content_is_served_and_streams_are_released() { + let lab = Lab::start(Profile::PqPure); + let (sharer, fetcher, sharing, fetching) = content_pools(&lab, "announce").await; + let mcid = sharer + .share_content(&REALM, &pattern(300_000), "two.bin") + .await + .unwrap(); + let (found, _) = fetcher + .find_records(&record::content_key(&mcid).unwrap()) + .await + .unwrap(); + assert_eq!(found.len(), 1); + let a = record::read_content_announcement(found[0].record()).unwrap(); + assert_eq!(a.announcer_node, sharer.node_id()); + assert_eq!(a.realm_id, REALM); + assert_eq!(a.serving_station, sharing.node_id); + assert_eq!( + a.procedure, + record::own_procedure(&sharer.node_id(), CONTENT_PROCEDURE) + ); + assert_eq!(a.chunk_count, Some(2)); + let one_at_a_time = ContentOptions { + parallel: 1, + ..ContentOptions::default() + }; + fetcher + .get_content(&REALM, &mcid, one_at_a_time) + .await + .unwrap(); + within( + Duration::from_secs(2), + "every content stream released", + || lab.relayed(&sharing) == 0 && lab.relayed(&fetching) == 0, + ) + .await; +} + +#[tokio::test(flavor = "multi_thread")] +async fn a_manifest_and_its_chunks_are_checked_against_their_content_ids() { + let lab = Lab::start(Profile::PqPure); + let (_, fetcher, sharing, _) = content_pools(&lab, "manifests").await; + let (real, real_chunks) = + manifest::create(&pattern(600_000), "real.bin", DEFAULT_CHUNK_SIZE).unwrap(); + let (other, other_chunks) = manifest::create( + &pattern(700_000)[100_000..], + "other.bin", + DEFAULT_CHUNK_SIZE, + ) + .unwrap(); + + // Another content's manifest for this content id. + let (real_mcid, other_m) = (real.mcid, other.clone()); + let _first = liar_for(&sharing, &real.mcid, move |want, asked| { + if want == "root" { + return manifest_body(&real_mcid, &other_m); + } + block_body(&asked, &chunk_of(&other_m, &other_chunks, &asked)) + }) + .await; + let got = fetcher + .get_content(&REALM, &real.mcid, ContentOptions::default()) + .await; + assert!( + matches!(&got, Err(PoolError::ContentUnavailable(tried)) if matches!(&tried[0].1, PoolError::ContentMismatch(why) if why.contains("manifest"))), + "{got:?}" + ); + + // The right manifest, then chunks that are not theirs: the first one ends + // the fetch, before another is asked for. + let asked_count = Arc::new(AtomicUsize::new(0)); + let (counter, real_m) = (asked_count.clone(), real.clone()); + let _second = liar_for(&sharing, &real.mcid, move |want, asked| { + if want == "root" { + return manifest_body(&real_m.mcid, &real_m); + } + counter.fetch_add(1, Ordering::SeqCst); + let mut b = chunk_of(&real_m, &real_chunks, &asked); + b[0] ^= 0xff; + block_body(&asked, &b) + }) + .await; + let one_at_a_time = ContentOptions { + parallel: 1, + ..ContentOptions::default() + }; + let got = fetcher.get_content(&REALM, &real.mcid, one_at_a_time).await; + let chunk_mismatch = |tried: &Vec<(_, PoolError)>| { + tried + .iter() + .any(|(_, e)| matches!(e, PoolError::ContentMismatch(why) if why.contains("chunk"))) + }; + assert!( + matches!(&got, Err(PoolError::ContentUnavailable(tried)) if chunk_mismatch(tried)), + "{got:?}" + ); + assert_eq!( + asked_count.load(Ordering::SeqCst), + 1, + "chunks asked for after the first did not match" + ); +} + +#[tokio::test(flavor = "multi_thread")] +async fn content_is_fetched_only_in_the_realm_it_is_shared_in() { + let lab = Lab::start(Profile::PqPure); + let (sharer, fetcher, _, _) = content_pools(&lab, "realms").await; + let mcid = sharer + .share_content(&REALM, &pattern(1_000), "r.bin") + .await + .unwrap(); + let got = fetcher + .get_content(&[0x0d; 32], &mcid, ContentOptions::default()) + .await; + assert!(matches!(got, Err(PoolError::NotShared)), "{got:?}"); +}