Qualify durable purchase and media primitives and preserve app launch paths

This commit is contained in:
archipelago
2026-10-06 20:50:44 -04:00
parent a876dc3d0b
commit b52214f7a0
31 changed files with 4417 additions and 83 deletions
@@ -349,6 +349,18 @@ impl ApiHandler {
));
}
if let Err(error) =
content_server::ensure_payment_source_available(&self.config.data_dir, item).await
{
return Ok(build_response(
StatusCode::CONFLICT,
"application/json",
hyper::Body::from(serde_json::to_vec(
&serde_json::json!({ "error": error.to_string(), "payment_started": false }),
)?),
));
}
let memo = format!("Archipelago peer file {content_id}");
match self
.rpc_handler
@@ -468,6 +480,30 @@ impl ApiHandler {
}
};
// Match the node wallet's existing sendcoins minimum before exposing a
// payable address for an amount its own payment flow cannot broadcast.
if let Err(error) = content_server::validate_onchain_payment_price(price_sats) {
return Ok(build_response(
StatusCode::BAD_REQUEST,
"application/json",
hyper::Body::from(serde_json::to_vec(
&serde_json::json!({ "error": error.to_string(), "payment_started": false }),
)?),
));
}
if let Err(error) =
content_server::ensure_payment_source_available(&self.config.data_dir, item).await
{
return Ok(build_response(
StatusCode::CONFLICT,
"application/json",
hyper::Body::from(serde_json::to_vec(
&serde_json::json!({ "error": error.to_string(), "payment_started": false }),
)?),
));
}
match self.rpc_handler.new_onchain_address().await {
Ok(address) if !address.is_empty() => {
crate::content_invoice::record_pending_method(
+7
View File
@@ -1,6 +1,7 @@
mod blob;
mod cdp;
mod content;
mod registered_media;
mod dwn;
mod model_proxy;
mod node_message;
@@ -584,6 +585,12 @@ impl ApiHandler {
Self::handle_blob_download(&self.blob_store, p, &query_string).await
}
// Immutable registered rentals use durable seller receipts and their
// first-open window, never legacy mutable filename shares.
(Method::GET, p) if p.starts_with("/content/registered_") && p.contains("/rental/") => {
self.handle_registered_rental(p, &headers).await
}
// Content preview — degraded previews for paid content (no auth, no payment)
(Method::GET, p) if p.starts_with("/content/") && p.ends_with("/preview") => {
Self::handle_content_preview(p, &self.config).await
@@ -0,0 +1,310 @@
//! Authenticated immutable rental streaming, separate from legacy mutable shares.
use super::{build_response, ApiHandler};
use crate::{content_server::ByteRange, identity::NodeIdentity, registered_media::OpenedMedia};
use anyhow::{Context, Result};
use hyper::{Body, HeaderMap, Response, StatusCode};
use std::sync::Arc;
use tokio::io::{AsyncReadExt, AsyncSeekExt};
fn route(path: &str) -> Result<(&str, &str)> {
let (content, purchase) = path
.strip_prefix("/content/")
.and_then(|value| value.split_once("/rental/"))
.context("Invalid rental route")?;
anyhow::ensure!(
content.starts_with("registered_") && !content.contains('/') && !purchase.contains('/'),
"Invalid rental identifiers"
);
let id = uuid::Uuid::parse_str(purchase)?;
anyhow::ensure!(
id.to_string() == purchase && id.get_version_num() == 4,
"Invalid purchase identifier"
);
Ok((content, purchase))
}
fn bounds(range: Option<ByteRange>, total: u64) -> Result<Option<(u64, u64)>> {
let Some(range) = range else {
anyhow::ensure!(total > 0, "Registered media is empty");
return Ok(None);
};
let last = total.checked_sub(1).context("Registered media is empty")?;
let (start, end) = match range {
ByteRange::From { start, end } => (start, end.unwrap_or(last).min(last)),
ByteRange::Suffix(count) => {
anyhow::ensure!(count > 0, "Invalid suffix range");
(total.saturating_sub(count), last)
}
};
anyhow::ensure!(start <= end && start < total, "Invalid rental byte range");
Ok(Some((start, end)))
}
fn clock() -> u64 {
u64::try_from(chrono::Utc::now().timestamp()).unwrap_or(0)
}
fn denied(message: &'static str) -> Response<Body> {
build_response(StatusCode::FORBIDDEN, "text/plain", Body::from(message))
}
impl ApiHandler {
pub(super) async fn handle_registered_rental(
&self,
path: &str,
headers: &HeaderMap,
) -> Result<Response<Body>> {
let (content, purchase) = match route(path) {
Ok(ids) => ids,
Err(_) => {
return Ok(build_response(
StatusCode::BAD_REQUEST,
"text/plain",
Body::from("Invalid rental route"),
))
}
};
let audience = crate::identity::did_key_from_pubkey_hex(&self.self_pubkey_hex)?;
let buyer = match crate::content_auth::incoming(
headers,
&audience,
path,
chrono::Utc::now().timestamp(),
) {
Ok(Some(buyer)) => buyer,
_ => return Ok(denied("Authenticated node proof is required")),
};
let capability = match headers
.get("x-content-capability")
.and_then(|v| v.to_str().ok())
{
Some(value) if value.len() == 64 => value.to_owned(),
_ => return Ok(denied("Original purchase capability is required")),
};
let requested_range = match headers.get("range") {
None => None,
Some(value) => match value
.to_str()
.ok()
.and_then(crate::content_server::parse_range_header)
{
Some(range) => Some(range),
None => {
return Ok(build_response(
StatusCode::RANGE_NOT_SATISFIABLE,
"text/plain",
Body::from("Invalid byte range"),
))
}
},
};
let identity =
Arc::new(NodeIdentity::load_existing(&self.config.data_dir.join("identity")).await?);
anyhow::ensure!(identity.did_key()? == audience, "Node identity changed");
let data = self.config.data_dir.clone();
let selected = content.to_owned();
let key = identity.clone();
let metadata = tokio::task::spawn_blocking(move || {
crate::registered_media::registered_terms(&data, &key, &selected)
})
.await?;
let (receipt, _) = match metadata {
Ok(value) => value,
Err(_) => {
return Ok(build_response(
StatusCode::NOT_FOUND,
"text/plain",
Body::from("Registered content is unavailable"),
))
}
};
let total = receipt.size_bytes.parse::<u64>()?;
let range = match bounds(requested_range, total) {
Ok(value) => value,
Err(_) => {
return Ok(Response::builder()
.status(StatusCode::RANGE_NOT_SATISFIABLE)
.header("Content-Range", format!("bytes */{total}"))
.body(Body::from("Invalid byte range"))?)
}
};
// All malformed/out-of-bounds requests are rejected before first-open
// rental creation. No payment or new receipt is attempted by this route.
let opened = match crate::registered_media::open_paid(
self.config.data_dir.clone(),
identity,
content.into(),
purchase.into(),
buyer,
capability,
)
.await
{
Ok(opened) => opened,
Err(_) => {
return Ok(denied(
"This purchase is not settled, does not match, or its rental has expired",
))
}
};
rental_response(opened, range, Arc::new(clock)).await
}
}
async fn rental_response(
opened: OpenedMedia,
range: Option<(u64, u64)>,
now: Arc<dyn Fn() -> u64 + Send + Sync>,
) -> Result<Response<Body>> {
anyhow::ensure!(
opened.still_authorized(now()),
"Rental expired before streaming"
);
let total = opened.size_bytes;
let started = opened.started_at;
let expires = opened.expires_at;
let (start, length) = range.map_or((0, total), |(start, end)| (start, end - start + 1));
let mut file = tokio::fs::File::from_std(opened.file);
file.seek(std::io::SeekFrom::Start(start)).await?;
let chunks = futures_util::stream::try_unfold(
(file, length, now),
move |(mut file, left, now)| async move {
if left == 0 {
return Ok::<_, std::io::Error>(None);
}
let instant = now();
if instant < started || instant >= expires {
return Err(std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"Rental window ended",
));
}
let mut bytes = vec![0; left.min(64 * 1024) as usize];
let count = tokio::time::timeout(
std::time::Duration::from_secs(expires - instant),
file.read(&mut bytes),
)
.await
.map_err(|_| {
std::io::Error::new(std::io::ErrorKind::TimedOut, "Rental window ended")
})??;
let instant = now();
if instant < started || instant >= expires {
return Err(std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"Rental window ended",
));
}
if count == 0 {
return Err(std::io::Error::new(
std::io::ErrorKind::UnexpectedEof,
"Registered snapshot ended early",
));
}
bytes.truncate(count);
Ok(Some((bytes, (file, left - count as u64, now))))
},
);
let mut response = Response::builder()
.status(if range.is_some() {
StatusCode::PARTIAL_CONTENT
} else {
StatusCode::OK
})
.header("Content-Type", opened.mime_type)
.header("Content-Length", length)
.header("Accept-Ranges", "bytes")
.header("X-Content-Type-Options", "nosniff")
.header("Cache-Control", "private, no-store")
.header("X-Rental-Expires-At", expires);
if let Some((start, end)) = range {
response = response.header("Content-Range", format!("bytes {start}-{end}/{total}"));
}
Ok(response.body(Body::wrap_stream(chunks))?)
}
#[cfg(test)]
mod tests {
use super::*;
use hyper::body::HttpBody;
use std::sync::atomic::{AtomicU64, Ordering};
#[test]
fn invalid_routes_and_ranges_cannot_reach_rental_creation() {
let id = uuid::Uuid::new_v4();
assert!(route(&format!("/content/registered_{id}/rental/{id}")).is_ok());
for path in [
format!("/content/registered_{id}/rental/{id}/extra"),
format!("/content/../rental/{id}"),
format!("/content/registered_{id}/rental/not-a-purchase"),
] {
assert!(route(&path).is_err());
}
assert!(bounds(
Some(ByteRange::From {
start: 20,
end: None
}),
20
)
.is_err());
assert!(bounds(
Some(ByteRange::From {
start: 9,
end: Some(8)
}),
20
)
.is_err());
assert_eq!(
bounds(Some(ByteRange::Suffix(5)), 20).unwrap(),
Some((15, 19))
);
}
fn opened(size: u64) -> OpenedMedia {
let file = tempfile::tempfile().unwrap();
file.set_len(size).unwrap();
OpenedMedia {
file,
size_bytes: size,
mime_type: "video/mp4".into(),
started_at: 1000,
expires_at: 1060,
}
}
#[tokio::test]
async fn bounded_stream_stops_at_persisted_deadline_without_restarting_window() {
let clock = Arc::new(AtomicU64::new(1000));
let read_clock = clock.clone();
let mut response = rental_response(
opened(200_000),
None,
Arc::new(move || read_clock.load(Ordering::SeqCst)),
)
.await
.unwrap();
assert_eq!(response.headers()["x-rental-expires-at"], "1060");
assert_eq!(
response.body_mut().data().await.unwrap().unwrap().len(),
64 * 1024
);
clock.store(1060, Ordering::SeqCst);
assert!(response.body_mut().data().await.unwrap().is_err());
}
#[tokio::test]
async fn suffix_response_has_exact_length_and_expired_or_rollback_stream_denies() {
let mut response = rental_response(opened(20), Some((15, 19)), Arc::new(|| 1000))
.await
.unwrap();
assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
assert_eq!(response.headers()["content-range"], "bytes 15-19/20");
assert_eq!(
hyper::body::to_bytes(response.body_mut())
.await
.unwrap()
.len(),
5
);
assert!(rental_response(opened(20), None, Arc::new(|| 1060))
.await
.is_err());
assert!(rental_response(opened(20), None, Arc::new(|| 999))
.await
.is_err());
}
}
+66 -6
View File
@@ -896,6 +896,62 @@ impl RpcHandler {
Ok(result)
}
/// Owner-authenticated local recovery lookup. No mint, peer request or
/// wallet mutation occurs here, and private tokens/capabilities are omitted.
pub(super) async fn handle_content_payment_status(
&self,
params: Option<serde_json::Value>,
) -> Result<serde_json::Value> {
let params = params.context("Missing payment lookup parameters")?;
let onion = params
.get("onion")
.and_then(|v| v.as_str())
.context("Missing seller address")?;
let content_id = params
.get("content_id")
.and_then(|v| v.as_str())
.context("Missing content identifier")?;
anyhow::ensure!(is_valid_v3_onion(onion), "Invalid seller address");
crate::content_owned::validate_identity(onion, content_id)?;
let peer =
crate::federation::load_unique_payment_peer(&self.config.data_dir, onion).await?;
let (data, _) = self.state_manager.get_snapshot().await;
let buyer_did = crate::identity::did_key_from_pubkey_hex(&data.server_info.pubkey)?;
let journal = crate::content_purchase::Journal::open(&self.config.data_dir).await?;
let records = if let Some(id) = params.get("operation_id") {
let id = id
.as_str()
.context("Invalid purchase operation identifier")?;
match journal.buyer(id).await? {
Some(record) => {
anyhow::ensure!(
record.contract.buyer_did == buyer_did
&& record.contract.seller_did == peer.did
&& record.contract.content_id == content_id,
"Purchase belongs to another buyer, seller or content item"
);
vec![record]
}
None => Vec::new(),
}
} else {
journal
.find_buyers(&buyer_did, &peer.did, content_id)
.await?
};
let attempts: Vec<_> = records
.iter()
.map(|record| record.public_status())
.collect();
Ok(serde_json::json!({
"state": if attempts.is_empty() { "unknown" } else { "recorded" },
"attempts": attempts,
// Absence is not evidence that a legacy payment failed/reclaimed.
"can_start_new_payment": false,
"legacy_recovery_unresolved": attempts.is_empty(),
}))
}
/// Buyer side (#46): ask the selling node to mint a Lightning invoice for a
/// paid item so the buyer can pay from any external wallet. Returns the
/// bolt11 invoice + payment hash to render as a QR and poll for settlement.
@@ -943,9 +999,11 @@ impl RpcHandler {
};
if !response.status().is_success() {
return Ok(serde_json::json!({
"error": format!("Seller could not create an invoice ({}).", response.status())
}));
let status = response.status();
let body = bounded_seller_error(response).await;
return Ok(
serde_json::json!({ "error": seller_error_message(status, &body), "payment_started": false }),
);
}
let body: serde_json::Value = response
.json()
@@ -1196,9 +1254,11 @@ impl RpcHandler {
}
};
if !response.status().is_success() {
return Ok(serde_json::json!({
"error": format!("Seller could not provide an address ({}).", response.status())
}));
let status = response.status();
let body = bounded_seller_error(response).await;
return Ok(
serde_json::json!({ "error": seller_error_message(status, &body), "payment_started": false }),
);
}
let body: serde_json::Value = response
.json()
@@ -332,6 +332,7 @@ impl RpcHandler {
"content.download-peer-paid" => self.handle_content_download_peer_paid(params).await,
"content.indeehub-projects" => self.handle_content_indeehub_projects().await,
"content.browse-all-peers" => self.handle_content_browse_all_peers().await,
"content.payment-status" => self.handle_content_payment_status(params).await,
"content.owned-list" => self.handle_content_owned_list().await,
"content.owned-get" => self.handle_content_owned_get(params).await,
"content.request-invoice" => self.handle_content_request_invoice(params).await,
@@ -512,6 +512,82 @@ mod lifecycle_regression_tests {
assert_eq!(main.lan_config.as_deref(), Some("kept"));
}
#[test]
fn installed_manifest_entry_path_survives_scans_without_changing_runtime_origin() {
let mut entry = installing_fixture();
entry.state = PackageState::Stopped;
entry.installed = Some(InstalledPackageDataEntry {
current_dependents: HashMap::new(),
current_dependencies: HashMap::new(),
last_backup: None,
interface_addresses: HashMap::from([(
"main".into(),
InterfaceAddress {
lan_address: Some("https://localhost:7443/old?gate=retained#position".into()),
tor_address: "existing.onion".into(),
},
)]),
status: ServiceStatus::Stopped,
});
let manifest = serde_json::json!({"app":{"interfaces":{"main":{
"type":"ui", "port":80, "protocol":"http", "path":"/browse"
}}}});
for _ in 0..2 {
apply_manifest_value(&manifest, &mut entry);
let main = &entry.installed.as_ref().unwrap().interface_addresses["main"];
assert_eq!(
main.lan_address.as_deref(),
Some("https://localhost:7443/browse?gate=retained#position")
);
assert_eq!(main.tor_address, "existing.onion");
assert_eq!(entry.state, PackageState::Stopped);
assert_eq!(entry.ui_ready, Some(false));
}
entry
.installed
.as_mut()
.unwrap()
.interface_addresses
.get_mut("main")
.unwrap()
.lan_address = None;
apply_manifest_value(&manifest, &mut entry);
assert!(
entry.installed.as_ref().unwrap().interface_addresses["main"]
.lan_address
.is_none()
);
}
#[test]
fn manifest_entry_path_preserves_authority_and_optional_query() {
assert_eq!(
manifest_launch_path("http://localhost:7475", "/browse"),
Some("http://localhost:7475/browse".into())
);
assert_eq!(
manifest_launch_path(
"https://localhost:7443/old?keep=1#old",
"/browse?view=all#new"
),
Some("https://localhost:7443/browse?view=all#new".into())
);
for path in [
"https://foreign.invalid/",
"//foreign.invalid/",
"/\\foreign.invalid/",
"browse",
"/bad\npath",
] {
assert!(
manifest_launch_path("https://localhost:7443", path).is_none(),
"{path:?}"
);
}
assert!(manifest_launch_path("", "/browse").is_none());
assert!(manifest_launch_path("file:///tmp/local", "/browse").is_none());
}
#[test]
fn btcpay_aliases_share_one_package_without_promoting_dependencies() {
for name in ["btcpay", "btcpayserver", "btcpay-server", "archy-btcpay"] {
@@ -761,6 +837,38 @@ fn apply_manifest_value(value: &serde_json::Value, entry: &mut PackageDataEntry)
entry.manifest.tier = Some(tier);
}
}
// Runtime discovery owns origins, ports and availability; the reviewed
// manifest owns the UI's entry path. Otherwise every scan resets a declared
// /browse or /admin entry point to the container root.
if let (Some(installed), Some(interfaces)) = (
entry.installed.as_mut(),
app.get("interfaces").and_then(|v| v.as_object()),
) {
for (id, address) in &mut installed.interface_addresses {
let Some(interface) = interfaces.get(id) else {
continue;
};
if interface
.get("type")
.and_then(|v| v.as_str())
.unwrap_or("ui")
!= "ui"
{
continue;
}
let Some(path) = interface.get("path").and_then(|v| v.as_str()) else {
continue;
};
if let Some(lan) = address.lan_address.as_mut() {
if let Some(updated) = manifest_launch_path(lan, path) {
*lan = updated;
}
}
if let Some(updated) = manifest_launch_path(&address.tor_address, path) {
address.tor_address = updated;
}
}
}
// Once installed, the scanner owns UI detection (including companion UIs).
// Only seed classification while there is no observed runtime package.
if entry.installed.is_some() {
@@ -790,6 +898,33 @@ fn apply_manifest_value(value: &serde_json::Value, entry: &mut PackageDataEntry)
}
}
/// Apply a local entry path without changing the scanner-confirmed authority.
fn manifest_launch_path(address: &str, path: &str) -> Option<String> {
if !path.starts_with('/')
|| path.starts_with("//")
|| path.contains('\\')
|| path.chars().any(char::is_control)
{
return None;
}
let base = reqwest::Url::parse(address).ok()?;
if !matches!(base.scheme(), "http" | "https") {
return None;
}
let mut updated = base.join(path).ok()?;
if updated.origin() != base.origin() {
return None;
}
// A plain path must not discard an existing gate/deep-link query.
if !path.contains('?') {
updated.set_query(base.query());
}
if !path.contains('#') {
updated.set_fragment(base.fragment());
}
Some(updated.into())
}
fn get_app_metadata(app_id: &str) -> AppMetadata {
let mut meta = match app_id {
"bitcoin-core" => AppMetadata {
+1
View File
@@ -18,6 +18,7 @@ pub mod npm;
pub mod prod_orchestrator;
pub mod quadlet;
pub mod registry;
pub mod registration_pin;
pub mod secrets;
pub mod traits;
pub mod ui_detection;
@@ -3915,6 +3915,27 @@ impl ProdContainerOrchestrator {
}
async fn resolve_dynamic_env(&self, manifest: &mut AppManifest) -> Result<()> {
if manifest.app.container.media_registration_identity {
// Only opted-in manifests read the existing appliance identity. A
// missing key/public-key mismatch is an error, never key generation.
let identity =
crate::identity::NodeIdentity::load_existing(&self.data_dir.join("identity"))
.await?;
anyhow::ensure!(
self.node_pubkey_hex().await? == identity.pubkey_hex(),
"Existing node public key disagrees with its signing identity"
);
let data_dir = self.data_dir.clone();
let app_id = manifest.app.id.clone();
let pin = tokio::task::spawn_blocking(move || {
crate::container::registration_pin::ensure_for_installation(
&data_dir, &app_id, &identity,
)
})
.await
.context("Application registration pin worker failed")??;
crate::container::registration_pin::apply_environment(manifest, &pin)?;
}
if manifest.app.id == "nginx-proxy-manager" {
crate::container::npm::resolve_storage()
.await?
@@ -6433,6 +6454,100 @@ app:
.unwrap();
}
#[tokio::test]
async fn opted_in_media_registration_pins_preserve_audience_through_env_reconciliation() {
let rt = Arc::new(MockRuntime::default());
let root = tempfile::tempdir().unwrap();
let mut orch =
ProdContainerOrchestrator::with_runtime(rt, PathBuf::from("/nonexistent-for-tests"));
orch.set_data_dir(root.path().to_owned());
orch.set_secrets_dir(root.path().join("secrets"));
let identity = crate::identity::NodeIdentity::load_or_create(&root.path().join("identity"))
.await
.unwrap();
let key_before = tokio::fs::read(root.path().join("identity/node_key"))
.await
.unwrap();
let mut app = pull_manifest("indeedhub-api", "fixture:1");
app.app.container.media_registration_identity = true;
app.app
.environment
.push("ARCHIPELAGO_REGISTRATION_AUDIENCE=untrusted-manifest-value".into());
orch.resolve_dynamic_env(&mut app).await.unwrap();
let first = app.app.environment.clone();
orch.resolve_dynamic_env(&mut app).await.unwrap();
assert_eq!(app.app.environment, first);
let pin = crate::container::registration_pin::load_existing(
root.path(),
"indeedhub-api",
&identity,
)
.unwrap();
assert!(first.contains(&format!(
"ARCHIPELAGO_REGISTRATION_NODE_PUBLIC_KEY={}",
identity.pubkey_hex()
)));
assert!(first.contains(&format!(
"ARCHIPELAGO_REGISTRATION_NODE_DID={}",
identity.did_key().unwrap()
)));
assert!(first.contains(&format!(
"ARCHIPELAGO_REGISTRATION_AUDIENCE={}",
pin.app_audience
)));
assert!(!first.iter().any(|entry| entry.contains("_ENABLED=")));
assert_eq!(
tokio::fs::read(root.path().join("identity/node_key"))
.await
.unwrap(),
key_before
);
}
#[tokio::test]
async fn unrelated_apps_do_not_require_identity_or_get_registration_pins() {
let rt = Arc::new(MockRuntime::default());
let root = tempfile::tempdir().unwrap();
let mut orch =
ProdContainerOrchestrator::with_runtime(rt, PathBuf::from("/nonexistent-for-tests"));
orch.set_data_dir(root.path().to_owned());
orch.set_secrets_dir(root.path().join("secrets"));
let mut app = pull_manifest("ordinary-app", "fixture:1");
orch.resolve_dynamic_env(&mut app).await.unwrap();
assert!(!root.path().join("identity").exists());
assert!(!root.path().join("app-registration-pins").exists());
assert!(!app
.app
.environment
.iter()
.any(|entry| entry.starts_with("ARCHIPELAGO_REGISTRATION_")));
}
#[tokio::test]
async fn registration_pins_require_existing_matching_node_keys() {
let rt = Arc::new(MockRuntime::default());
let root = tempfile::tempdir().unwrap();
let mut orch = ProdContainerOrchestrator::with_runtime(
rt.clone(),
PathBuf::from("/nonexistent-for-tests"),
);
orch.set_data_dir(root.path().to_owned());
orch.set_secrets_dir(root.path().join("secrets"));
let mut app = pull_manifest("indeedhub-api", "fixture:1");
app.app.container.media_registration_identity = true;
assert!(orch.resolve_dynamic_env(&mut app).await.is_err());
assert!(!root.path().join("identity").exists());
crate::identity::NodeIdentity::load_or_create(&root.path().join("identity"))
.await
.unwrap();
tokio::fs::write(root.path().join("identity/node_key.pub"), [3u8; 32])
.await
.unwrap();
assert!(orch.resolve_dynamic_env(&mut app).await.is_err());
assert!(!root.path().join("app-registration-pins").exists());
assert!(rt.calls().is_empty());
}
#[tokio::test]
async fn node_identity_pubkeys_placeholder_renders_the_signable_identities() {
// The owners must be exactly the identities the app signer offers:
@@ -0,0 +1,374 @@
//! Installer-owned public identity bindings for opted-in media-registration apps.
//! No identity generation, app enablement, signing permission or data deletion.
use anyhow::{Context, Result};
use archipelago_container::AppManifest;
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use std::fs::{File, OpenOptions};
use std::io::{Read, Write};
use std::os::fd::AsRawFd;
use std::os::unix::fs::{DirBuilderExt, MetadataExt, OpenOptionsExt};
use std::path::{Path, PathBuf};
use std::time::{Duration, Instant};
const ROOT: &str = "app-registration-pins";
const MAX_RECORD: u64 = 16 * 1024;
const ENV_NAMES: [&str; 3] = [
"ARCHIPELAGO_REGISTRATION_NODE_PUBLIC_KEY",
"ARCHIPELAGO_REGISTRATION_NODE_DID",
"ARCHIPELAGO_REGISTRATION_AUDIENCE",
];
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct RegistrationPin {
version: u8,
pub app_id: String,
pub node_public_key: String,
pub node_did: String,
pub app_audience: String,
}
#[derive(Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct Marker {
version: u8,
app_id: String,
pin_sha256: String,
}
fn valid_app_id(id: &str) -> bool {
!id.is_empty()
&& id.len() <= 128
&& !id.starts_with('-')
&& !id.ends_with('-')
&& !id.contains("--")
&& id
.bytes()
.all(|v| v.is_ascii_lowercase() || v.is_ascii_digit() || v == b'-')
}
fn paths(data_dir: &Path, app_id: &str) -> Result<(PathBuf, PathBuf, PathBuf)> {
anyhow::ensure!(
valid_app_id(app_id),
"Invalid application registration scope"
);
let root = data_dir.join(ROOT);
Ok((
root.clone(),
root.join(format!("{app_id}.json")),
root.join(format!("{app_id}.initialized.json")),
))
}
fn private_root(path: &Path) -> Result<File> {
let dir = OpenOptions::new()
.read(true)
.custom_flags(libc::O_DIRECTORY | libc::O_NOFOLLOW | libc::O_CLOEXEC)
.open(path)?;
let metadata = dir.metadata()?;
anyhow::ensure!(
metadata.uid() == unsafe { libc::geteuid() } && metadata.mode() & 0o077 == 0,
"Application registration pins must remain private to the node service"
);
Ok(dir)
}
fn lock(root: &File) -> Result<()> {
let deadline = Instant::now() + Duration::from_secs(30);
loop {
if unsafe { libc::flock(root.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) } == 0 {
return Ok(());
}
let error = std::io::Error::last_os_error();
if !matches!(
error.kind(),
std::io::ErrorKind::WouldBlock | std::io::ErrorKind::Interrupted
) {
return Err(error.into());
}
anyhow::ensure!(
Instant::now() < deadline,
"Application registration provisioning is busy; retry the same install"
);
std::thread::sleep(Duration::from_millis(20));
}
}
fn read<T: serde::de::DeserializeOwned>(path: &Path) -> Result<Option<T>> {
let file = match OpenOptions::new()
.read(true)
.custom_flags(libc::O_NOFOLLOW | libc::O_NONBLOCK | libc::O_CLOEXEC)
.open(path)
{
Ok(file) => file,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
Err(error) => {
return Err(error)
.context("Application registration pin is unavailable; preserve it for recovery")
}
};
let metadata = file.metadata()?;
anyhow::ensure!(
metadata.is_file() && metadata.len() <= MAX_RECORD && metadata.mode() & 0o077 == 0,
"Invalid application registration pin storage"
);
let mut bytes = Vec::new();
file.take(MAX_RECORD + 1).read_to_end(&mut bytes)?;
anyhow::ensure!(
bytes.len() as u64 <= MAX_RECORD,
"Application registration pin is too large"
);
Ok(Some(serde_json::from_slice(&bytes).context(
"Damaged application registration pin; do not replace it",
)?))
}
fn persist<T: Serialize>(root: &Path, path: &Path, value: &T) -> Result<()> {
let bytes = serde_json::to_vec(value)?;
let mut pending = tempfile::NamedTempFile::new_in(root)?;
pending.write_all(&bytes)?;
pending.as_file().sync_all()?;
pending.persist_noclobber(path).map_err(|error| {
anyhow::anyhow!(
"Could not commit application registration pin: {}",
error.error
)
})?;
private_root(root)?.sync_all()?;
Ok(())
}
fn validate(
pin: &RegistrationPin,
app_id: &str,
identity: &crate::identity::NodeIdentity,
) -> Result<()> {
let audience =
uuid::Uuid::parse_str(&pin.app_audience).context("Invalid saved application audience")?;
anyhow::ensure!(pin.version == 1 && pin.app_id == app_id
&& pin.node_public_key == identity.pubkey_hex() && pin.node_did == identity.did_key()?
&& audience.get_version_num() == 4 && audience.get_variant() == uuid::Variant::RFC4122
&& audience.to_string() == pin.app_audience,
"Application registration identity changed; preserve the existing pin and migrate explicitly");
Ok(())
}
fn commitment(pin: &RegistrationPin) -> Result<String> {
Ok(hex::encode(Sha256::digest(serde_json::to_vec(pin)?)))
}
fn validate_marker(marker: &Marker, pin: &RegistrationPin) -> Result<()> {
anyhow::ensure!(
marker.version == 1 && marker.app_id == pin.app_id && marker.pin_sha256 == commitment(pin)?,
"Application registration initialization record changed; preserve both records"
);
Ok(())
}
/// Installer-only provisioning. Run on a blocking worker. Existing data is never
/// replaced or repaired by generating another audience. Does not load/create keys.
pub fn ensure_for_installation(
data_dir: &Path,
app_id: &str,
identity: &crate::identity::NodeIdentity,
) -> Result<RegistrationPin> {
let (root, path, marker_path) = paths(data_dir, app_id)?;
let mut builder = std::fs::DirBuilder::new();
builder.mode(0o700);
if let Err(error) = builder.create(&root) {
if error.kind() != std::io::ErrorKind::AlreadyExists {
return Err(error.into());
}
}
let held = private_root(&root)?;
File::open(data_dir)?.sync_all()?;
lock(&held)?;
let marker: Option<Marker> = read(&marker_path)?;
let pin = match read::<RegistrationPin>(&path)? {
Some(pin) => {
validate(&pin, app_id, identity)?;
pin
}
None => {
anyhow::ensure!(marker.is_none(), "Previously provisioned application pin is missing; recover it instead of creating another audience");
let pin = RegistrationPin {
version: 1,
app_id: app_id.into(),
node_public_key: identity.pubkey_hex(),
node_did: identity.did_key()?,
app_audience: uuid::Uuid::new_v4().to_string(),
};
persist(&root, &path, &pin)?;
pin
}
};
if let Some(marker) = marker {
validate_marker(&marker, &pin)?;
} else {
persist(
&root,
&marker_path,
&Marker {
version: 1,
app_id: app_id.into(),
pin_sha256: commitment(&pin)?,
},
)?;
}
held.sync_all()?;
Ok(pin)
}
/// Caller-side lookup never creates installation state. Require this before
/// accepting a native registration request for the installed application scope.
pub fn load_existing(
data_dir: &Path,
app_id: &str,
identity: &crate::identity::NodeIdentity,
) -> Result<RegistrationPin> {
let (root, path, marker_path) = paths(data_dir, app_id)?;
let held = private_root(&root)?;
lock(&held)?;
let pin: RegistrationPin =
read(&path)?.context("Application registration identity is not provisioned")?;
validate(&pin, app_id, identity)?;
let marker: Marker =
read(&marker_path)?.context("Application registration provisioning is incomplete")?;
validate_marker(&marker, &pin)?;
Ok(pin)
}
/// Installer identity pins are authoritative. No flag enabling registration or
/// publication is set, and non-opted-in apps are not modified.
pub fn apply_environment(manifest: &mut AppManifest, pin: &RegistrationPin) -> Result<()> {
anyhow::ensure!(
manifest.app.container.media_registration_identity && manifest.app.id == pin.app_id,
"Registration pin does not belong to this opted-in application"
);
anyhow::ensure!(!manifest.app.container.derived_env.iter().any(|entry| ENV_NAMES.contains(&entry.key.as_str()))
&& !manifest.app.container.secret_env.iter().any(|entry| ENV_NAMES.contains(&entry.key.as_str()))
&& !manifest.app.container.secret_env_refs.iter().any(|entry| ENV_NAMES.contains(&entry.env_key.as_str())),
"Registration identity variables cannot be replaced by derived or secret environment entries");
manifest.app.environment.retain(|entry| {
!entry
.split_once('=')
.is_some_and(|(key, _)| ENV_NAMES.contains(&key))
});
manifest.app.environment.extend([
format!("{}={}", ENV_NAMES[0], pin.node_public_key),
format!("{}={}", ENV_NAMES[1], pin.node_did),
format!("{}={}", ENV_NAMES[2], pin.app_audience),
]);
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use std::os::unix::fs::PermissionsExt;
use std::sync::Arc;
async fn fixture() -> (tempfile::TempDir, crate::identity::NodeIdentity) {
let root = tempfile::tempdir().unwrap();
let identity = crate::identity::NodeIdentity::load_or_create(&root.path().join("identity"))
.await
.unwrap();
(root, identity)
}
#[tokio::test]
async fn audience_is_stable_across_retries_reconstruction_and_other_apps_are_distinct() {
let (root, identity) = fixture().await;
assert!(load_existing(root.path(), "indeedhub-api", &identity).is_err());
assert!(!root.path().join(ROOT).exists());
let original_key = std::fs::read(root.path().join("identity/node_key")).unwrap();
let first = ensure_for_installation(root.path(), "indeedhub-api", &identity).unwrap();
let reloaded = crate::identity::NodeIdentity::load_existing(&root.path().join("identity"))
.await
.unwrap();
assert_eq!(
ensure_for_installation(root.path(), "indeedhub-api", &reloaded).unwrap(),
first
);
assert_eq!(
load_existing(root.path(), "indeedhub-api", &reloaded).unwrap(),
first
);
assert_ne!(
ensure_for_installation(root.path(), "another-app", &reloaded)
.unwrap()
.app_audience,
first.app_audience
);
assert_eq!(
std::fs::read(root.path().join("identity/node_key")).unwrap(),
original_key
);
}
#[tokio::test]
async fn concurrent_installers_persist_one_audience() {
let (root, identity) = fixture().await;
let identity = Arc::new(identity);
let workers: Vec<_> = (0..4)
.map(|_| {
let identity = identity.clone();
let path = root.path().to_owned();
std::thread::spawn(move || {
ensure_for_installation(&path, "indeedhub-api", &identity).unwrap()
})
})
.collect();
let results: Vec<_> = workers
.into_iter()
.map(|worker| worker.join().unwrap())
.collect();
assert!(results.iter().all(|pin| pin == &results[0]));
}
#[tokio::test]
async fn corrupt_missing_or_foreign_pin_is_preserved_without_rotation() {
let (root, identity) = fixture().await;
let first = ensure_for_installation(root.path(), "indeedhub-api", &identity).unwrap();
let (_, path, marker) = paths(root.path(), "indeedhub-api").unwrap();
let bytes = std::fs::read(&path).unwrap();
let marker_bytes = std::fs::read(&marker).unwrap();
let (_other_root, other_identity) = fixture().await;
assert!(ensure_for_installation(root.path(), "indeedhub-api", &other_identity).is_err());
assert_eq!(std::fs::read(&path).unwrap(), bytes);
std::fs::write(&path, b"damaged").unwrap();
assert!(ensure_for_installation(root.path(), "indeedhub-api", &identity).is_err());
assert_eq!(std::fs::read(&path).unwrap(), b"damaged");
std::fs::remove_file(&path).unwrap();
assert!(ensure_for_installation(root.path(), "indeedhub-api", &identity).is_err());
assert!(!path.exists());
assert_eq!(std::fs::read(&marker).unwrap(), marker_bytes);
std::fs::write(&path, &bytes).unwrap();
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o600)).unwrap();
assert_eq!(
load_existing(root.path(), "indeedhub-api", &identity).unwrap(),
first
);
}
#[tokio::test]
async fn interrupted_marker_commit_finishes_with_the_same_pin() {
let (root, identity) = fixture().await;
let first = ensure_for_installation(root.path(), "indeedhub-api", &identity).unwrap();
let (_, path, marker) = paths(root.path(), "indeedhub-api").unwrap();
let bytes = std::fs::read(&path).unwrap();
std::fs::remove_file(marker).unwrap();
assert!(load_existing(root.path(), "indeedhub-api", &identity).is_err());
assert_eq!(
ensure_for_installation(root.path(), "indeedhub-api", &identity).unwrap(),
first
);
assert_eq!(std::fs::read(path).unwrap(), bytes);
}
#[tokio::test]
async fn no_pin_symlink_or_manifest_override_is_followed() {
let (root, identity) = fixture().await;
let pin = ensure_for_installation(root.path(), "indeedhub-api", &identity).unwrap();
let (_, path, _) = paths(root.path(), "indeedhub-api").unwrap();
std::fs::remove_file(&path).unwrap();
std::os::unix::fs::symlink(root.path().join("identity/node_key"), &path).unwrap();
assert!(ensure_for_installation(root.path(), "indeedhub-api", &identity).is_err());
let mut manifest = AppManifest::parse("app:\n id: indeedhub-api\n name: Fixture\n version: '1'\n container:\n image: fixture:1\n media_registration_identity: true\n environment:\n - ARCHIPELAGO_REGISTRATION_AUDIENCE=body-value\n").unwrap();
apply_environment(&mut manifest, &pin).unwrap();
let first = manifest.app.environment.clone();
apply_environment(&mut manifest, &pin).unwrap();
assert_eq!(manifest.app.environment, first);
assert!(first.contains(&format!(
"ARCHIPELAGO_REGISTRATION_AUDIENCE={}",
pin.app_audience
)));
assert!(!first.iter().any(|entry| entry.contains("_ENABLED=")));
manifest.app.container.media_registration_identity = false;
assert!(apply_environment(&mut manifest, &pin).is_err());
}
}
+289 -8
View File
@@ -130,6 +130,27 @@ impl Contract {
}
}
/// Accepted seller liability has no automatic expiry or garbage collection.
/// The seller must retain immutable snapshot eligibility until settlement or an
/// explicit future cancellation/refund protocol resolves this commitment.
#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct Acceptance {
pub contract_hash: String,
pub accepted_at: i64,
}
impl Acceptance {
fn validate(&self, contract: &Contract) -> Result<()> {
anyhow::ensure!(
self.contract_hash == contract.context_hash()?
&& self.accepted_at >= contract.offered_at
&& self.accepted_at < contract.expires_at,
"Seller acceptance does not match the offer"
);
Ok(())
}
}
/// Private delivery capability; do not log or expose it to another buyer.
/// The wire layer must authenticate this receipt before a buyer stores it.
#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
@@ -206,6 +227,7 @@ impl PreparedToken {
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) enum BuyerPhase {
Intent,
AcceptanceSaved,
TokenPrepared,
ReceiptSaved,
Delivered,
@@ -215,10 +237,34 @@ pub(crate) enum BuyerPhase {
pub(crate) struct BuyerRecord {
pub contract: Contract,
pub phase: BuyerPhase,
acceptance: Option<Acceptance>,
token: Option<PreparedToken>,
receipt: Option<Receipt>,
}
impl BuyerRecord {
pub fn public_status(&self) -> serde_json::Value {
let (state, settlement_confirmed, delivered) = match self.phase {
BuyerPhase::Intent => ("intent", false, false),
BuyerPhase::AcceptanceSaved => ("accepted_payment_unconfirmed", false, false),
BuyerPhase::TokenPrepared => ("token_prepared_settlement_unconfirmed", false, false),
BuyerPhase::ReceiptSaved => ("settled_delivery_pending", true, false),
BuyerPhase::Delivered => ("delivered", true, true),
};
serde_json::json!({
"operation_id": self.contract.id,
"content_id": self.contract.content_id,
"seller_did": self.contract.seller_did,
"state": state,
"gross_sats": self.contract.gross_token_sats,
"minimum_net_sats": self.contract.minimum_net_sats,
"settlement_confirmed": settlement_confirmed,
"amount_received": self.receipt().map(|receipt| receipt.amount_received),
"delivered": delivered,
"recovery_required": !delivered,
"can_start_new_payment": false,
})
}
pub fn token(&self) -> Option<&str> {
self.token.as_ref().map(|token| token.encoded.as_str())
}
@@ -227,6 +273,9 @@ impl BuyerRecord {
}
fn validate(&self) -> Result<()> {
self.contract.validate()?;
if let Some(acceptance) = &self.acceptance {
acceptance.validate(&self.contract)?;
}
if let Some(token) = &self.token {
token.validate(&self.contract)?;
}
@@ -235,10 +284,14 @@ impl BuyerRecord {
}
anyhow::ensure!(
match self.phase {
BuyerPhase::Intent => self.token.is_none() && self.receipt.is_none(),
BuyerPhase::TokenPrepared => self.token.is_some() && self.receipt.is_none(),
BuyerPhase::Intent =>
self.acceptance.is_none() && self.token.is_none() && self.receipt.is_none(),
BuyerPhase::AcceptanceSaved =>
self.acceptance.is_some() && self.token.is_none() && self.receipt.is_none(),
BuyerPhase::TokenPrepared =>
self.acceptance.is_some() && self.token.is_some() && self.receipt.is_none(),
BuyerPhase::ReceiptSaved | BuyerPhase::Delivered =>
self.token.is_some() && self.receipt.is_some(),
self.acceptance.is_some() && self.token.is_some() && self.receipt.is_some(),
},
"Invalid buyer purchase transition"
);
@@ -256,11 +309,27 @@ pub(crate) enum SellerPhase {
#[serde(deny_unknown_fields)]
pub(crate) struct SellerRecord {
pub contract: Contract,
pub accepted_at: i64,
token_hash: Option<String>,
pub phase: SellerPhase,
}
impl SellerRecord {
pub fn acceptance(&self) -> Result<Acceptance> {
Ok(Acceptance {
contract_hash: self.contract.context_hash()?,
accepted_at: self.accepted_at,
})
}
fn validate(&self) -> Result<()> {
self.contract.validate()?;
self.acceptance()?.validate(&self.contract)?;
if let Some(token_hash) = &self.token_hash {
anyhow::ensure!(valid_hash(token_hash), "Invalid seller token hash");
}
anyhow::ensure!(
matches!(self.phase, SellerPhase::Intent) || self.token_hash.is_some(),
"Seller token was not durably bound"
);
match &self.phase {
SellerPhase::Intent => (),
SellerPhase::Settled { amount_received } => {
@@ -458,16 +527,73 @@ impl Journal {
}
Ok(result)
}
/// Discover existing node-owned intent after browser storage loss. The
/// journal lock makes this lookup and prepare_buyer's duplicate guard one
/// serialized decision; caller-supplied fresh UUIDs cannot bypass it.
pub async fn find_buyers(
&self,
buyer_did: &str,
seller_did: &str,
content_id: &str,
) -> Result<Vec<BuyerRecord>> {
crate::identity::pubkey_bytes_from_did_key(buyer_did)?;
crate::identity::pubkey_bytes_from_did_key(seller_did)?;
let mut entries = fs::read_dir(&self.directory).await?;
let mut result = Vec::new();
let mut count = 0usize;
while let Some(entry) = entries.next_entry().await? {
let name = entry.file_name();
let Some(name) = name.to_str() else { continue };
let Some(id) = name
.strip_prefix("buyer-")
.and_then(|v| v.strip_suffix(".json"))
else {
continue;
};
count += 1;
anyhow::ensure!(
count <= 10000,
"Purchase recovery index requires maintenance; do not pay again"
);
let record = self
.buyer(id)
.await?
.context("Purchase recovery disappeared")?;
if record.contract.buyer_did == buyer_did
&& record.contract.seller_did == seller_did
&& record.contract.content_id == content_id
{
result.push(record);
}
}
result.sort_by(|a, b| a.contract.id.cmp(&b.contract.id));
Ok(result)
}
pub async fn prepare_buyer(&self, contract: &Contract, now: i64) -> Result<BuyerRecord> {
contract.validate()?;
if let Some(record) = self.buyer(&contract.id).await? {
anyhow::ensure!(&record.contract == contract, "Buyer purchase terms changed");
return Ok(record);
}
let pending = self
.find_buyers(
&contract.buyer_did,
&contract.seller_did,
&contract.content_id,
)
.await?;
anyhow::ensure!(
pending
.iter()
.all(|record| record.phase == BuyerPhase::Delivered),
"An existing purchase must be recovered before a new operation is created"
);
contract.validate_new_at(now)?;
let record = BuyerRecord {
contract: contract.clone(),
phase: BuyerPhase::Intent,
acceptance: None,
token: None,
receipt: None,
};
@@ -486,6 +612,8 @@ impl Journal {
contract.validate_new_at(now)?;
let record = SellerRecord {
contract: contract.clone(),
accepted_at: now,
token_hash: None,
phase: SellerPhase::Intent,
};
self.write("seller", &contract.id, &record).await?;
@@ -510,6 +638,56 @@ impl Journal {
);
Ok(record)
}
/// The transport layer must verify response provenance before passing the
/// verified seller DID. A claimed DID or client mint timestamp is insufficient.
pub async fn record_acceptance(
&self,
contract: &Contract,
acceptance: &Acceptance,
verified_seller_did: &str,
) -> Result<BuyerRecord> {
anyhow::ensure!(
verified_seller_did == contract.seller_did,
"Acceptance is from another seller"
);
let mut record = self.bound_buyer(contract).await?;
acceptance.validate(contract)?;
if let Some(previous) = &record.acceptance {
anyhow::ensure!(previous == acceptance, "Seller acceptance changed");
return Ok(record);
}
anyhow::ensure!(
record.phase == BuyerPhase::Intent,
"Buyer acceptance phase changed"
);
record.acceptance = Some(acceptance.clone());
record.phase = BuyerPhase::AcceptanceSaved;
record.validate()?;
self.write("buyer", &contract.id, &record).await?;
Ok(record)
}
/// Bind exact incoming bytes before wallet settlement. Receipt/status replay
/// remains possible without bearer token bytes, but settlement cannot change them.
pub async fn record_incoming_token(
&self,
contract: &Contract,
encoded: &str,
) -> Result<SellerRecord> {
let mut record = self.bound_seller(contract).await?;
let token = PreparedToken::new(contract, encoded.into())?;
if let Some(previous) = &record.token_hash {
anyhow::ensure!(previous == &token.sha256, "Seller incoming token changed");
return Ok(record);
}
anyhow::ensure!(
matches!(record.phase, SellerPhase::Intent),
"Seller settlement phase changed"
);
record.token_hash = Some(token.sha256);
record.validate()?;
self.write("seller", &contract.id, &record).await?;
Ok(record)
}
/// Call only with the original correlated recoverable-send result.
pub async fn record_token(&self, contract: &Contract, encoded: &str) -> Result<BuyerRecord> {
let mut record = self.bound_buyer(contract).await?;
@@ -522,8 +700,8 @@ impl Journal {
return Ok(record);
}
anyhow::ensure!(
record.phase == BuyerPhase::Intent,
"Cannot prepare token in this phase"
record.phase == BuyerPhase::AcceptanceSaved,
"Authenticated seller acceptance is not durable"
);
record.token = Some(token);
record.phase = BuyerPhase::TokenPrepared;
@@ -710,7 +888,19 @@ mod tests {
assert!(journal.record_token(&contract, &encoded).await.is_err());
assert!(journal.record_settlement(&contract, 7).await.is_err());
journal.prepare_buyer(&contract, 1500).await.unwrap();
journal.prepare_seller(&contract, 1500).await.unwrap();
let seller = journal.prepare_seller(&contract, 1500).await.unwrap();
journal
.record_acceptance(
&contract,
&&seller.acceptance().unwrap(),
&contract.seller_did,
)
.await
.unwrap();
journal
.record_incoming_token(&contract, &token(&contract, "seller-original"))
.await
.unwrap();
assert!(journal.issue_receipt(&contract).await.is_err());
journal.record_token(&contract, &encoded).await.unwrap();
journal.record_settlement(&contract, 7).await.unwrap();
@@ -769,7 +959,19 @@ mod tests {
let contract = contract();
let journal = Journal::open(root.path()).await.unwrap();
journal.prepare_buyer(&contract, 1500).await.unwrap();
journal.prepare_seller(&contract, 1500).await.unwrap();
let seller = journal.prepare_seller(&contract, 1500).await.unwrap();
journal
.record_acceptance(
&contract,
&&seller.acceptance().unwrap(),
&contract.seller_did,
)
.await
.unwrap();
journal
.record_incoming_token(&contract, &token(&contract, "seller-original"))
.await
.unwrap();
journal
.record_token(&contract, &token(&contract, "first"))
.await
@@ -870,6 +1072,15 @@ mod tests {
let contract = contract();
let mut journal = Journal::open(root.path()).await.unwrap();
journal.prepare_buyer(&contract, 1500).await.unwrap();
let seller = journal.prepare_seller(&contract, 1500).await.unwrap();
journal
.record_acceptance(
&contract,
&&seller.acceptance().unwrap(),
&contract.seller_did,
)
.await
.unwrap();
let reached = std::sync::Arc::new(tokio::sync::Notify::new());
let resume = std::sync::Arc::new(tokio::sync::Notify::new());
journal.before_commit = Some((reached.clone(), resume.clone()));
@@ -886,7 +1097,7 @@ mod tests {
let journal = Journal::open(root.path()).await.unwrap();
assert_eq!(
journal.buyer(&contract.id).await.unwrap().unwrap().phase,
BuyerPhase::Intent
BuyerPhase::AcceptanceSaved
);
let current_token = token(&contract, "current");
journal
@@ -906,4 +1117,74 @@ mod tests {
assert!(!entry.file_name().to_string_lossy().ends_with(".tmp"));
}
}
#[tokio::test]
async fn node_lookup_retains_pending_purchase_when_client_generates_a_fresh_id() {
let root = tempfile::tempdir().unwrap();
let first = contract();
let journal = Journal::open(root.path()).await.unwrap();
journal.prepare_buyer(&first, 1100).await.unwrap();
let mut duplicate = first.clone();
duplicate.id = uuid::Uuid::new_v4().to_string();
assert!(journal.prepare_buyer(&duplicate, 1100).await.is_err());
assert!(journal.buyer(&duplicate.id).await.unwrap().is_none());
let records = journal
.find_buyers(&first.buyer_did, &first.seller_did, &first.content_id)
.await
.unwrap();
assert_eq!(records.len(), 1);
assert_eq!(records[0].contract, first);
duplicate.content_id = "other-file".into();
journal.prepare_buyer(&duplicate, 1100).await.unwrap();
let records = journal
.find_buyers(&first.buyer_did, &first.seller_did, &first.content_id)
.await
.unwrap();
assert_eq!(records.len(), 1);
let mut bytes = fs::read(journal.path("buyer", &first.id).unwrap())
.await
.unwrap();
bytes[0] = b'!';
fs::write(journal.path("buyer", &first.id).unwrap(), bytes)
.await
.unwrap();
assert!(journal
.find_buyers(&first.buyer_did, &first.seller_did, &first.content_id)
.await
.is_err());
}
#[tokio::test]
async fn public_purchase_status_never_exposes_token_or_delivery_capability() {
let root = tempfile::tempdir().unwrap();
let contract = contract();
let journal = Journal::open(root.path()).await.unwrap();
let buyer = journal.prepare_buyer(&contract, 1100).await.unwrap();
assert_eq!(buyer.public_status()["state"], "intent");
let seller = journal.prepare_seller(&contract, 1100).await.unwrap();
journal
.record_acceptance(
&contract,
&seller.acceptance().unwrap(),
&contract.seller_did,
)
.await
.unwrap();
let encoded = token(&contract, "private-status-proof");
let buyer = journal.record_token(&contract, &encoded).await.unwrap();
let status = buyer.public_status();
assert_eq!(status["state"], "token_prepared_settlement_unconfirmed");
assert_eq!(status["settlement_confirmed"], false);
assert!(!status.to_string().contains(&encoded));
journal
.record_incoming_token(&contract, &encoded)
.await
.unwrap();
journal.record_settlement(&contract, 8).await.unwrap();
let receipt = journal.issue_receipt(&contract).await.unwrap();
let buyer = journal.record_receipt(&contract, &receipt).await.unwrap();
let status = buyer.public_status();
assert_eq!(status["state"], "settled_delivery_pending");
assert_eq!(status["amount_received"], 8);
assert_eq!(status["can_start_new_payment"], false);
assert!(!status.to_string().contains(&receipt.capability));
}
}
@@ -0,0 +1,89 @@
//! Purchase/wallet sequencing only. The caller handles authenticated transport,
//! node-side pending-operation lookup and retained immutable content snapshots.
//! No legacy purchase RPC should call this until those boundaries are connected.
use crate::{
content_purchase::{BuyerPhase, Contract, Journal, Receipt, SellerPhase},
wallet::ecash,
};
use anyhow::{Context, Result};
use std::path::Path;
/// Buyer Intent and authenticated AcceptanceSaved must already be durable.
/// This function never allocates a purchase UUID or accepts a client mint time.
pub(crate) async fn prepare_buyer_token(data_dir: &Path, contract: &Contract) -> Result<String> {
contract.validate()?;
{
let journal = Journal::open(data_dir).await?;
let record = journal
.buyer(&contract.id)
.await?
.context("Buyer intent is not durable")?;
anyhow::ensure!(&record.contract == contract, "Buyer purchase terms changed");
if let Some(token) = record.token() {
return Ok(token.to_owned());
}
anyhow::ensure!(
record.phase == BuyerPhase::AcceptanceSaved,
"Authenticated seller acceptance is not durable"
);
}
// Deadline is enforced inside the wallet after recovering original results,
// immediately before a fresh exact send or mint POST. Do not pre-reject an
// expired contract here: a previous wallet commit may need to be recovered.
let token = ecash::send_token_recoverable_before(
data_dir,
&contract.id,
contract.network,
&contract.mint_url,
contract.gross_token_sats,
&contract.context_hash()?,
contract.expires_at,
)
.await?;
let journal = Journal::open(data_dir).await?;
let record = journal.record_token(contract, &token).await?;
Ok(record
.token()
.context("Prepared buyer token is missing")?
.to_owned())
}
/// Caller verified the buyer's signed exact-body request and bound it to this
/// contract. Seller acceptance must predate offer expiry; an existing accepted
/// liability honors delayed first settlement, not merely an already minted claim.
pub(crate) async fn settle_seller_token(
data_dir: &Path,
contract: &Contract,
encoded: &str,
verified_buyer_did: &str,
) -> Result<Receipt> {
anyhow::ensure!(
verified_buyer_did == contract.buyer_did,
"Payment is from another buyer"
);
contract.validate()?;
{
let journal = Journal::open(data_dir).await?;
// Missing acceptance rejects, regardless of current time; this endpoint
// cannot manufacture seller acceptance from an arriving payment token.
let record = journal.record_incoming_token(contract, encoded).await?;
match record.phase {
SellerPhase::ReceiptSaved(receipt) => return Ok(receipt),
SellerPhase::Settled { .. } => return journal.issue_receipt(contract).await,
SellerPhase::Intent => (),
}
}
let received = ecash::receive_token_recoverable(
data_dir,
&contract.id,
contract.network,
&contract.mint_url,
encoded,
contract.minimum_net_sats,
&contract.context_hash()?,
)
.await?;
let journal = Journal::open(data_dir).await?;
journal.record_settlement(contract, received).await?;
journal.issue_receipt(contract).await
}
+122 -42
View File
@@ -115,23 +115,6 @@ async fn save_catalog_unlocked(data_dir: &Path, catalog: &ContentCatalog) -> Res
Ok(())
}
/// Removes `id` from the on-disk catalog. Best-effort: a failure here just
/// means the entry gets pruned again next time it's requested, so errors are
/// logged rather than propagated.
async fn prune_missing_content_entry(data_dir: &Path, id: &str) {
let _lock = CATALOG_WRITES.lock().await;
let Ok(mut catalog) = load_catalog(data_dir).await else {
return;
};
let before = catalog.items.len();
catalog.items.retain(|i| i.id != id);
if catalog.items.len() != before {
if let Err(e) = save_catalog_unlocked(data_dir, &catalog).await {
warn!(error = %e, content_id = %id, "failed to save catalog after pruning missing content entry");
}
}
}
/// Get the full filesystem path for a content item.
/// Checks the dedicated content/files/ directory first, then falls back to the
/// FileBrowser data directory (where users manage files via the web UI).
@@ -155,6 +138,46 @@ pub fn content_file_path(data_dir: &Path, item: &ContentItem) -> PathBuf {
primary
}
pub(crate) fn validate_onchain_payment_price(price_sats: u64) -> Result<()> {
anyhow::ensure!(
price_sats >= 546,
"On-chain payment requires at least 546 sats. Choose Lightning or ecash for this file."
);
Ok(())
}
/// Read-only preflight before issuing a payable invoice/address. This catches
/// missing/replaced files but does not substitute for an immutable purchase
/// snapshot: later delivery must still preserve the original accepted contract.
pub(crate) async fn ensure_payment_source_available(
data_dir: &Path,
item: &ContentItem,
) -> Result<()> {
let path = content_file_path(data_dir, item);
let canonical = fs::canonicalize(&path)
.await
.context("The shared file is currently unavailable; no payment request was created")?;
let mut inside_root = false;
for root in [data_dir.join(CONTENT_DIR), data_dir.join("filebrowser")] {
if let Ok(root) = fs::canonicalize(root).await {
inside_root |= canonical.starts_with(root);
}
}
anyhow::ensure!(inside_root, "The shared file is outside the content roots");
let mut options = fs::OpenOptions::new();
options.read(true);
#[cfg(unix)]
options.custom_flags(libc::O_NOFOLLOW | libc::O_NONBLOCK);
let file = options
.open(&canonical)
.await
.context("The shared file cannot be opened; no payment request was created")?;
let metadata = file.metadata().await?;
anyhow::ensure!(metadata.is_file() && metadata.len() > 0 && metadata.len() == item.size_bytes,
"The shared file has changed or is unavailable; refresh its catalog before accepting payment");
Ok(())
}
/// Add a content item to the catalog.
///
/// Idempotent per FILE, not just per id: `content.add` mints a fresh UUID on
@@ -404,20 +427,10 @@ where
let file_path = content_file_path(data_dir, item);
if !file_path.exists() {
// The catalog entry survived (it's a separate JSON file) but its
// backing file is gone — most likely lost in an unrelated data-dir
// reset (a shared filebrowser file, 2026-07-01: two catalog entries
// outlived a filebrowser reinstall that wiped the files themselves).
// Leaving the entry in place would keep advertising it as available
// to every peer forever, each hitting the exact same dead end this
// one just did. Prune it so it stops being offered.
warn!(
content_id = %id,
filename = %item.filename,
"content catalog entry's file is missing on disk — pruning the stale entry"
);
prune_missing_content_entry(data_dir, id).await;
return Ok(ServeResult::NotFound);
// A disconnected mount, moved file or permission failure is not an
// instruction to unshare content or erase its purchase metadata.
warn!(content_id = %id, "Shared content is temporarily unavailable; catalog retained");
return Ok(ServeResult::Unavailable);
}
// Refuse unauthorized viewers before opening or reading any bytes.
@@ -981,11 +994,11 @@ mod faststart_tests {
}
#[cfg(test)]
mod prune_missing_content_tests {
mod unavailable_content_tests {
use super::*;
#[tokio::test]
async fn serve_content_prunes_catalog_entry_whose_file_is_missing() {
async fn unavailable_file_retains_identity_and_recovers_when_storage_returns() {
// Simulates a catalog entry that outlived its backing file (a shared
// filebrowser file lost in an unrelated data-dir reset, 2026-07-01) —
// every peer request for it would otherwise 404 forever with no way
@@ -1010,17 +1023,24 @@ mod prune_missing_content_tests {
let result = serve_content(data_dir, "missing-item", None, None, None, None, false)
.await
.unwrap();
assert!(matches!(result, ServeResult::NotFound));
assert!(matches!(result, ServeResult::Unavailable));
let reloaded = load_catalog(data_dir).await.unwrap();
assert!(
reloaded.items.is_empty(),
"stale entry should have been pruned after the 404"
);
assert_eq!(reloaded.items.len(), 1);
assert_eq!(reloaded.items[0].id, "missing-item");
fs::create_dir_all(data_dir.join("filebrowser"))
.await
.unwrap();
fs::write(data_dir.join("filebrowser/gone.mp4"), b"recovered")
.await
.unwrap();
let result = serve_content(data_dir, "missing-item", None, None, None, None, false)
.await
.unwrap();
assert!(matches!(result, ServeResult::Ok(bytes, _) if bytes == b"recovered"));
}
#[tokio::test]
async fn serve_content_leaves_other_entries_untouched_when_pruning() {
async fn unavailable_file_does_not_rewrite_any_catalog_entries() {
let dir = tempfile::tempdir().unwrap();
let data_dir = dir.path();
let missing = ContentItem {
@@ -1062,8 +1082,9 @@ mod prune_missing_content_tests {
.unwrap();
let reloaded = load_catalog(data_dir).await.unwrap();
assert_eq!(reloaded.items.len(), 1);
assert_eq!(reloaded.items[0].id, "present-item");
assert_eq!(reloaded.items.len(), 2);
assert_eq!(reloaded.items[0].id, "missing-item");
assert_eq!(reloaded.items[1].id, "present-item");
}
}
@@ -1817,3 +1838,62 @@ mod preview_boundary_tests {
);
}
}
#[cfg(test)]
mod payment_source_tests {
use super::*;
#[test]
fn onchain_quote_preserves_wallet_minimum_boundary() {
assert!(validate_onchain_payment_price(545).is_err());
assert!(validate_onchain_payment_price(546).is_ok());
}
#[tokio::test]
async fn payment_preflight_rejects_missing_changed_directory_and_escaped_sources() {
let root = tempfile::tempdir().unwrap();
let item = ContentItem {
id: "paid".into(),
filename: "Photos/clip.mp4".into(),
mime_type: "video/mp4".into(),
size_bytes: 4,
description: String::new(),
access: AccessControl::Paid {
price_sats: 2,
accepted: vec![],
},
availability: Availability::AllPeers,
added_at: String::new(),
};
assert!(ensure_payment_source_available(root.path(), &item)
.await
.is_err());
fs::create_dir_all(root.path().join("filebrowser/Photos"))
.await
.unwrap();
let path = root.path().join("filebrowser/Photos/clip.mp4");
fs::write(&path, b"film").await.unwrap();
ensure_payment_source_available(root.path(), &item)
.await
.unwrap();
fs::write(&path, b"changed").await.unwrap();
assert!(ensure_payment_source_available(root.path(), &item)
.await
.is_err());
fs::remove_file(&path).await.unwrap();
fs::create_dir(&path).await.unwrap();
assert!(ensure_payment_source_available(root.path(), &item)
.await
.is_err());
fs::remove_dir(&path).await.unwrap();
let outside = tempfile::tempdir().unwrap();
let outside_file = outside.path().join("film");
fs::write(&outside_file, b"film").await.unwrap();
#[cfg(unix)]
{
std::os::unix::fs::symlink(&outside_file, &path).unwrap();
assert!(ensure_payment_source_available(root.path(), &item)
.await
.is_err());
}
}
}
+481
View File
@@ -0,0 +1,481 @@
//! Retained versioned snapshots for explicitly shared Cloud files. Reuses the
//! confined descriptor and durable-copy primitives of media registration.
//! Call from spawn_blocking; ownership/visibility is authenticated by the caller.
use crate::media_registration::{self as io, Limits, SourceStamp};
use anyhow::{Context, Result};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use std::{
fs::File,
os::unix::{
fs::{MetadataExt, PermissionsExt},
io::AsRawFd,
},
path::Path,
time::{Duration, Instant},
};
#[derive(Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct Record {
version: u8,
content_id: String,
source: SourceStamp,
snapshot: SourceStamp,
sha256: String,
size: u64,
}
#[derive(Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct Envelope {
record: Record,
checksum: String,
}
fn read_manifest(version: &File) -> Result<Option<Record>> {
let Some(saved) = io::read_record::<Envelope>(version, "snapshot.json")? else {
return Ok(None);
};
anyhow::ensure!(
saved.checksum == hash(&serde_json::to_vec(&saved.record)?),
"Snapshot metadata is damaged; preserve accepted purchases"
);
Ok(Some(saved.record))
}
pub(crate) struct Snapshot {
pub file: File,
pub key: String,
pub sha256: String,
pub size: u64,
}
fn hash(bytes: &[u8]) -> String {
hex::encode(Sha256::digest(bytes))
}
/// Only callers holding a catalog item selected from this node's owner-managed
/// catalog may call this. Paths are confined to the supplied configured root.
/// max_total_bytes counts incomplete files too; failed quotes cannot fill disk.
pub(crate) fn prepare(
data_dir: &Path,
cloud_root: &Path,
content_id: &str,
relative: &Path,
limits: &Limits<'_>,
max_total_bytes: u64,
minimum_free_bytes: u64,
mut progress: impl FnMut(u64) -> Result<()>,
) -> Result<Snapshot> {
anyhow::ensure!(
!content_id.is_empty()
&& content_id.len() <= 256
&& content_id
.bytes()
.all(|c| c.is_ascii_alphanumeric() || b"_-".contains(&c)),
"Invalid shared content identity"
);
let data = io::open_directory(&data_dir.canonicalize()?)?;
let root = io::private_directory(&data, "content-snapshots")?;
// One copy admission decision at a time also makes quota accounting exact.
io::lock_operation(&root, limits, Instant::now() + Duration::from_secs(30))?;
let cloud = io::open_directory(&cloud_root.canonicalize()?)?;
let mut source = io::open_cloud_file(&cloud, relative)?;
let source_stamp = SourceStamp::read(&source)?;
let size = source.metadata()?.len();
anyhow::ensure!(
size > 0 && size <= limits.max_bytes,
"Shared file exceeds snapshot limits"
);
let key = hash(&serde_json::to_vec(&(
"shared-content-snapshot-v1",
content_id,
relative,
&source_stamp,
))?);
let mut versions = 0usize;
let mut existing_version = false;
for entry in std::fs::read_dir(format!("/proc/self/fd/{}", root.as_raw_fd()))? {
let entry = entry?;
versions += 1;
existing_version |= entry.file_name().to_str() == Some(&key);
}
anyhow::ensure!(
existing_version || versions < 4096,
"Shared snapshot version budget is full; retain existing purchases"
);
let version = io::private_directory(&root, &key)?;
if let Some(record) = read_manifest(&version)? {
anyhow::ensure!(
record.version == 1 && record.content_id == content_id && record.source == source_stamp,
"Shared snapshot identity changed; preserve it for recovery"
);
let file = io::open_at(&version, "media", libc::O_RDONLY | libc::O_NONBLOCK, 0)?;
anyhow::ensure!(
SourceStamp::read(&file)? == record.snapshot
&& file.metadata()?.mode() & 0o7777 == 0o400
&& file.metadata()?.len() == record.size,
"Retained shared snapshot changed; preserve accepted purchases"
);
return Ok(Snapshot {
file,
key,
sha256: record.sha256,
size: record.size,
});
}
// A crash can publish immutable bytes before the small manifest commit.
// Recover those exact bytes by checking both held descriptors; never replace
// an existing media file or manufacture a hash from a path/size alone.
match io::open_at(&version, "media", libc::O_RDONLY | libc::O_NONBLOCK, 0) {
Ok(mut file) => {
let stamp = SourceStamp::read(&file)?;
anyhow::ensure!(
file.metadata()?.mode() & 0o7777 == 0o400 && file.metadata()?.len() == size,
"Incomplete retained snapshot; preserve it for recovery"
);
let (snapshot_hash, snapshot_size) =
io::hash_file(&mut file, None, limits, &mut progress)?;
let (source_hash, source_size) =
io::hash_file(&mut source, None, limits, &mut progress)?;
anyhow::ensure!(
snapshot_size == size
&& source_size == size
&& snapshot_hash == source_hash
&& SourceStamp::read(&source)? == source_stamp
&& SourceStamp::read(&file)? == stamp,
"Interrupted snapshot does not match the selected shared version"
);
let record = Record {
version: 1,
content_id: content_id.into(),
source: source_stamp,
snapshot: stamp,
sha256: snapshot_hash.clone(),
size,
};
let checksum = hash(&serde_json::to_vec(&record)?);
io::save_record(&version, "snapshot.json", &Envelope { record, checksum })?;
use std::io::{Seek, SeekFrom};
file.seek(SeekFrom::Start(0))?;
return Ok(Snapshot {
file,
key,
sha256: snapshot_hash,
size,
});
}
Err(error)
if error
.downcast_ref::<std::io::Error>()
.is_some_and(|e| e.kind() == std::io::ErrorKind::NotFound) =>
{
()
}
Err(error) => return Err(error),
}
let used = storage_bytes(&root)?;
anyhow::ensure!(
used.checked_add(size)
.and_then(|v| v.checked_add(128 * 1024))
.is_some_and(|total| total <= max_total_bytes),
"Shared snapshot storage budget is full; no new purchase accepted"
);
let mut stat = std::mem::MaybeUninit::<libc::statvfs>::uninit();
anyhow::ensure!(
unsafe { libc::fstatvfs(root.as_raw_fd(), stat.as_mut_ptr()) } == 0,
"Could not verify snapshot disk space"
);
let stat = unsafe { stat.assume_init() };
let free = (stat.f_bavail as u64).saturating_mul(stat.f_frsize as u64);
anyhow::ensure!(
free >= size.saturating_add(minimum_free_bytes),
"Not enough free storage for a retained purchase snapshot"
);
let (name, mut destination) = io::temporary(&version)?;
let mut temporary = Temporary {
directory: &version,
name: name.clone(),
};
let (sha256, copied) =
io::hash_file(&mut source, Some(&mut destination), limits, &mut progress)?;
anyhow::ensure!(
copied == size && SourceStamp::read(&source)? == source_stamp,
"Shared file changed during snapshot creation; no purchase accepted"
);
destination.set_permissions(std::fs::Permissions::from_mode(0o400))?;
destination.sync_all()?;
io::publish_file(&version, &name, "media")?;
temporary.name.clear();
let file = io::open_at(&version, "media", libc::O_RDONLY | libc::O_NONBLOCK, 0)?;
let record = Record {
version: 1,
content_id: content_id.into(),
source: source_stamp,
snapshot: SourceStamp::read(&file)?,
sha256: sha256.clone(),
size,
};
let checksum = hash(&serde_json::to_vec(&record)?);
io::save_record(&version, "snapshot.json", &Envelope { record, checksum })?;
Ok(Snapshot {
file,
key,
sha256,
size,
})
}
/// Open retained accepted bytes even after the original Cloud source vanishes.
/// Caller supplies the durable contract's identity/hash/size, not a browser path.
pub(crate) fn open_matching(
data_dir: &Path,
content_id: &str,
sha256: &str,
size: u64,
) -> Result<Snapshot> {
let data = io::open_directory(&data_dir.canonicalize()?)?;
let root = io::open_at(
&data,
"content-snapshots",
libc::O_RDONLY | libc::O_DIRECTORY,
0,
)?;
let mut count = 0usize;
for entry in std::fs::read_dir(format!("/proc/self/fd/{}", root.as_raw_fd()))? {
let entry = entry?;
count += 1;
anyhow::ensure!(count <= 10000, "Snapshot index requires maintenance");
let key = entry
.file_name()
.to_str()
.context("Invalid snapshot version")?
.to_owned();
anyhow::ensure!(
key.len() == 64 && key.bytes().all(|v| v.is_ascii_hexdigit()),
"Unexpected snapshot version"
);
let version = io::open_at(&root, &key, libc::O_RDONLY | libc::O_DIRECTORY, 0)?;
let Some(record) = read_manifest(&version)? else {
continue;
};
if record.content_id != content_id || record.sha256 != sha256 || record.size != size {
continue;
}
let file = io::open_at(&version, "media", libc::O_RDONLY | libc::O_NONBLOCK, 0)?;
anyhow::ensure!(
record.version == 1
&& SourceStamp::read(&file)? == record.snapshot
&& file.metadata()?.mode() & 0o7777 == 0o400,
"Accepted snapshot changed; retain purchase for recovery"
);
return Ok(Snapshot {
file,
key,
sha256: record.sha256,
size: record.size,
});
}
anyhow::bail!("Accepted content snapshot is unavailable; retain purchase for recovery")
}
fn storage_bytes(directory: &File) -> Result<u64> {
let mut total = 0u64;
for entry in std::fs::read_dir(format!("/proc/self/fd/{}", directory.as_raw_fd()))? {
let entry = entry?;
let name = entry.file_name();
let name = name.to_str().context("Invalid snapshot filename")?;
let metadata = std::fs::symlink_metadata(entry.path())?;
anyhow::ensure!(
!metadata.file_type().is_symlink(),
"Snapshot storage contains an unexpected symlink"
);
let size = if metadata.is_dir() {
storage_bytes(&io::open_at(
directory,
name,
libc::O_RDONLY | libc::O_DIRECTORY,
0,
)?)?
} else {
anyhow::ensure!(metadata.is_file(), "Unexpected snapshot storage entry");
metadata.len()
};
total = total
.checked_add(size)
.context("Snapshot storage accounting overflow")?;
}
Ok(total)
}
struct Temporary<'a> {
directory: &'a File,
name: String,
}
impl Drop for Temporary<'_> {
fn drop(&mut self) {
if let Ok(name) = std::ffi::CString::new(self.name.as_str()) {
if !self.name.is_empty() {
unsafe {
libc::unlinkat(self.directory.as_raw_fd(), name.as_ptr(), 0);
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::{io::Read, sync::atomic::AtomicBool};
#[test]
fn quotes_reuse_one_version_and_accepted_bytes_survive_source_replacement() {
let data = tempfile::tempdir().unwrap();
let cloud = tempfile::tempdir().unwrap();
std::fs::write(cloud.path().join("film.mp4"), b"old film").unwrap();
let cancelled = AtomicBool::new(false);
let limits = Limits {
max_bytes: 1024,
cancelled: &cancelled,
};
let first = prepare(
data.path(),
cloud.path(),
"film",
Path::new("film.mp4"),
&limits,
1024 * 1024,
0,
|_| Ok(()),
)
.unwrap();
let mut copied = 0;
let replay = prepare(
data.path(),
cloud.path(),
"film",
Path::new("film.mp4"),
&limits,
1024 * 1024,
0,
|bytes| {
copied = bytes;
Ok(())
},
)
.unwrap();
assert_eq!(replay.key, first.key);
assert_eq!(copied, 0, "a new buyer quote must not recopy the file");
std::fs::write(cloud.path().join("film.mp4"), b"new movie").unwrap();
let next = prepare(
data.path(),
cloud.path(),
"film",
Path::new("film.mp4"),
&limits,
1024 * 1024,
0,
|_| Ok(()),
)
.unwrap();
assert_ne!(first.key, next.key);
std::fs::remove_file(cloud.path().join("film.mp4")).unwrap();
let mut retained = open_matching(data.path(), "film", &first.sha256, first.size).unwrap();
let mut bytes = Vec::new();
retained.file.read_to_end(&mut bytes).unwrap();
assert_eq!(bytes, b"old film");
assert!(open_matching(data.path(), "other", &first.sha256, first.size).is_err());
}
#[test]
fn quotas_cancellation_and_source_mutation_never_issue_a_snapshot() {
let data = tempfile::tempdir().unwrap();
let cloud = tempfile::tempdir().unwrap();
let path = cloud.path().join("film");
std::fs::write(&path, vec![7u8; 100000]).unwrap();
let cancelled = AtomicBool::new(false);
let limits = Limits {
max_bytes: 200000,
cancelled: &cancelled,
};
assert!(prepare(
data.path(),
cloud.path(),
"film",
Path::new("film"),
&limits,
50000,
0,
|_| Ok(())
)
.is_err());
assert!(prepare(
data.path(),
cloud.path(),
"film",
Path::new("film"),
&limits,
1024 * 1024,
0,
|_| {
std::fs::write(&path, b"changed")?;
Ok(())
}
)
.is_err());
cancelled.store(true, std::sync::atomic::Ordering::Relaxed);
assert!(prepare(
data.path(),
cloud.path(),
"film",
Path::new("film"),
&limits,
1024 * 1024,
0,
|_| Ok(())
)
.is_err());
}
#[test]
fn interrupted_manifest_commit_recovers_same_snapshot_without_replacing_bytes() {
let data = tempfile::tempdir().unwrap();
let cloud = tempfile::tempdir().unwrap();
std::fs::write(cloud.path().join("film"), b"film").unwrap();
let cancelled = AtomicBool::new(false);
let limits = Limits {
max_bytes: 1024,
cancelled: &cancelled,
};
let first = prepare(
data.path(),
cloud.path(),
"film",
Path::new("film"),
&limits,
1024 * 1024,
0,
|_| Ok(()),
)
.unwrap();
let inode = first.file.metadata().unwrap().ino();
std::fs::remove_file(
data.path()
.join("content-snapshots")
.join(&first.key)
.join("snapshot.json"),
)
.unwrap();
let restored = prepare(
data.path(),
cloud.path(),
"film",
Path::new("film"),
&limits,
1024 * 1024,
0,
|_| Ok(()),
)
.unwrap();
assert_eq!(restored.file.metadata().unwrap().ino(), inode);
assert_eq!(restored.sha256, first.sha256);
}
}
+1
View File
@@ -20,6 +20,7 @@ pub(crate) use invites::notify_join;
// Crate-internal: peer-joined resolves the granted trust level by matching
// the acceptor's invite_token against our stored outgoing invites.
pub(crate) use storage::load_invites;
pub(crate) use storage::load_unique_payment_peer;
#[allow(unused_imports)]
pub use storage::{
add_node, fips_npub_for_onion, load_nodes, load_removed_dids, record_peer_transport,
@@ -71,6 +71,36 @@ pub async fn load_nodes(data_dir: &Path) -> Result<Vec<FederatedNode>> {
load_nodes_inner(data_dir).await
}
/// Resolve payment identity from persisted records before display deduplication
/// can merge fields from different identities sharing an address.
pub(crate) async fn load_unique_payment_peer(
data_dir: &Path,
onion: &str,
) -> Result<FederatedNode> {
let _guard = FEDERATION_STORE_LOCK.lock().await;
let content = fs::read(data_dir.join(FEDERATION_DIR).join(NODES_FILE))
.await
.context("Could not read payment peer bindings")?;
let file: NodesFile =
serde_json::from_slice(&content).context("Invalid payment peer bindings")?;
let target = onion.strip_suffix(".onion").unwrap_or(onion);
anyhow::ensure!(!target.is_empty(), "Missing payment peer address");
let mut matches = file
.nodes
.into_iter()
.filter(|peer| peer.onion.strip_suffix(".onion").unwrap_or(&peer.onion) == target);
let peer = matches.next().context("Payment peer binding is missing")?;
anyhow::ensure!(
matches.next().is_none(),
"Payment peer binding is ambiguous"
);
anyhow::ensure!(
crate::identity::did_key_from_pubkey_hex(&peer.pubkey)? == peer.did,
"Payment peer identity does not match its public key"
);
Ok(peer)
}
/// Lock-free body of `load_nodes`. Callers that already hold
/// `FEDERATION_STORE_LOCK` (i.e. other functions in this module composing a
/// multi-step critical section) must call this instead of `load_nodes` to
+195 -9
View File
@@ -511,6 +511,66 @@ impl<'a> PeerRequest<'a> {
self.authenticate_content(data_dir).await?.send_get().await
}
/// Send a purchase request with the exact serialized bytes bound to its peer
/// proof. The expected seller must match the existing authenticated binding.
/// An ambiguous reply is recovered by the purchase journal, never by this
/// transport replaying a token through another route.
pub(crate) async fn send_content_json<B: serde::Serialize>(
self,
data_dir: &std::path::Path,
expected_seller: &str,
body: &B,
) -> Result<(reqwest::Response, crate::transport::TransportKind)> {
let (request, encoded) = self
.prepare_content_json(data_dir, expected_seller, body)
.await?;
request.send_encoded_json(&encoded).await
}
async fn prepare_content_json<B: serde::Serialize>(
mut self,
data_dir: &std::path::Path,
expected_seller: &str,
body: &B,
) -> Result<(Self, Vec<u8>)> {
anyhow::ensure!(
self.path.starts_with("/content/purchase/"),
"Not a purchase route"
);
let peer = crate::federation::load_unique_payment_peer(data_dir, self.onion_host).await?;
anyhow::ensure!(
peer.did == expected_seller,
"Purchase seller does not match the authenticated peer binding"
);
anyhow::ensure!(
self.fips_npub.is_some_and(|npub| !npub.is_empty())
&& peer.fips_npub.as_deref() == self.fips_npub,
"Purchase mesh route does not match the authenticated peer binding"
);
let identity =
crate::identity::NodeIdentity::load_existing(&data_dir.join("identity")).await?;
let encoded = serde_json::to_vec(body).context("Encode purchase request")?;
anyhow::ensure!(
encoded.len() <= 1024 * 1024,
"Purchase request is too large"
);
let proof = crate::content_auth::sign_request(
&identity,
expected_seller,
&hyper::Method::POST,
self.path,
&encoded,
chrono::Utc::now().timestamp(),
)?;
self.headers
.retain(|(name, _)| !name.eq_ignore_ascii_case(crate::content_auth::REQUEST_HEADER));
self.headers
.push((crate::content_auth::REQUEST_HEADER, proof));
self.single_delivery = true;
self.require_fips = true;
Ok((self, encoded))
}
pub fn header(mut self, name: &'a str, value: impl Into<String>) -> Self {
self.headers.push((name, value.into()));
self
@@ -537,12 +597,21 @@ impl<'a> PeerRequest<'a> {
pub async fn send_json<B: serde::Serialize>(
&self,
body: &B,
) -> Result<(reqwest::Response, crate::transport::TransportKind)> {
let encoded = serde_json::to_vec(body).context("Encode peer JSON request")?;
self.send_encoded_json(&encoded).await
}
// Serialize once: FIPS attempts and any permitted fallback use identical bytes.
async fn send_encoded_json(
&self,
body: &[u8],
) -> Result<(reqwest::Response, crate::transport::TransportKind)> {
use crate::settings::transport::TransportPref;
let pref = self.preference().await;
// FIPS-only or Auto: try FIPS first.
if matches!(pref, TransportPref::Auto | TransportPref::Fips) {
match self.try_fips_post_json(body).await? {
match self.try_fips_post_bytes(body).await? {
Some(resp) => {
// Use the FIPS reply unless it's one a Tor retry could
// fix (404 path-not-served / 5xx) and we're allowed to
@@ -574,7 +643,7 @@ impl<'a> PeerRequest<'a> {
}
}
}
let resp = self.send_tor_post_json(body).await?;
let resp = self.send_tor_post_bytes(body).await?;
self.spawn_record(crate::transport::TransportKind::Tor);
Ok((resp, crate::transport::TransportKind::Tor))
}
@@ -618,10 +687,7 @@ impl<'a> PeerRequest<'a> {
Ok((resp, crate::transport::TransportKind::Tor))
}
async fn try_fips_post_json<B: serde::Serialize>(
&self,
body: &B,
) -> Result<Option<reqwest::Response>> {
async fn try_fips_post_bytes(&self, body: &[u8]) -> Result<Option<reqwest::Response>> {
let Some(npub) = self.fips_npub else {
telemetry::record_fallback(FallbackReason::NoNpub);
return Ok(None);
@@ -657,7 +723,10 @@ impl<'a> PeerRequest<'a> {
budget
};
let c = client_with_delivery_policy(per_attempt, self.single_delivery || self.require_fips);
let mut rb = c.post(&url).json(body);
let mut rb = c
.post(&url)
.header(reqwest::header::CONTENT_TYPE, "application/json")
.body(body.to_vec());
for (k, v) in &self.headers {
rb = rb.header(*k, v);
}
@@ -770,10 +839,13 @@ impl<'a> PeerRequest<'a> {
}
}
async fn send_tor_post_json<B: serde::Serialize>(&self, body: &B) -> Result<reqwest::Response> {
async fn send_tor_post_bytes(&self, body: &[u8]) -> Result<reqwest::Response> {
let url = self.tor_url();
let client = self.tor_client()?;
let mut rb = client.post(&url).json(body);
let mut rb = client
.post(&url)
.header(reqwest::header::CONTENT_TYPE, "application/json")
.body(body.to_vec());
for (k, v) in &self.headers {
rb = rb.header(*k, v);
}
@@ -815,6 +887,120 @@ impl<'a> PeerRequest<'a> {
mod tests {
use super::*;
#[tokio::test]
async fn purchase_post_serializes_once_and_binds_the_actual_bytes_to_the_peer() {
use std::sync::atomic::{AtomicUsize, Ordering};
struct Counted(AtomicUsize);
impl serde::Serialize for Counted {
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
serializer.serialize_u64(self.0.fetch_add(1, Ordering::SeqCst) as u64)
}
}
let dir = tempfile::tempdir().unwrap();
let identity = crate::identity::NodeIdentity::load_or_create(&dir.path().join("identity"))
.await
.unwrap();
let seller = crate::identity::did_key_from_pubkey_hex(&hex::encode([7; 32])).unwrap();
tokio::fs::create_dir(dir.path().join("federation"))
.await
.unwrap();
let peer = serde_json::json!({"did":seller,"pubkey":hex::encode([7;32]),"onion":"seller.onion",
"trust_level":"trusted","added_at":"2026-10-06T00:00:00Z","fips_npub":"npub-test-seller"});
tokio::fs::write(
dir.path().join("federation/nodes.json"),
serde_json::to_vec(&serde_json::json!({"nodes":[peer.clone()]})).unwrap(),
)
.await
.unwrap();
let body = Counted(AtomicUsize::new(0));
let path = "/content/purchase/settle";
let (request, bytes) = PeerRequest::new(Some("npub-test-seller"), "seller.onion", path)
.prepare_content_json(dir.path(), &seller, &body)
.await
.unwrap();
assert_eq!(body.0.load(Ordering::SeqCst), 1);
assert_eq!(bytes, b"0");
assert!(request.require_fips && request.single_delivery);
let mut headers = hyper::HeaderMap::new();
for (key, value) in request.headers {
headers.insert(
hyper::header::HeaderName::from_bytes(key.as_bytes()).unwrap(),
value.parse().unwrap(),
);
}
assert_eq!(
crate::content_auth::authenticate_request(
&headers,
&seller,
&hyper::Method::POST,
path,
&bytes,
chrono::Utc::now().timestamp()
)
.unwrap(),
identity.did_key().unwrap()
);
assert!(crate::content_auth::authenticate_request(
&headers,
&seller,
&hyper::Method::POST,
path,
b"1",
chrono::Utc::now().timestamp()
)
.is_err());
for (npub, did) in [
(Some("wrong-route"), seller.as_str()),
(None, seller.as_str()),
(Some("npub-test-seller"), "wrong-seller"),
] {
assert!(PeerRequest::new(npub, "seller.onion", path)
.prepare_content_json(dir.path(), did, &body)
.await
.is_err());
}
assert_eq!(body.0.load(Ordering::SeqCst), 1);
tokio::fs::write(
dir.path().join("federation/nodes.json"),
serde_json::to_vec(&serde_json::json!({"nodes":[peer.clone(),peer.clone()]})).unwrap(),
)
.await
.unwrap();
assert!(
PeerRequest::new(Some("npub-test-seller"), "seller.onion", path)
.prepare_content_json(dir.path(), &seller, &body)
.await
.is_err()
);
let mut conflicting = peer.clone();
conflicting["did"] = serde_json::json!(crate::identity::did_key_from_pubkey_hex(
&hex::encode([8; 32])
)
.unwrap());
conflicting["pubkey"] = serde_json::json!(hex::encode([8; 32]));
conflicting["onion"] = serde_json::json!("seller");
let mut wrong_key = peer.clone();
wrong_key["pubkey"] = serde_json::json!(hex::encode([8; 32]));
for nodes in [
serde_json::json!([peer, conflicting]),
serde_json::json!([wrong_key]),
] {
tokio::fs::write(
dir.path().join("federation/nodes.json"),
serde_json::to_vec(&serde_json::json!({"nodes": nodes})).unwrap(),
)
.await
.unwrap();
assert!(
PeerRequest::new(Some("npub-test-seller"), "seller.onion", path)
.prepare_content_json(dir.path(), &seller, &body)
.await
.is_err()
);
}
assert_eq!(body.0.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn required_media_never_falls_back_when_peer_has_no_fips_identity() {
let request = PeerRequest::new(None, "unreachable.onion", "/content/video").require_fips();
+3
View File
@@ -46,8 +46,11 @@ mod content_indeehub;
mod content_invoice;
mod content_owned;
mod content_purchase;
mod content_purchase_executor;
mod content_snapshot;
mod media_stream;
mod media_registration;
mod registered_media;
mod prepared_media;
mod content_server;
mod crash_recovery;
+21 -13
View File
@@ -120,7 +120,7 @@ pub struct PreparedRegistration {
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct SourceStamp {
pub(crate) struct SourceStamp {
device: u64,
inode: u64,
size: u64,
@@ -130,7 +130,7 @@ struct SourceStamp {
changed_nanos: i64,
}
impl SourceStamp {
fn read(file: &File) -> Result<Self> {
pub(crate) fn read(file: &File) -> Result<Self> {
let m = file.metadata()?;
anyhow::ensure!(
m.is_file(),
@@ -254,7 +254,12 @@ fn fd_result(fd: libc::c_int) -> Result<File> {
// 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<File> {
pub(crate) fn open_at(
dir: &File,
name: &str,
flags: libc::c_int,
mode: libc::mode_t,
) -> Result<File> {
let name = CString::new(name)?;
// SAFETY: descriptor and NUL-terminated name remain alive through the call.
fd_result(unsafe {
@@ -266,7 +271,7 @@ fn open_at(dir: &File, name: &str, flags: libc::c_int, mode: libc::mode_t) -> Re
)
})
}
fn open_directory(path: &Path) -> Result<File> {
pub(crate) fn open_directory(path: &Path) -> Result<File> {
let path = c_path(path)?;
// SAFETY: path is a valid NUL-terminated string.
fd_result(unsafe {
@@ -276,7 +281,7 @@ fn open_directory(path: &Path) -> Result<File> {
)
})
}
fn private_directory(parent: &File, name: &str) -> Result<File> {
pub(crate) fn private_directory(parent: &File, name: &str) -> Result<File> {
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) };
@@ -298,7 +303,7 @@ fn private_directory(parent: &File, name: &str) -> Result<File> {
}
#[cfg(target_os = "linux")]
fn open_cloud_file(root: &File, relative: &Path) -> Result<File> {
pub(crate) fn open_cloud_file(root: &File, relative: &Path) -> Result<File> {
#[repr(C)]
struct OpenHow {
flags: u64,
@@ -340,7 +345,7 @@ fn open_cloud_file(root: &File, relative: &Path) -> Result<File> {
Ok(file)
}
#[cfg(not(target_os = "linux"))]
fn open_cloud_file(_root: &File, _relative: &Path) -> Result<File> {
pub(crate) fn open_cloud_file(_root: &File, _relative: &Path) -> Result<File> {
anyhow::bail!("Safe Cloud registration currently requires Linux openat2")
}
@@ -351,7 +356,7 @@ fn cancelled(limits: &Limits<'_>) -> Result<()> {
);
Ok(())
}
fn lock_operation(dir: &File, limits: &Limits<'_>, deadline: Instant) -> Result<()> {
pub(crate) fn lock_operation(dir: &File, limits: &Limits<'_>, deadline: Instant) -> Result<()> {
loop {
cancelled(limits)?;
anyhow::ensure!(
@@ -371,7 +376,10 @@ fn lock_operation(dir: &File, limits: &Limits<'_>, deadline: Instant) -> Result<
std::thread::sleep(Duration::from_millis(20));
}
}
fn read_record<T: serde::de::DeserializeOwned>(dir: &File, name: &str) -> Result<Option<T>> {
pub(crate) fn read_record<T: serde::de::DeserializeOwned>(
dir: &File,
name: &str,
) -> Result<Option<T>> {
let file = match open_at(dir, name, libc::O_RDONLY | libc::O_NONBLOCK, 0) {
Ok(file) => file,
Err(error)
@@ -398,7 +406,7 @@ fn read_record<T: serde::de::DeserializeOwned>(dir: &File, name: &str) -> Result
"Damaged registration record; preserve it for recovery",
)?))
}
fn temporary(dir: &File) -> Result<(String, File)> {
pub(crate) fn temporary(dir: &File) -> Result<(String, File)> {
let name = format!("pending-{}", uuid::Uuid::new_v4());
Ok((
name.clone(),
@@ -410,7 +418,7 @@ fn temporary(dir: &File) -> Result<(String, File)> {
)?,
))
}
fn publish_file(dir: &File, temporary: &str, final_name: &str) -> Result<()> {
pub(crate) 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
@@ -434,7 +442,7 @@ fn publish_file(dir: &File, temporary: &str, final_name: &str) -> Result<()> {
dir.sync_all()?;
Ok(())
}
fn save_record<T: Serialize>(dir: &File, name: &str, value: &T) -> Result<()> {
pub(crate) fn save_record<T: Serialize>(dir: &File, name: &str, value: &T) -> Result<()> {
let bytes = serde_json::to_vec(value)?;
anyhow::ensure!(
bytes.len() as u64 <= MAX_RECORD_BYTES,
@@ -446,7 +454,7 @@ fn save_record<T: Serialize>(dir: &File, name: &str, value: &T) -> Result<()> {
publish_file(dir, &temporary, name)
}
fn hash_file(
pub(crate) fn hash_file(
file: &mut File,
mut output: Option<&mut File>,
limits: &Limits<'_>,
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,26 @@
{
"description": "Independent Node.js ordered-array SHA256 fixture using the public registration test vector.",
"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.registered-media.terms.v1\",\"did:key:z6MkvDqGT54cXesYGvABpF1UapVNwjCqRcafi4Px6Thv5T3Z\",\"fixture-indeehub\",\"cdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcd\",\"fixture-project\",\"registered_00000000-0000-4000-8000-000000000001\",\"bbe573fdac96c464b67257a97c8667713c124c80927c01bab75b7178f4c0a661\",\"23\",15,3600,[\"cashu\",\"lightning-cashu\"]]",
"sha256": "a0d104ece82fd9d2094dbe4e6fd8285889ccb030407c42d9098d3c74faf02608"
}
+62 -1
View File
@@ -813,6 +813,61 @@ pub async fn send_token_recoverable(
mint_url: &str,
amount_sats: u64,
context_hash: &str,
) -> Result<String> {
send_token_recoverable_with_deadline(
data_dir,
operation_id,
network,
mint_url,
amount_sats,
context_hash,
None,
|| chrono::Utc::now().timestamp(),
)
.await
}
/// Recover original results after expiry, but never initiate a fresh exact send
/// or swap POST after the immutable purchase deadline. The caller binds this
/// deadline into context_hash and must already have seller acceptance saved.
pub async fn send_token_recoverable_before(
data_dir: &Path,
operation_id: &str,
network: EcashNetwork,
mint_url: &str,
amount_sats: u64,
context_hash: &str,
expires_at: i64,
) -> Result<String> {
anyhow::ensure!(expires_at > 0, "Invalid purchase spend deadline");
send_token_recoverable_with_deadline(
data_dir,
operation_id,
network,
mint_url,
amount_sats,
context_hash,
Some(expires_at),
|| chrono::Utc::now().timestamp(),
)
.await
}
fn ensure_fresh_payment_allowed(expires_at: Option<i64>, now: i64) -> Result<()> {
anyhow::ensure!(
expires_at.is_none_or(|deadline| now < deadline),
"The offer expired; recover the original purchase instead of starting another payment"
);
Ok(())
}
async fn send_token_recoverable_with_deadline(
data_dir: &Path,
operation_id: &str,
network: EcashNetwork,
mint_url: &str,
amount_sats: u64,
context_hash: &str,
expires_at: Option<i64>,
now: impl Fn() -> i64 + Send + Sync,
) -> Result<String> {
use super::send_journal::{Binding, Journal, Outcome, Phase, Request};
let held = super::mutation::guard(data_dir).await?;
@@ -837,6 +892,7 @@ pub async fn send_token_recoverable(
);
record
} else {
ensure_fresh_payment_allowed(expires_at, now())?;
anyhow::ensure!(amount_sats > 0, "Payment amount must be positive");
let wallet = load_wallet(data_dir).await?;
let (indices, excess) = wallet
@@ -868,6 +924,7 @@ pub async fn send_token_recoverable(
);
Request::Swap(prepared)
};
ensure_fresh_payment_allowed(expires_at, now())?;
journal.prepare(binding.clone(), request).await?
};
if matches!(record.phase, Phase::Result(_) | Phase::Committed(_)) {
@@ -875,7 +932,10 @@ pub async fn send_token_recoverable(
}
journal.reserve_wallet(&binding).await?;
let (send, change) = match &record.request {
Request::Exact { proofs } => (proofs.clone(), vec![]),
Request::Exact { proofs } => {
ensure_fresh_payment_allowed(expires_at, now())?;
(proofs.clone(), vec![])
}
Request::Swap(prepared) => {
// All recovery material is already durable. Do not derive new
// outputs or release reservations after an ambiguous response.
@@ -901,6 +961,7 @@ pub async fn send_token_recoverable(
"This payment is still pending at the mint; do not pay again"
);
}
ensure_fresh_payment_allowed(expires_at, now())?;
client.execute_prepared_swap(prepared).await.map_err(|_| anyhow::anyhow!(
"The mint did not confirm this payment; retry this same operation to recover it"))?
};
+2
View File
@@ -13,3 +13,5 @@ pub mod nut13;
pub mod profits;
mod send_journal;
mod receive_journal;
pub(crate) mod purchase_fee_plan;
+363 -4
View File
@@ -23,6 +23,9 @@ struct Mint {
lose_swap_reply: Arc<std::sync::atomic::AtomicBool>,
restore_reply: Arc<Mutex<Option<Value>>>,
state_reply: Arc<Mutex<Option<Value>>>,
// Per-fixture clock advancement proves the spend deadline is checked after
// mint preflight; no process-global clock or timing-sensitive sleep.
restore_clock: Arc<Mutex<Option<(Arc<std::sync::atomic::AtomicI64>, i64)>>>,
}
impl Drop for Mint {
fn drop(&mut self) {
@@ -65,6 +68,8 @@ impl Mint {
let restore_override = restore_reply.clone();
let state_reply = Arc::new(Mutex::new(None::<Value>));
let state_override = state_reply.clone();
let restore_clock = Arc::new(Mutex::new(None::<(Arc<std::sync::atomic::AtomicI64>, i64)>));
let advance_clock = restore_clock.clone();
let service = make_service_fn(move |_| {
let seen = seen.clone();
let rejection = rejection.clone();
@@ -73,6 +78,7 @@ impl Mint {
let lose_reply = lose_reply.clone();
let restore_override = restore_override.clone();
let state_override = state_override.clone();
let advance_clock = advance_clock.clone();
async move {
Ok::<_, Infallible>(service_fn(move |req: Request<Body>| {
let seen = seen.clone();
@@ -82,6 +88,7 @@ impl Mint {
let lose_reply = lose_reply.clone();
let restore_override = restore_override.clone();
let state_override = state_override.clone();
let advance_clock = advance_clock.clone();
async move {
let mut status = 200;
let body = match req.uri().path() {
@@ -164,6 +171,11 @@ impl Mint {
})
}
"/v1/restore" => {
if let Some((clock, deadline)) =
advance_clock.lock().unwrap().as_ref()
{
clock.store(*deadline, std::sync::atomic::Ordering::SeqCst);
}
let body: Value = serde_json::from_slice(
&hyper::body::to_bytes(req.into_body()).await.unwrap(),
)
@@ -219,6 +231,7 @@ impl Mint {
lose_swap_reply,
restore_reply,
state_reply,
restore_clock,
}
}
async fn wallet(&self) -> tempfile::TempDir {
@@ -486,10 +499,17 @@ async fn paid_file_gate_delivers_bytes_only_after_payment_and_does_not_charge_mi
}
assert_eq!(load_wallet(seller.path()).await.unwrap().balance(), 128);
} else {
assert!(matches!(
result,
ServeResult::NotFound | ServeResult::PaymentRequired(_)
));
if !exists {
assert!(matches!(result, ServeResult::Unavailable));
assert!(content_server::load_catalog(seller.path())
.await
.unwrap()
.items
.iter()
.any(|item| item.id == "paid-test"));
} else {
assert!(matches!(result, ServeResult::PaymentRequired(_)));
}
assert_eq!(load_wallet(seller.path()).await.unwrap().balance(), 0);
assert!(mint.requests.lock().unwrap().is_empty());
}
@@ -1574,3 +1594,342 @@ async fn recoverable_receive_unspent_retry_reuses_exact_request_and_waits_for_ve
assert_eq!(wallet.transactions.len(), 1);
assert_eq!(wallet.receive_commits.len(), 1);
}
#[tokio::test]
async fn purchase_deadline_rejects_fresh_exact_send_but_recovers_prior_committed_token() {
let root = tempfile::tempdir().unwrap();
let mint = "https://unused-mint.invalid";
let mut wallet = WalletState::default();
wallet.mint_url = mint.into();
wallet.add_proofs(mint, vec![proof(ACTIVE, 8)]);
save_wallet(root.path(), &wallet).await.unwrap();
let id = uuid::Uuid::new_v4().to_string();
let context = "ab".repeat(32);
assert!(send_token_recoverable_before(
root.path(),
&id,
EcashNetwork::Mainnet,
mint,
8,
&context,
1
)
.await
.is_err());
assert_eq!(load_wallet(root.path()).await.unwrap().balance(), 8);
// Fixture a prior completed send; an expired caller must recover its result,
// not require another payment or credit the old proofs back to the purse.
let original =
send_token_recoverable(root.path(), &id, EcashNetwork::Mainnet, mint, 8, &context)
.await
.unwrap();
assert_eq!(
send_token_recoverable_before(
root.path(),
&id,
EcashNetwork::Mainnet,
mint,
8,
&context,
1
)
.await
.unwrap(),
original
);
assert_eq!(load_wallet(root.path()).await.unwrap().balance(), 0);
}
#[tokio::test]
async fn expired_purchase_recovers_issued_swap_but_never_posts_still_unspent_inputs() {
for issued in [false, true] {
let mint = Mint::start(0, if issued { None } else { Some(503) }).await;
let root = mint.wallet().await;
let mut wallet = load_wallet(root.path()).await.unwrap();
wallet.add_proofs(&mint.url, vec![proof(ACTIVE, 8)]);
save_wallet(root.path(), &wallet).await.unwrap();
let id = uuid::Uuid::new_v4().to_string();
let context = "ab".repeat(32);
mint.lose_swap_reply
.store(issued, std::sync::atomic::Ordering::SeqCst);
assert!(send_token_recoverable(
root.path(),
&id,
EcashNetwork::Mainnet,
&mint.url,
4,
&context
)
.await
.is_err());
mint.failure.store(0, std::sync::atomic::Ordering::SeqCst);
let result = send_token_recoverable_before(
root.path(),
&id,
EcashNetwork::Mainnet,
&mint.url,
4,
&context,
1,
)
.await;
if issued {
assert_eq!(
CashuToken::deserialize(&result.unwrap())
.unwrap()
.total_amount(),
4
);
assert_eq!(load_wallet(root.path()).await.unwrap().balance(), 4);
} else {
assert!(result.unwrap_err().to_string().contains("expired"));
assert_eq!(load_wallet(root.path()).await.unwrap().balance(), 0);
assert!(load_wallet(root.path())
.await
.unwrap()
.proofs
.iter()
.any(|p| p.reserved));
}
assert_eq!(mint.requests.lock().unwrap().len(), 1);
}
}
#[tokio::test]
async fn purchase_executor_recovers_prior_buyer_spend_and_delayed_accepted_seller_settlement() {
use crate::content_purchase::{BuyerPhase, Contract, Journal};
use crate::content_purchase_executor::{prepare_buyer_token, settle_seller_token};
let mint = Mint::start(0, None).await;
let buyer = mint.wallet().await;
let seller = mint.wallet().await;
let contract = Contract {
version: 1,
id: uuid::Uuid::new_v4().to_string(),
buyer_did: crate::identity::did_key_from_pubkey_hex(&hex::encode([1; 32])).unwrap(),
seller_did: crate::identity::did_key_from_pubkey_hex(&hex::encode([2; 32])).unwrap(),
content_id: "film-1".into(),
content_sha256: "ab".repeat(32),
content_size: 1024,
terms_sha256: "cd".repeat(32),
network: EcashNetwork::Mainnet,
mint_url: mint.url.clone(),
gross_token_sats: 8,
minimum_net_sats: 8,
offered_at: 1000,
expires_at: 2000,
};
let mut wallet = load_wallet(buyer.path()).await.unwrap();
wallet.add_proofs(&mint.url, vec![proof(ACTIVE, 8)]);
save_wallet(buyer.path(), &wallet).await.unwrap();
// An expired offer without durable seller acceptance cannot redeem a token
// or create settlement state, even when the token and buyer are otherwise valid.
let unaccepted_token = CashuToken::new(&mint.url, vec![proof(ACTIVE, 8)])
.serialize()
.unwrap();
let seller_wallet_before = std::fs::read(seller.path().join("wallet/ecash.json")).ok();
let rejected = settle_seller_token(
seller.path(),
&contract,
&unaccepted_token,
&contract.buyer_did,
)
.await
.err()
.unwrap();
assert!(rejected
.to_string()
.contains("Seller intent is not durable"));
assert!(mint.requests.lock().unwrap().is_empty());
assert_eq!(
std::fs::read(seller.path().join("wallet/ecash.json")).ok(),
seller_wallet_before
);
{
let journal = Journal::open(seller.path()).await.unwrap();
assert!(journal.seller(&contract.id).await.unwrap().is_none());
}
// Deterministically fixture durable acceptance before the historical deadline.
let accepted = {
let journal = Journal::open(seller.path()).await.unwrap();
journal
.prepare_seller(&contract, 1500)
.await
.unwrap()
.acceptance()
.unwrap()
};
{
let journal = Journal::open(buyer.path()).await.unwrap();
journal.prepare_buyer(&contract, 1500).await.unwrap();
assert!(journal
.record_acceptance(&contract, &accepted, &contract.buyer_did)
.await
.is_err());
journal
.record_acceptance(&contract, &accepted, &contract.seller_did)
.await
.unwrap();
}
// Fixture a prior wallet commit, interrupted before the buyer journal saved
// its token. The actual executor must recover that result after expiry.
let original = send_token_recoverable(
buyer.path(),
&contract.id,
contract.network,
&mint.url,
8,
&contract.context_hash().unwrap(),
)
.await
.unwrap();
assert_eq!(
prepare_buyer_token(buyer.path(), &contract).await.unwrap(),
original
);
assert!(mint.requests.lock().unwrap().is_empty());
mint.lose_swap_reply
.store(true, std::sync::atomic::Ordering::SeqCst);
assert!(
settle_seller_token(seller.path(), &contract, &original, &contract.buyer_did)
.await
.is_err()
);
assert_eq!(mint.requests.lock().unwrap().len(), 1);
let receipt = settle_seller_token(seller.path(), &contract, &original, &contract.buyer_did)
.await
.unwrap();
assert!(
settle_seller_token(seller.path(), &contract, &original, &contract.buyer_did)
.await
.unwrap()
== receipt
);
assert!(
settle_seller_token(seller.path(), &contract, &original, &contract.seller_did)
.await
.is_err()
);
let altered = CashuToken::new(&mint.url, vec![proof(V2, 8)])
.serialize()
.unwrap();
assert!(
settle_seller_token(seller.path(), &contract, &altered, &contract.buyer_did)
.await
.is_err()
);
assert_eq!(mint.requests.lock().unwrap().len(), 1);
assert_eq!(load_wallet(buyer.path()).await.unwrap().balance(), 0);
assert_eq!(load_wallet(seller.path()).await.unwrap().balance(), 8);
let journal = Journal::open(buyer.path()).await.unwrap();
journal.record_receipt(&contract, &receipt).await.unwrap();
assert_eq!(
journal.buyer(&contract.id).await.unwrap().unwrap().phase,
BuyerPhase::ReceiptSaved
);
}
#[tokio::test]
async fn purchase_expiring_during_mint_preflight_never_reserves_or_posts_swap() {
use std::sync::atomic::{AtomicI64, Ordering};
let mint = Mint::start(0, None).await;
let root = mint.wallet().await;
let mut wallet = load_wallet(root.path()).await.unwrap();
wallet.add_proofs(&mint.url, vec![proof(ACTIVE, 8)]);
save_wallet(root.path(), &wallet).await.unwrap();
let wallet_before = std::fs::read(root.path().join("wallet/ecash.json")).unwrap();
let clock = Arc::new(AtomicI64::new(1000));
*mint.restore_clock.lock().unwrap() = Some((clock.clone(), 2000));
let id = uuid::Uuid::new_v4().to_string();
let error = send_token_recoverable_with_deadline(
root.path(),
&id,
EcashNetwork::Mainnet,
&mint.url,
4,
&"ab".repeat(32),
Some(2000),
|| clock.load(Ordering::SeqCst),
)
.await
.unwrap_err();
assert_eq!(
clock.load(Ordering::SeqCst),
2000,
"preflight must have completed"
);
assert!(error.to_string().contains("expired"));
assert!(
mint.requests.lock().unwrap().is_empty(),
"no swap POST after expiry"
);
assert_eq!(
std::fs::read(root.path().join("wallet/ecash.json")).unwrap(),
wallet_before
);
assert_eq!(load_wallet(root.path()).await.unwrap().balance(), 8);
let held = crate::wallet::mutation::guard(root.path()).await.unwrap();
assert!(crate::wallet::send_journal::Journal::new(&held)
.load(&id)
.await
.unwrap()
.is_none());
}
// Append to wallet/payment_tests.rs when the next payment-plan batch is applied.
#[tokio::test]
async fn stale_planned_inputs_do_not_reselect_wallet_coins_or_post_to_mint() {
use crate::wallet::{mutation, send_journal};
let mint = Mint::start(0, None).await;
let root = tempfile::tempdir().unwrap();
let mut wallet = WalletState::default();
wallet.mint_url = mint.url.clone();
let original = proof(ACTIVE, 8);
wallet.add_proofs(&mint.url, vec![original.clone()]);
save_wallet(root.path(), &wallet).await.unwrap();
let id = uuid::Uuid::new_v4().to_string();
let context = "ab".repeat(32);
let prepared = MintClient::new(&mint.url)
.unwrap()
.prepare_swap_at_least(&[original], &[4, 4], 4)
.await
.unwrap();
{
let guard = mutation::guard(root.path()).await.unwrap();
send_journal::Journal::new(&guard)
.prepare(
send_journal::Binding {
id: id.clone(),
network: EcashNetwork::Mainnet,
mint_url: mint.url.clone(),
amount_sats: 4,
context_hash: context.clone(),
},
send_journal::Request::Swap(prepared),
)
.await
.unwrap();
}
// Simulate a selected proof becoming unavailable before reservation. Other
// sufficient coins exist, but the accepted immutable plan cannot select them.
wallet.proofs.clear();
let mut replacement = proof(ACTIVE, 8);
replacement.secret = "replacement-not-in-plan".into();
wallet.add_proofs(&mint.url, vec![replacement]);
save_wallet(root.path(), &wallet).await.unwrap();
let before = std::fs::read(root.path().join("wallet/ecash.json")).unwrap();
assert!(send_token_recoverable(
root.path(),
&id,
EcashNetwork::Mainnet,
&mint.url,
4,
&context
)
.await
.is_err());
assert_eq!(
std::fs::read(root.path().join("wallet/ecash.json")).unwrap(),
before
);
assert!(mint.requests.lock().unwrap().is_empty());
}
@@ -0,0 +1,293 @@
//! Quote the actual outgoing proof shape; never assume every exact-send proof
//! belongs to the current active keyset. The caller supplies mint-verified fees.
use super::cashu::{matches_stored_keyset_id, CashuToken, KeysetInfo, Proof};
use anyhow::{Context, Result};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use std::collections::{BTreeMap, HashSet};
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct KeysetPlan {
pub keyset_id: String,
pub denominations: Vec<u64>,
pub input_fee_ppk: u64,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct FeePlan {
pub mint_url: String,
pub keysets: Vec<KeysetPlan>,
pub gross_sats: u64,
pub fee_sats: u64,
pub net_sats: u64,
}
impl FeePlan {
pub fn from_shape(mint_url: &str, mut keysets: Vec<KeysetPlan>) -> Result<Self> {
let parsed = reqwest::Url::parse(mint_url)?;
anyhow::ensure!(
matches!(parsed.scheme(), "http" | "https")
&& parsed.host_str().is_some()
&& parsed.username().is_empty()
&& parsed.password().is_none()
&& parsed.query().is_none()
&& parsed.fragment().is_none()
&& parsed.to_string().trim_end_matches('/') == mint_url,
"Invalid canonical payment mint"
);
anyhow::ensure!(
!keysets.is_empty() && keysets.len() <= 64,
"Invalid payment keyset plan"
);
keysets.sort_by(|a, b| a.keyset_id.cmp(&b.keyset_id));
let mut identifiers = HashSet::new();
let mut count = 0usize;
let mut gross = 0u64;
let mut fee_ppk = 0u64;
for group in &mut keysets {
anyhow::ensure!(
matches!(group.keyset_id.len(), 16 | 66)
&& hex::decode(&group.keyset_id).is_ok()
&& group.keyset_id == group.keyset_id.to_ascii_lowercase()
&& identifiers.insert(group.keyset_id.clone()),
"Invalid or duplicate payment keyset"
);
// Full v2 identities are required in a plan; compact wire IDs only
// match an already known unambiguous group at validation time.
anyhow::ensure!(
group.keyset_id.len() != 16 || !group.keyset_id.starts_with("01"),
"Payment plan requires full v2 keyset identity"
);
anyhow::ensure!(!group.denominations.is_empty(), "Empty payment proof group");
group.denominations.sort_unstable();
count = count
.checked_add(group.denominations.len())
.context("Proof count overflow")?;
anyhow::ensure!(count <= 1024, "Too many outgoing payment proofs");
for amount in &group.denominations {
anyhow::ensure!(amount.is_power_of_two(), "Invalid payment denomination");
gross = gross
.checked_add(*amount)
.context("Payment amount overflow")?;
fee_ppk = fee_ppk
.checked_add(group.input_fee_ppk)
.context("Payment fee overflow")?;
}
}
let fee = fee_ppk.div_ceil(1000);
let net = gross
.checked_sub(fee)
.context("Payment fees exceed outgoing value")?;
anyhow::ensure!(net > 0, "Payment has no net value");
Ok(Self {
mint_url: mint_url.into(),
keysets,
gross_sats: gross,
fee_sats: fee,
net_sats: net,
})
}
pub fn validate(&self) -> Result<()> {
anyhow::ensure!(
&Self::from_shape(&self.mint_url, self.keysets.clone())? == self,
"Payment fee plan is inconsistent or noncanonical"
);
Ok(())
}
pub fn commitment(&self) -> Result<String> {
self.validate()?;
let groups: Vec<_> = self
.keysets
.iter()
.map(|group| {
serde_json::json!([group.keyset_id, group.input_fee_ppk, group.denominations])
})
.collect();
Ok(hex::encode(Sha256::digest(serde_json::to_vec(
&serde_json::json!([
"content-payment-fee-plan-v1",
self.mint_url,
self.gross_sats,
self.fee_sats,
self.net_sats,
groups
]),
)?)))
}
/// Check authoritative mint metadata at acceptance, and again immediately
/// before a fresh swap POST. Inactive old proofs may still be exact-spent;
/// a prepared output keyset must remain active for a new issuance.
pub fn verify_mint_keysets(
&self,
keysets: &[KeysetInfo],
require_active_output: bool,
) -> Result<()> {
self.validate()?;
for group in &self.keysets {
let mut matching = keysets
.iter()
.filter(|keyset| keyset.id.eq_ignore_ascii_case(&group.keyset_id));
let keyset = matching
.next()
.context("Quoted payment keyset is missing")?;
anyhow::ensure!(
matching.next().is_none()
&& keyset.unit == "sat"
&& keyset.input_fee_ppk == group.input_fee_ppk
&& (!require_active_output || keyset.active),
"Quoted payment keyset or fees changed"
);
}
Ok(())
}
pub fn validate_proofs(&self, proofs: &[Proof]) -> Result<()> {
self.validate()?;
let mut actual: BTreeMap<String, Vec<u64>> = BTreeMap::new();
let mut secrets = HashSet::new();
for proof in proofs {
anyhow::ensure!(
secrets.insert(proof.secret.as_str()),
"Duplicate outgoing proof"
);
let mut groups = self
.keysets
.iter()
.filter(|group| matches_stored_keyset_id(&proof.id, &group.keyset_id));
let group = groups.next().context("Outgoing proof keyset changed")?;
anyhow::ensure!(groups.next().is_none(), "Ambiguous compact outgoing keyset");
actual
.entry(group.keyset_id.clone())
.or_default()
.push(proof.amount);
}
for amounts in actual.values_mut() {
amounts.sort_unstable();
}
let expected: BTreeMap<_, _> = self
.keysets
.iter()
.map(|group| (group.keyset_id.clone(), group.denominations.clone()))
.collect();
anyhow::ensure!(
actual == expected,
"Outgoing proof count or denominations changed"
);
Ok(())
}
pub fn validate_token(&self, encoded: &str) -> Result<()> {
let token = CashuToken::deserialize(encoded)?;
anyhow::ensure!(
token.token.len() == 1 && token.token[0].mint.trim_end_matches('/') == self.mint_url,
"Outgoing payment mint changed"
);
self.validate_proofs(&token.token[0].proofs)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn group(id: &str, denominations: Vec<u64>, fee: u64) -> KeysetPlan {
KeysetPlan {
keyset_id: id.into(),
denominations,
input_fee_ppk: fee,
}
}
fn proof(id: &str, amount: u64, secret: &str) -> Proof {
Proof {
id: id.into(),
amount,
secret: secret.into(),
c: "unused-in-shape-validation".into(),
}
}
#[test]
fn mixed_keysets_quote_actual_proof_counts_and_round_fee_once() {
let first = "0011223344556677";
let second = "008899aabbccddee";
let plan = FeePlan::from_shape(
"https://mint.invalid",
vec![group(first, vec![4, 4], 400), group(second, vec![2], 100)],
)
.unwrap();
assert_eq!((plan.gross_sats, plan.fee_sats, plan.net_sats), (10, 1, 9));
let proofs = vec![
proof(first, 4, "a"),
proof(first, 4, "b"),
proof(second, 2, "c"),
];
plan.validate_proofs(&proofs).unwrap();
assert!(plan
.validate_proofs(&[proof(first, 8, "a"), proof(second, 2, "c")])
.is_err());
assert!(plan
.validate_proofs(&[
proof(first, 4, "a"),
proof(first, 4, "a"),
proof(second, 2, "c")
])
.is_err());
let mut changed = plan.clone();
changed.keysets[0].input_fee_ppk += 1000;
assert!(changed.validate().is_err());
}
#[test]
fn compact_collision_and_changed_fee_quote_reject_before_token_transfer() {
let first = format!("01{}", "11".repeat(32));
let mut second = first.clone();
second.replace_range(64..66, "22");
let plan = FeePlan::from_shape(
"https://mint.invalid",
vec![group(&first, vec![4], 0), group(&second, vec![4], 0)],
)
.unwrap();
plan.validate_proofs(&[proof(&first, 4, "a"), proof(&second, 4, "b")])
.unwrap();
assert!(plan
.validate_proofs(&[proof(&first[..16], 4, "a"), proof(&second, 4, "b")])
.is_err());
let mut keysets: Vec<_> = [&first, &second]
.into_iter()
.map(|id| KeysetInfo {
id: id.clone(),
unit: "sat".into(),
active: true,
input_fee_ppk: 0,
})
.collect();
plan.verify_mint_keysets(&keysets, true).unwrap();
keysets[0].active = false;
assert!(plan.verify_mint_keysets(&keysets, true).is_err());
plan.verify_mint_keysets(&keysets, false).unwrap();
keysets[0].input_fee_ppk = 1;
assert!(plan.verify_mint_keysets(&keysets, false).is_err());
}
#[test]
fn fee_commitment_matches_shared_ordered_array_fixture() {
let fixture: serde_json::Value = serde_json::from_str(include_str!(
"../../../../tests/fixtures/purchase-fee-plan-v1.json"
))
.unwrap();
let plan: FeePlan = serde_json::from_value(fixture["plan"].clone()).unwrap();
assert_eq!(
plan.commitment().unwrap(),
fixture["sha256"].as_str().unwrap()
);
assert_eq!(
hex::encode(Sha256::digest(
fixture["preimage"].as_str().unwrap().as_bytes()
)),
fixture["sha256"].as_str().unwrap()
);
let mut reversed = plan.keysets.clone();
reversed.reverse();
assert_eq!(
FeePlan::from_shape(&plan.mint_url, reversed)
.unwrap()
.commitment()
.unwrap(),
plan.commitment().unwrap()
);
}
}
+28
View File
@@ -141,6 +141,12 @@ impl HookStep {
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct ContainerConfig {
/// Opt in to installer-provisioned public node identity and stable app audience
/// for the reviewed media-registration bridge. Does not enable registration,
/// publication or signing permissions. Missing/corrupt existing pins fail.
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub media_registration_identity: bool,
/// Pull source. Mutually exclusive with `build`. Exactly one of the two must be present.
#[serde(default)]
pub image: Option<String>,
@@ -1981,6 +1987,24 @@ app:
assert!(AppManifest::parse(yaml).is_err());
}
#[test]
fn media_registration_identity_is_explicit_and_preserved() {
let ordinary = AppManifest::parse("app:\n id: fixture\n name: Fixture\n version: '1'\n container:\n image: fixture:1\n").unwrap();
assert!(!ordinary.app.container.media_registration_identity);
assert!(!serde_yaml::to_string(&ordinary)
.unwrap()
.contains("media_registration_identity"));
let opted_in = AppManifest::parse("app:\n id: fixture\n name: Fixture\n version: '1'\n container:\n image: fixture:1\n media_registration_identity: true\n").unwrap();
assert!(opted_in.app.container.media_registration_identity);
assert!(
AppManifest::parse(&serde_yaml::to_string(&opted_in).unwrap())
.unwrap()
.app
.container
.media_registration_identity
);
}
#[test]
fn test_manifest_parse() {
let yaml = r#"
@@ -2473,6 +2497,7 @@ app:
#[test]
fn resolve_derived_env_renders_host_facts() {
let c = ContainerConfig {
media_registration_identity: false,
image: Some("x:latest".to_string()),
image_signature: None,
pull_policy: "if-not-present".to_string(),
@@ -2532,6 +2557,7 @@ app:
#[test]
fn resolve_secret_env_reads_from_provider() {
let c = ContainerConfig {
media_registration_identity: false,
image: Some("x:latest".to_string()),
image_signature: None,
pull_policy: "if-not-present".to_string(),
@@ -2580,6 +2606,7 @@ app:
#[test]
fn resolve_secret_env_rejects_empty_value() {
let c = ContainerConfig {
media_registration_identity: false,
image: Some("x:latest".to_string()),
image_signature: None,
pull_policy: "if-not-present".to_string(),
@@ -2616,6 +2643,7 @@ app:
#[test]
fn resolve_secret_env_skips_missing_or_empty_optional_entries() {
let c = ContainerConfig {
media_registration_identity: false,
image: Some("x:latest".to_string()),
image_signature: None,
pull_policy: "if-not-present".to_string(),