Draft explicit rental readiness and start with verified chunk delivery

This commit is contained in:
archipelago
2026-10-07 00:23:43 -04:00
parent 697aabeec3
commit ed96df0ac3
25 changed files with 1914 additions and 129 deletions
+42 -15
View File
@@ -1,16 +1,16 @@
mod purchase;
mod cloud_purchase;
mod blob;
mod cdp;
mod cloud_purchase;
mod content;
mod registered_media;
mod rental_playback;
mod dwn;
mod model_proxy;
mod node_message;
mod proxy;
mod purchase;
mod registered_media;
mod remote_input;
mod remote_relay;
mod rental_playback;
mod routstr_proxy;
mod websocket;
@@ -389,7 +389,9 @@ impl ApiHandler {
let method = req.method().clone();
if path.starts_with("/api/rental-playback/") {
return self.handle_local_rental_request(&method, &path, req.headers()).await;
return self
.handle_local_rental_request(&method, &path, req.headers())
.await;
}
// Handle CORS preflight for all routes
@@ -453,13 +455,28 @@ impl ApiHandler {
}
// Purchase routes bound the original body before the generic buffer.
if method == Method::POST && matches!(path.as_str(),
crate::content_purchase_protocol::OFFER_ROUTE | crate::content_purchase_protocol::ACCEPT_ROUTE
| crate::content_purchase_protocol::SETTLE_ROUTE | crate::content_purchase_protocol::STATUS_ROUTE
| crate::content_purchase_protocol::CANCEL_ROUTE) {
if method == Method::POST
&& matches!(
path.as_str(),
crate::content_purchase_protocol::PREPARE_OFFER_ROUTE
| crate::content_purchase_protocol::OFFER_ROUTE
| crate::content_purchase_protocol::ACCEPT_ROUTE
| crate::content_purchase_protocol::SETTLE_ROUTE
| crate::content_purchase_protocol::STATUS_ROUTE
| crate::content_purchase_protocol::CANCEL_ROUTE
)
{
return self.handle_purchase_request(req).await;
}
if method == Method::POST
&& path.starts_with("/content/registered_")
&& path.contains("/rental/")
&& (path.ends_with("/prepare") || path.ends_with("/start"))
{
return self.handle_rental_control(req).await;
}
// Convert body to bytes for non-WS routes
let headers = req.headers().clone();
let query_string = req.uri().query().map(|s| s.to_string()).unwrap_or_default();
@@ -649,21 +666,31 @@ impl ApiHandler {
// falls outside `connect-src`). Session-authenticated so only
// the logged-in node owner can spin up fetches.
(Method::GET, "/api/node-app-catalog") => {
if !self.is_authenticated(&headers).await { return Ok(Self::unauthorized()); }
if !self.is_authenticated(&headers).await {
return Ok(Self::unauthorized());
}
let data_dir = self.config.data_dir.clone();
let result = tokio::task::spawn_blocking(move || {
crate::container::node_catalog::verified_body(&data_dir)
}).await.unwrap_or_else(|error| Err(anyhow::anyhow!(error)));
})
.await
.unwrap_or_else(|error| Err(anyhow::anyhow!(error)));
let (status, body) = match result {
Ok(Some(body)) => (StatusCode::OK, body),
Ok(None) => (StatusCode::NOT_FOUND, "{}".to_owned()),
Err(error) => {
tracing::warn!("Node demo catalog rejected: {error}");
(StatusCode::CONFLICT, "{\"error\":\"Node demo catalog is unavailable\"}".to_owned())
},
(
StatusCode::CONFLICT,
"{\"error\":\"Node demo catalog is unavailable\"}".to_owned(),
)
}
};
Ok(Response::builder().status(status).header("Content-Type", "application/json")
.header("Cache-Control", "private, no-store").body(hyper::Body::from(body))?)
Ok(Response::builder()
.status(status)
.header("Content-Type", "application/json")
.header("Cache-Control", "private, no-store")
.body(hyper::Body::from(body))?)
}
(Method::GET, "/api/app-catalog") => {
+28 -1
View File
@@ -23,7 +23,8 @@ impl ApiHandler {
request.method() == Method::POST
&& matches!(
path.as_str(),
protocol::OFFER_ROUTE
protocol::PREPARE_OFFER_ROUTE
| protocol::OFFER_ROUTE
| protocol::ACCEPT_ROUTE
| protocol::SETTLE_ROUTE
| protocol::STATUS_ROUTE
@@ -59,6 +60,32 @@ impl ApiHandler {
)?;
let data_dir = &self.config.data_dir;
let result = match path.as_str() {
protocol::PREPARE_OFFER_ROUTE => {
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct Prepare {
content_id: String,
#[serde(default)]
retry: bool,
}
let input: Prepare = serde_json::from_slice(&bytes)?;
let identity = std::sync::Arc::new(
NodeIdentity::load_existing(&data_dir.join("identity")).await?,
);
anyhow::ensure!(identity.did_key()? == audience, "Node identity changed");
let root = data_dir.clone();
serde_json::to_value(
tokio::task::spawn_blocking(move || {
crate::registered_media::prepare_registered(
root,
identity,
&input.content_id,
input.retry,
)
})
.await??,
)?
}
protocol::OFFER_ROUTE => {
let body: OfferRequest = serde_json::from_slice(&bytes)?;
anyhow::ensure!(
@@ -4,7 +4,6 @@ use crate::{content_server::ByteRange, identity::NodeIdentity, registered_media:
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
@@ -46,6 +45,125 @@ fn denied(message: &'static str) -> Response<Body> {
}
impl ApiHandler {
pub(super) async fn handle_rental_control(
&self,
mut request: hyper::Request<Body>,
) -> Result<Response<Body>> {
use hyper::body::HttpBody;
#[derive(serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct Control {
capability: String,
ready_id: Option<String>,
#[serde(default)]
retry: bool,
}
let path = request.uri().path().to_owned();
let (base, action) = path.rsplit_once('/').context("Invalid rental action")?;
anyhow::ensure!(
matches!(action, "prepare" | "start") && request.method() == hyper::Method::POST,
"Invalid rental action"
);
let (content, purchase) = route(base)?;
let bytes = tokio::time::timeout(std::time::Duration::from_secs(15), async {
let mut bytes = Vec::new();
while let Some(chunk) = request.body_mut().data().await {
let chunk = chunk?;
anyhow::ensure!(
bytes
.len()
.checked_add(chunk.len())
.is_some_and(|n| n <= 16 * 1024),
"Rental request too large"
);
bytes.extend_from_slice(&chunk);
}
Ok::<_, anyhow::Error>(bytes)
})
.await
.context("Rental request timed out")??;
let audience = crate::identity::did_key_from_pubkey_hex(&self.self_pubkey_hex)?;
let buyer = crate::content_auth::authenticate_request(
request.headers(),
&audience,
&hyper::Method::POST,
&path,
&bytes,
chrono::Utc::now().timestamp(),
)?;
let input: Control = serde_json::from_slice(&bytes)?;
let identity =
Arc::new(NodeIdentity::load_existing(&self.config.data_dir.join("identity")).await?);
anyhow::ensure!(identity.did_key()? == audience, "Node identity changed");
let result = if action == "prepare" {
anyhow::ensure!(input.ready_id.is_none(), "Prepare does not start a rental");
let prior = crate::registered_media::paid_window(
&self.config.data_dir,
&identity,
content,
purchase,
&buyer,
&input.capability,
)
.await?;
let metadata = crate::registered_media::registered_metadata(
&self.config.data_dir,
&identity,
content,
)?;
anyhow::ensure!(
prior
.as_ref()
.is_none_or(|window| clock() >= window.started_at),
"Rental clock moved backwards"
);
if let Some(window) = prior.as_ref().filter(|window| clock() >= window.expires_at) {
serde_json::json!({"state":"expired", "viewing_seconds":metadata.0.viewing_seconds,"started_at":window.started_at,"expires_at":window.expires_at})
} else {
let state = crate::registered_media::prepare_paid(
self.config.data_dir.clone(),
identity,
content.into(),
purchase.into(),
buyer,
input.capability,
input.retry,
)
.await?;
let mut result = serde_json::to_value(state)?;
result["viewing_seconds"] = serde_json::json!(metadata.0.viewing_seconds);
result["started_at"] =
serde_json::json!(prior.as_ref().map(|window| window.started_at));
result["expires_at"] =
serde_json::json!(prior.as_ref().map(|window| window.expires_at));
result
}
} else {
anyhow::ensure!(!input.retry, "Start cannot retry verification");
let ready_id = input
.ready_id
.context("Media must be ready before explicit Start")?;
let window = crate::registered_media::start_paid(
self.config.data_dir.clone(),
identity,
content.into(),
purchase.into(),
buyer,
input.capability,
ready_id,
)
.await?;
anyhow::ensure!(clock() >= window.started_at, "Rental clock moved backwards");
serde_json::json!({"state": if clock() >= window.expires_at {"expired"} else {"started"},
"started_at":window.started_at,"expires_at":window.expires_at})
};
Ok(build_response(
StatusCode::OK,
"application/json",
Body::from(serde_json::to_vec(&result)?),
))
}
pub(super) async fn handle_registered_rental(
&self,
path: &str,
@@ -102,7 +220,7 @@ impl ApiHandler {
let selected = content.to_owned();
let key = identity.clone();
let metadata = tokio::task::spawn_blocking(move || {
crate::registered_media::registered_terms(&data, &key, &selected)
crate::registered_media::registered_metadata(&data, &key, &selected)
})
.await?;
let (receipt, _) = match metadata {
@@ -160,11 +278,9 @@ async fn rental_response(
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 {
(opened.file, opened.verification, start, length, now),
move |(mut file, verification, position, left, now)| async move {
if left == 0 {
return Ok::<_, std::io::Error>(None);
}
@@ -175,15 +291,20 @@ async fn rental_response(
"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 read = tokio::task::spawn_blocking(move || {
let bytes = verification
.index
.read_slice(&mut file, position, left.min(64 * 1024) as usize)
.map_err(std::io::Error::other)?;
Ok::<_, std::io::Error>((file, verification, bytes))
});
let (file, verification, bytes) =
tokio::time::timeout(std::time::Duration::from_secs(expires - instant), read)
.await
.map_err(|_| {
std::io::Error::new(std::io::ErrorKind::TimedOut, "Rental window ended")
})?
.map_err(std::io::Error::other)??;
let instant = now();
if instant < started || instant >= expires {
return Err(std::io::Error::new(
@@ -191,14 +312,11 @@ async fn rental_response(
"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 count = bytes.len() as u64;
Ok(Some((
bytes,
(file, verification, position + count, left - count, now),
)))
},
);
let mut response = Response::builder()
@@ -257,10 +375,19 @@ mod tests {
);
}
fn opened(size: u64) -> OpenedMedia {
let file = tempfile::tempfile().unwrap();
use sha2::{Digest, Sha256};
let mut file = tempfile::tempfile().unwrap();
file.set_len(size).unwrap();
let binding = crate::rental_chunk_index::Binding {
content_id: format!("registered_{}", uuid::Uuid::new_v4()),
receipt_sha256: "ab".repeat(32),
full_sha256: hex::encode(Sha256::digest(vec![0; size as usize])),
size,
};
let index = crate::rental_chunk_index::Index::scan(&mut file, binding, |_| Ok(())).unwrap();
OpenedMedia {
file,
verification: crate::rental_readiness::Ready::fixture(index),
size_bytes: size,
mime_type: "video/mp4".into(),
started_at: 1000,
@@ -333,6 +333,8 @@ impl RpcHandler {
"content.indeehub-projects" => self.handle_content_indeehub_projects().await,
"content.browse-all-peers" => self.handle_content_browse_all_peers().await,
"content.playback-handle" => self.handle_playback_handle(params, session_token).await,
"content.playback-prepare" => self.handle_playback_prepare(params, session_token).await,
"content.playback-start" => self.handle_playback_start(params, session_token).await,
"content.playback-status" => self.handle_playback_status(params, session_token).await,
"content.rental-purchase" => self.handle_content_rental_purchase(params).await,
"content.purchase" => self.handle_content_purchase(params).await,
+5 -1
View File
@@ -20,8 +20,8 @@ mod interfaces;
pub(crate) mod lnd;
mod marketplace;
mod media_registration;
mod purchase;
mod playback;
mod purchase;
// pub(crate): 13-10's `assistant::backends::select_backend` reuses
// `mesh::assistant::detect_ollama()` (D-04) rather than re-probing —
// matches the existing `pub(crate) mod bitcoin_relay;`/`pub(crate) mod
@@ -113,6 +113,8 @@ fn native_consent_origin_allowed(method: &str, headers: &hyper::HeaderMap, dev_m
| "content.cancel-purchase"
| "content.playback-handle"
| "content.playback-status"
| "content.playback-prepare"
| "content.playback-start"
) || nostr_signing_origin_allowed(headers, dev_mode)
}
@@ -831,6 +833,8 @@ mod nostr_signing_origin_tests {
"content.cancel-purchase",
"content.playback-handle",
"content.playback-status",
"content.playback-prepare",
"content.playback-start",
] {
assert!(!native_consent_origin_allowed(
method,
+177 -3
View File
@@ -55,9 +55,7 @@ impl RpcHandler {
let expires = self
.playback_handles()
.expiry(&handle, &session, &context)?;
Ok(
serde_json::json!({"playback_url":format!("/api/rental-playback/{handle}"), "expires_at":expires}),
)
Ok(serde_json::json!({"handle":handle,"expires_at":expires}))
}
pub(super) async fn handle_playback_status(
&self,
@@ -72,3 +70,179 @@ impl RpcHandler {
Ok(serde_json::json!({"expires_at":expires}))
}
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct Prepare {
handle: String,
#[serde(default)]
retry: bool,
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct Start {
handle: String,
ready_id: String,
}
impl RpcHandler {
async fn playback_control(
&self,
handle: &str,
ready_id: Option<&str>,
retry: bool,
session: &Option<String>,
) -> Result<serde_json::Value> {
let (session, context) = self.playback_context(session).await?;
let binding = self.playback_handles().lookup(handle, &session, &context)?;
let (capability, duration) = {
let journal = crate::content_purchase::Journal::open(&self.config.data_dir).await?;
let record = journal
.buyer(&binding.contract.id)
.await?
.context("Original purchase missing")?;
anyhow::ensure!(
record.contract == binding.contract,
"Original purchase changed"
);
let capability = record
.receipt()
.context("Original payment is not settled")?
.capability
.clone();
let envelope = journal
.protocol_envelope("buyer", &binding.contract.id)
.await?
.context("Original rental terms missing")?;
(
capability,
envelope
.offer
.viewing_seconds
.context("Purchase is not a timed rental")?,
)
};
let peer = crate::federation::load_unique_payment_peer(
&self.config.data_dir,
&binding.seller_onion,
)
.await?;
anyhow::ensure!(
peer.did == binding.contract.seller_did,
"Purchased seller identity changed"
);
let mesh = peer.fips_npub.context("Seller mesh binding unavailable")?;
let action = if ready_id.is_some() {
"start"
} else {
"prepare"
};
let path = format!(
"/content/{}/rental/{}/{action}",
binding.contract.content_id, binding.contract.id
);
let body = serde_json::json!({"capability":capability,"ready_id":ready_id,"retry":retry});
let (mut response, _) =
crate::fips::dial::PeerRequest::new(Some(&mesh), &binding.seller_onion, &path)
.require_fips()
.single_delivery()
.timeout(std::time::Duration::from_secs(20))
.send_content_json(&self.config.data_dir, &binding.contract.seller_did, &body)
.await?;
anyhow::ensure!(
response.status().is_success(),
"Original rental is unavailable; retry this purchase without paying again"
);
let mut bytes = Vec::new();
while let Some(chunk) = response.chunk().await? {
anyhow::ensure!(
bytes
.len()
.checked_add(chunk.len())
.is_some_and(|n| n <= 16 * 1024),
"Rental control response too large"
);
bytes.extend_from_slice(&chunk);
}
let remote: serde_json::Value = serde_json::from_slice(&bytes)?;
let state = remote["state"].as_str().context("Missing rental state")?;
let window = match (remote["started_at"].as_u64(), remote["expires_at"].as_u64()) {
(None, None) if remote["started_at"].is_null() && remote["expires_at"].is_null() => {
None
}
(Some(started), Some(expires))
if started > 0 && started.checked_add(duration) == Some(expires) =>
{
Some((started, expires))
}
_ => anyhow::bail!("Seller changed the original rental window"),
};
let (started_at, expires_at) =
window.map_or((None, None), |(start, end)| (Some(start), Some(end)));
let result = match state {
"preparing" if ready_id.is_none() => {
let completed = remote["completed_bytes"]
.as_u64()
.context("Invalid verification progress")?;
anyhow::ensure!(
remote["total_bytes"].as_u64() == Some(binding.contract.content_size)
&& completed <= binding.contract.content_size
&& remote["viewing_seconds"].as_u64() == Some(duration),
"Rental preparation terms changed"
);
serde_json::json!({"state":"preparing","completed_bytes":completed,"total_bytes":binding.contract.content_size,
"viewing_seconds":duration,"started_at":started_at,"expires_at":expires_at})
}
"ready" if ready_id.is_none() => {
let id = remote["ready_id"]
.as_str()
.context("Missing readiness identifier")?;
let parsed = uuid::Uuid::parse_str(id)?;
anyhow::ensure!(
parsed.to_string() == id
&& parsed.get_version_num() == 4
&& remote["total_bytes"].as_u64() == Some(binding.contract.content_size)
&& remote["viewing_seconds"].as_u64() == Some(duration),
"Rental readiness changed"
);
serde_json::json!({"state":"ready","ready_id":id,"handle":handle,"viewing_seconds":duration,
"started_at":started_at,"expires_at":expires_at})
}
"unavailable" if ready_id.is_none() => {
serde_json::json!({"state":"unavailable","expires_at":expires_at})
}
"expired" if window.is_some() => {
serde_json::json!({"state":"expired","started_at":started_at,"expires_at":expires_at})
}
"started" if ready_id.is_some() && window.is_some() => {
serde_json::json!({"state":"started","started_at":started_at,
"expires_at":expires_at,"playback_url":format!("/api/rental-playback/{handle}")})
}
_ => {
anyhow::bail!("Unexpected rental state; recover this purchase without paying again")
}
};
if let Some((_, expires)) = window {
self.playback_handles()
.note_expiry(handle, &binding, expires)?;
}
Ok(result)
}
pub(super) async fn handle_playback_prepare(
&self,
params: Option<serde_json::Value>,
session: &Option<String>,
) -> Result<serde_json::Value> {
let input: Prepare = serde_json::from_value(params.context("Missing playback handle")?)?;
self.playback_control(&input.handle, None, input.retry, session)
.await
}
pub(super) async fn handle_playback_start(
&self,
params: Option<serde_json::Value>,
session: &Option<String>,
) -> Result<serde_json::Value> {
let input: Start = serde_json::from_value(params.context("Missing readiness identifier")?)?;
self.playback_control(&input.handle, Some(&input.ready_id), false, session)
.await
}
}
+13 -2
View File
@@ -51,6 +51,9 @@ impl RpcHandler {
)
.await?;
match result {
ReadyPurchase::Preparing { .. } => {
anyhow::bail!("Ordinary purchase cannot prepare a timed rental")
}
ReadyPurchase::AwaitingConfirmation {
operation_id,
envelope_sha256,
@@ -127,6 +130,8 @@ impl RpcHandler {
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct RentalParams {
#[serde(default)]
retry_preparation: bool,
seller_did: String,
content_id: String,
expected_sha256: String,
@@ -168,8 +173,9 @@ impl RpcHandler {
&params.seller_did,
)
.await?;
let transport =
FipsPurchaseTransport::load(self.config.data_dir.clone(), onion.clone()).await?;
let transport = FipsPurchaseTransport::load(self.config.data_dir.clone(), onion.clone())
.await?
.retry_preparation(params.retry_preparation);
let expected = caller::ExpectedRental {
seller_did: params.seller_did,
content_id: params.content_id.clone(),
@@ -189,6 +195,11 @@ impl RpcHandler {
)
.await?
{
ReadyPurchase::Preparing {
completed_bytes,
total_bytes,
} => Ok(serde_json::json!({
"state":"preparing", "completed_bytes":completed_bytes,"total_bytes":total_bytes})),
ReadyPurchase::AwaitingConfirmation {
operation_id,
envelope_sha256,
@@ -12,6 +12,12 @@ pub(crate) trait PurchaseTransport: Send + Sync {
/// This identity is the independently verified peer binding, not response JSON.
fn seller_did(&self) -> &str;
fn seller_onion(&self) -> &str;
fn prepare_offer(
&self,
_content_id: &str,
) -> impl Future<Output = Result<Option<(u64, u64)>>> + Send {
async { Ok(None) }
}
fn offer(&self, id: &str, content_id: &str) -> impl Future<Output = Result<Offer>> + Send;
fn accept(&self, envelope: &Envelope) -> impl Future<Output = Result<Accepted>> + Send;
fn status(&self, envelope: &Envelope) -> impl Future<Output = Result<SellerStatus>> + Send;
@@ -20,6 +26,10 @@ pub(crate) trait PurchaseTransport: Send + Sync {
}
#[derive(Clone)]
pub(crate) enum ReadyPurchase {
Preparing {
completed_bytes: u64,
total_bytes: u64,
},
AwaitingConfirmation {
operation_id: String,
envelope_sha256: String,
@@ -210,6 +220,16 @@ pub(crate) async fn purchase_bound(
!unresolved_owned,
"Prior purchase delivery needs recovery; no new payment operation was created"
);
if content_id.starts_with("registered_") {
if let Some((completed_bytes, total_bytes)) =
transport.prepare_offer(content_id).await?
{
return Ok(ReadyPurchase::Preparing {
completed_bytes,
total_bytes,
});
}
}
let id = uuid::Uuid::new_v4().to_string();
let offer = transport.offer(&id, content_id).await?;
offer.validate()?;
@@ -8,6 +8,7 @@ use anyhow::{Context, Result};
use serde::{Deserialize, Serialize};
use std::path::Path;
pub(crate) const PREPARE_OFFER_ROUTE: &str = "/content/purchase/v1/prepare-offer";
pub(crate) const OFFER_ROUTE: &str = "/content/purchase/v1/offer";
pub(crate) const ACCEPT_ROUTE: &str = "/content/purchase/v1/accept";
pub(crate) const SETTLE_ROUTE: &str = "/content/purchase/v1/settle";
@@ -15,6 +15,7 @@ pub(crate) struct FipsPurchaseTransport {
onion: String,
seller_did: String,
fips_npub: String,
retry_preparation: bool,
}
impl FipsPurchaseTransport {
pub async fn load(data_dir: PathBuf, onion: String) -> Result<Self> {
@@ -31,8 +32,13 @@ impl FipsPurchaseTransport {
onion,
seller_did: peer.did,
fips_npub,
retry_preparation: false,
})
}
pub fn retry_preparation(mut self, retry: bool) -> Self {
self.retry_preparation = retry;
self
}
async fn post<B: Serialize + Sync, R: DeserializeOwned>(
&self,
route: &str,
@@ -73,6 +79,30 @@ impl PurchaseTransport for FipsPurchaseTransport {
fn seller_did(&self) -> &str {
&self.seller_did
}
async fn prepare_offer(&self, content_id: &str) -> Result<Option<(u64, u64)>> {
let state: crate::rental_readiness::Status = self
.post(
protocol::PREPARE_OFFER_ROUTE,
&serde_json::json!({"content_id":content_id,"retry":self.retry_preparation}),
)
.await?;
match state {
crate::rental_readiness::Status::Ready { .. } => Ok(None),
crate::rental_readiness::Status::Preparing {
completed_bytes,
total_bytes,
} => {
anyhow::ensure!(
completed_bytes <= total_bytes && total_bytes > 0,
"Invalid media preparation progress"
);
Ok(Some((completed_bytes, total_bytes)))
}
crate::rental_readiness::Status::Unavailable => anyhow::bail!(
"Media verification failed; no payment started. Retry media preparation."
),
}
}
async fn offer(&self, id: &str, content_id: &str) -> Result<Offer> {
#[derive(Serialize)]
struct Request<'a> {
+2
View File
@@ -83,6 +83,8 @@ mod port_allocator;
mod prepared_media;
mod rate_limit;
mod registered_media;
mod rental_chunk_index;
mod rental_readiness;
pub mod seed;
mod server;
mod session;
+498 -35
View File
@@ -86,6 +86,7 @@ struct Lease {
/// stream at expiry, and must not call open_paid for HTTP HEAD/preflight.
pub(crate) struct OpenedMedia {
pub file: File,
pub verification: Arc<crate::rental_readiness::Ready>,
pub size_bytes: u64,
pub mime_type: String,
pub started_at: u64,
@@ -395,6 +396,13 @@ fn verify(
Ok(())
}
fn open_snapshot(data_dir: &Path, record: &Registered) -> Result<File> {
open_snapshot_with_stamp(data_dir, record, true)
}
fn open_snapshot_with_stamp(
data_dir: &Path,
record: &Registered,
require_original_stamp: bool,
) -> Result<File> {
use std::ffi::CString;
use std::os::fd::FromRawFd;
// Resolve from the configured data directory in one kernel operation. No
@@ -435,7 +443,8 @@ fn open_snapshot(data_dir: &Path, record: &Registered) -> Result<File> {
);
let file = unsafe { File::from_raw_fd(fd as i32) };
anyhow::ensure!(
Stamp::from_file(&file)? == record.stamp,
Stamp::from_file(&file)?.size == record.stamp.size
&& (!require_original_stamp || Stamp::from_file(&file)? == record.stamp),
"Registered immutable snapshot changed"
);
Ok(file)
@@ -561,11 +570,11 @@ pub(crate) fn resolve_registration(
/// No request chooses app scope or storage path. This is an offer prerequisite,
/// not advertisement: the future offer creator must authenticate the peer and
/// bind all returned terms into the purchase contract before seller acceptance.
pub(crate) fn registered_terms(
fn metadata_record(
data_dir: &Path,
identity: &NodeIdentity,
content_id: &str,
) -> Result<(Receipt, String)> {
) -> Result<Registered> {
let id = registration_id(content_id)?;
let pin = registration_pin::load_existing(data_dir, APP_ID, identity)?;
let held = held(data_dir, false)?;
@@ -573,32 +582,198 @@ pub(crate) fn registered_terms(
.context("Registered content is unavailable")?;
drop(held);
verify(&record, &pin, identity)?;
let mut file = open_snapshot(data_dir, &record)?;
// A quote must not invite payment for altered bytes, including a same-tick
// metadata collision. This scan completes before any offer is accepted.
ensure_verified_for_use(data_dir, &record, &mut file, true)?;
Ok(record)
}
fn chunk_binding(record: &Registered) -> Result<crate::rental_chunk_index::Binding> {
Ok(crate::rental_chunk_index::Binding {
content_id: record.receipt.content_id.clone(),
receipt_sha256: hash(&serde_json::to_vec(&record.receipt)?),
full_sha256: record.receipt.sha256.clone(),
size: record.stamp.size,
})
}
fn readiness_key(data_dir: &Path, record: &Registered) -> Result<String> {
Ok(format!(
"{}:{}",
data_dir.canonicalize()?.display(),
hash(&serde_json::to_vec(&record.receipt)?)
))
}
/// Signed metadata only: no whole-file scan or lease is permitted here.
pub(crate) fn registered_metadata(
data_dir: &Path,
identity: &NodeIdentity,
content_id: &str,
) -> Result<(Receipt, String)> {
let record = metadata_record(data_dir, identity, content_id)?;
Ok((record.receipt, record.terms_sha256))
}
/// Starts/reuses bounded background verification, before quoting or playing.
/// The signed index survives restart; range requests never invoke this scan.
pub(crate) fn prepare_registered(
data_dir: PathBuf,
identity: Arc<NodeIdentity>,
content_id: &str,
retry: bool,
) -> Result<crate::rental_readiness::Status> {
let record = metadata_record(&data_dir, &identity, content_id)?;
let binding = chunk_binding(&record)?;
let key = readiness_key(&data_dir, &record)?;
let manager = crate::rental_readiness::shared();
if retry {
manager.retry_failed(&key, &binding)?;
}
manager.prepare(key, binding.clone(), move |progress| {
let held = keyed(&data_dir, "verify", &record.receipt.request_id)?;
let path = held
.path
.join(format!("chunks-{}.bin", record.receipt.request_id));
let mut file = open_snapshot_with_stamp(&data_dir, &record, false)?;
if let Ok(Some(index)) = crate::rental_chunk_index::Index::load(&path, &binding, &identity)
{
return Ok(index);
}
// A damaged derived index is not evidence about the original bytes.
// Rebuild only after a complete scan matches the signed producer SHA;
// this avoids making an editable cache a permanent recovery dead end.
let index = crate::rental_chunk_index::Index::scan(&mut file, binding, |n| {
progress.store(n, std::sync::atomic::Ordering::SeqCst);
Ok(())
})?;
index.save(&path, &identity)?;
Ok(index)
})
}
#[derive(Clone, Serialize)]
pub(crate) struct RentalWindow {
pub started_at: u64,
pub expires_at: u64,
}
fn verify_lease(
lease: &Lease,
contract: &Contract,
capability: &str,
duration: u64,
) -> Result<RentalWindow> {
anyhow::ensure!(
lease.version == 1
&& lease.purchase_id == contract.id
&& lease.buyer_did == contract.buyer_did
&& lease.content_id == contract.content_id
&& lease.contract_hash == contract.context_hash()?
&& lease.capability_hash == hash(capability.as_bytes())
&& lease.started_at.checked_add(duration) == Some(lease.expires_at),
"Persisted rental terms changed"
);
Ok(RentalWindow {
started_at: lease.started_at,
expires_at: lease.expires_at,
})
}
fn verify_settled_terms(record: &Registered, contract: &Contract) -> Result<()> {
anyhow::ensure!(
contract.content_sha256 == record.receipt.sha256
&& contract.content_size == record.stamp.size
&& contract.terms_sha256 == record.terms_sha256
&& record.receipt.price_sats > 0
&& record
.receipt
.payment_methods
.iter()
.any(|method| method == "cashu")
&& contract.minimum_net_sats == record.receipt.price_sats,
"Settled purchase does not match registered immutable terms"
);
Ok(())
}
/// Only explicit authenticated Start invokes this commit. Losing its response
/// cannot reset the original clock; repeated Start returns the same window.
fn start_ready(
data_dir: &Path,
identity: &NodeIdentity,
contract: &Contract,
capability: &str,
ready_id: &str,
now: impl Fn() -> Result<u64>,
) -> Result<RentalWindow> {
let record = metadata_record(data_dir, identity, &contract.content_id)?;
verify_settled_terms(&record, contract)?;
{
let held = keyed(data_dir, "lease", &contract.id)?;
if let Some(lease) = read::<Lease>(&held.path.join(format!("lease-{}.json", contract.id)))?
{
return verify_lease(&lease, contract, capability, record.receipt.viewing_seconds);
}
}
let ready = crate::rental_readiness::shared().ready(
&readiness_key(data_dir, &record)?,
&chunk_binding(&record)?,
Some(ready_id),
)?;
let mut file = open_snapshot_with_stamp(data_dir, &record, false)?;
// A bounded actual byte read before the first clock. Later chunks are
// verified immediately before delivery, including seek/range requests.
ready.index.read_slice(&mut file, 0, 1)?;
let held = keyed(data_dir, "lease", &contract.id)?;
let path = held.path.join(format!("lease-{}.json", contract.id));
if let Some(lease) = read::<Lease>(&path)? {
return verify_lease(&lease, contract, capability, record.receipt.viewing_seconds);
}
let started_at = now()?;
let expires_at = started_at
.checked_add(record.receipt.viewing_seconds)
.context("Rental expiry overflow")?;
let lease = Lease {
version: 1,
purchase_id: contract.id.clone(),
buyer_did: contract.buyer_did.clone(),
content_id: contract.content_id.clone(),
contract_hash: contract.context_hash()?,
capability_hash: hash(capability.as_bytes()),
started_at,
expires_at,
};
persist(&held, &path, &lease)?;
Ok(RentalWindow {
started_at,
expires_at,
})
}
pub(crate) fn registered_terms(
data_dir: &Path,
identity: &NodeIdentity,
content_id: &str,
) -> Result<(Receipt, String)> {
let record = metadata_record(data_dir, identity, content_id)?;
let _ready = crate::rental_readiness::shared().ready(
&readiness_key(data_dir, &record)?,
&chunk_binding(&record)?,
None,
)?;
let _file = open_snapshot_with_stamp(data_dir, &record, false)?;
Ok((record.receipt, record.terms_sha256))
}
/// Only authenticated peer GET/range routes may call this. The capability is
/// checked against this node's durable seller journal; client receipts do not
/// establish payment. A lease is persisted before any bytes can be returned.
pub(crate) async fn open_paid(
data_dir: PathBuf,
identity: Arc<NodeIdentity>,
content_id: String,
purchase_id: String,
authenticated_buyer: String,
capability: String,
) -> Result<OpenedMedia> {
async fn paid_contract(
data_dir: &Path,
identity: &NodeIdentity,
content_id: &str,
purchase_id: &str,
authenticated_buyer: &str,
capability: &str,
) -> Result<Contract> {
uuid(&purchase_id)?;
registration_id(&content_id)?;
anyhow::ensure!(
capability.len() == 64 && capability.bytes().all(|c| c.is_ascii_hexdigit()),
"Invalid delivery capability"
);
let (contract, saved_capability) = {
let journal = Journal::open(&data_dir).await?;
let contract = {
let journal = Journal::open(data_dir).await?;
let record = journal
.seller(&purchase_id)
.await?
@@ -614,18 +789,157 @@ pub(crate) async fn open_paid(
_ => anyhow::bail!("Seller receipt is not durable"),
};
anyhow::ensure!(
constant_equal(&capability, &receipt.capability),
constant_equal(capability, &receipt.capability),
"Delivery capability does not match this purchase"
);
(record.contract, receipt.capability)
record.contract
};
Ok(contract)
}
pub(crate) async fn prepare_paid(
data_dir: PathBuf,
identity: Arc<NodeIdentity>,
content_id: String,
purchase_id: String,
buyer: String,
capability: String,
retry: bool,
) -> Result<crate::rental_readiness::Status> {
paid_contract(
&data_dir,
&identity,
&content_id,
&purchase_id,
&buyer,
&capability,
)
.await?;
tokio::task::spawn_blocking(move || prepare_registered(data_dir, identity, &content_id, retry))
.await?
}
pub(crate) async fn paid_window(
data_dir: &Path,
identity: &NodeIdentity,
content_id: &str,
purchase_id: &str,
buyer: &str,
capability: &str,
) -> Result<Option<RentalWindow>> {
let contract = paid_contract(
data_dir,
identity,
content_id,
purchase_id,
buyer,
capability,
)
.await?;
let record = metadata_record(data_dir, identity, content_id)?;
verify_settled_terms(&record, &contract)?;
let held = keyed(data_dir, "lease", purchase_id)?;
read::<Lease>(&held.path.join(format!("lease-{purchase_id}.json")))?
.map(|lease| {
verify_lease(
&lease,
&contract,
capability,
record.receipt.viewing_seconds,
)
})
.transpose()
}
pub(crate) async fn start_paid(
data_dir: PathBuf,
identity: Arc<NodeIdentity>,
content_id: String,
purchase_id: String,
buyer: String,
capability: String,
ready_id: String,
) -> Result<RentalWindow> {
let contract = paid_contract(
&data_dir,
&identity,
&content_id,
&purchase_id,
&buyer,
&capability,
)
.await?;
tokio::task::spawn_blocking(move || {
open_settled(&data_dir, &identity, &contract, &saved_capability, || {
start_ready(
&data_dir,
&identity,
&contract,
&capability,
&ready_id,
|| u64::try_from(chrono::Utc::now().timestamp()).context("Invalid node clock"),
)
})
.await?
}
pub(crate) async fn open_paid(
data_dir: PathBuf,
identity: Arc<NodeIdentity>,
content_id: String,
purchase_id: String,
authenticated_buyer: String,
capability: String,
) -> Result<OpenedMedia> {
let contract = paid_contract(
&data_dir,
&identity,
&content_id,
&purchase_id,
&authenticated_buyer,
&capability,
)
.await?;
tokio::task::spawn_blocking(move || {
open_started(&data_dir, &identity, &contract, &capability, || {
u64::try_from(chrono::Utc::now().timestamp()).context("Invalid node clock")
})
})
.await?
}
fn open_started(
data_dir: &Path,
identity: &NodeIdentity,
contract: &Contract,
capability: &str,
now: impl Fn() -> Result<u64>,
) -> Result<OpenedMedia> {
let record = metadata_record(data_dir, identity, &contract.content_id)?;
verify_settled_terms(&record, contract)?;
let lease = {
let held = keyed(data_dir, "lease", &contract.id)?;
read::<Lease>(&held.path.join(format!("lease-{}.json", contract.id)))?
.context("Explicit Start is required; no rental was started")?
};
let window = verify_lease(&lease, contract, capability, record.receipt.viewing_seconds)?;
let instant = now()?;
anyhow::ensure!(
instant >= window.started_at && instant < window.expires_at,
"Rental expired or node clock moved backwards"
);
let verification = crate::rental_readiness::shared().ready(
&readiness_key(data_dir, &record)?,
&chunk_binding(&record)?,
None,
)?;
let file = open_snapshot_with_stamp(data_dir, &record, false)?;
Ok(OpenedMedia {
file,
verification,
size_bytes: record.stamp.size,
mime_type: record.mime_type,
started_at: window.started_at,
expires_at: window.expires_at,
})
}
// Existing clock/receipt invariant fixtures exercise an explicit local start;
// public prepare/start/open separation is tested independently below.
#[cfg(test)]
fn open_settled(
data_dir: &Path,
identity: &NodeIdentity,
@@ -701,8 +1015,12 @@ fn open_settled(
now >= lease.started_at && now < lease.expires_at,
"Rental expired or node clock moved backwards"
);
let index =
crate::rental_chunk_index::Index::scan(&mut file, chunk_binding(&record)?, |_| Ok(()))?;
let verification = crate::rental_readiness::Ready::fixture(index);
Ok(OpenedMedia {
file,
verification,
size_bytes: record.stamp.size,
mime_type: record.mime_type,
started_at: lease.started_at,
@@ -786,6 +1104,30 @@ mod tests {
|_| Ok(()),
)
}
async fn ready(&self, receipt: &Receipt) -> String {
tokio::time::timeout(std::time::Duration::from_secs(5), async {
loop {
match prepare_registered(
self.root.path().into(),
self.identity.clone(),
&receipt.content_id,
false,
)
.unwrap()
{
crate::rental_readiness::Status::Ready { ready_id, .. } => return ready_id,
crate::rental_readiness::Status::Preparing { .. } => {
tokio::time::sleep(std::time::Duration::from_millis(5)).await
}
crate::rental_readiness::Status::Unavailable => {
panic!("fixture verification failed")
}
}
}
})
.await
.expect("bounded fixture preparation")
}
fn contract(&self, receipt: &Receipt) -> Contract {
Contract {
version: 1,
@@ -872,6 +1214,7 @@ mod tests {
std::fs::remove_file(fixture.root.path().join("cloud/film.mp4")).unwrap();
assert_eq!(fixture.register(2000).unwrap(), receipt);
assert_eq!(std::fs::read(&mapping).unwrap(), original);
fixture.ready(&receipt).await;
let found =
registered_terms(fixture.root.path(), &fixture.identity, &receipt.content_id).unwrap();
assert_eq!(found, (receipt.clone(), terms(&receipt).unwrap()));
@@ -1029,6 +1372,57 @@ mod tests {
.join(STORE)
.join(format!("lease-{}.json", contract.id))
.exists());
let ready_id = fixture.ready(&receipt).await;
// Preparation, settlement and a GET are not permission to start the clock.
assert!(open_paid(
fixture.root.path().into(),
fixture.identity.clone(),
receipt.content_id.clone(),
contract.id.clone(),
contract.buyer_did.clone(),
settled.capability.clone()
)
.await
.is_err());
assert!(paid_window(
fixture.root.path(),
&fixture.identity,
&receipt.content_id,
&contract.id,
&contract.buyer_did,
&settled.capability
)
.await
.unwrap()
.is_none());
let window = start_paid(
fixture.root.path().into(),
fixture.identity.clone(),
receipt.content_id.clone(),
contract.id.clone(),
contract.buyer_did.clone(),
settled.capability.clone(),
ready_id,
)
.await
.unwrap();
// Lost Start response: a retry replays the durable original window even
// if its ephemeral readiness identifier no longer exists after restart.
let replay = start_paid(
fixture.root.path().into(),
fixture.identity.clone(),
receipt.content_id.clone(),
contract.id.clone(),
contract.buyer_did.clone(),
settled.capability.clone(),
uuid::Uuid::new_v4().to_string(),
)
.await
.unwrap();
assert_eq!(
(window.started_at, window.expires_at),
(replay.started_at, replay.expires_at)
);
let opened = open_paid(
fixture.root.path().into(),
fixture.identity.clone(),
@@ -1041,6 +1435,33 @@ mod tests {
.unwrap();
assert_eq!(opened.expires_at - opened.started_at, 60);
assert_eq!(opened.size_bytes, 24);
for offset in [0, 8, 23] {
let mut range = open_paid(
fixture.root.path().into(),
fixture.identity.clone(),
receipt.content_id.clone(),
contract.id.clone(),
contract.buyer_did.clone(),
settled.capability.clone(),
)
.await
.unwrap();
assert_eq!(
range
.verification
.index
.read_slice(&mut range.file, offset, 1)
.unwrap()
.len(),
1
);
assert_eq!(range.started_at, window.started_at);
}
assert_eq!(
crate::rental_chunk_index::scans(&receipt.content_id),
1,
"quote preparation and repeated GET/ranges share one original SHA scan"
);
}
#[tokio::test]
async fn altered_snapshot_or_overflow_never_creates_first_rental() {
@@ -1249,21 +1670,69 @@ mod tests {
.exists());
}
#[tokio::test]
async fn offer_preflight_rebuilds_verified_cache_without_creating_rental_and_rejects_corruption(
) {
async fn signed_index_preparation_is_lease_free_and_unsigned_cache_cannot_change_it() {
let fixture = Fixture::new().await;
let receipt = fixture.register(1100).unwrap();
let path = fixture
assert!(
registered_terms(fixture.root.path(), &fixture.identity, &receipt.content_id).is_err()
);
let ready_id = fixture.ready(&receipt).await;
let record =
metadata_record(fixture.root.path(), &fixture.identity, &receipt.content_id).unwrap();
let index_path = fixture
.root
.path()
.join(STORE)
.join(format!("chunks-{}.bin", receipt.request_id));
let index = crate::rental_chunk_index::Index::load(
&index_path,
&chunk_binding(&record).unwrap(),
&fixture.identity,
)
.unwrap()
.unwrap();
assert_eq!(index.binding.full_sha256, receipt.sha256);
let cache = fixture
.root
.path()
.join(STORE)
.join(format!("verified-{}.json", receipt.request_id));
let original = std::fs::read(&path).unwrap();
std::fs::remove_file(&path).unwrap();
let (terms, _) =
registered_terms(fixture.root.path(), &fixture.identity, &receipt.content_id).unwrap();
assert_eq!(terms, receipt);
assert_eq!(std::fs::read(&path).unwrap(), original);
std::fs::write(cache, b"untrusted corrupt old metadata").unwrap();
assert_eq!(
registered_terms(fixture.root.path(), &fixture.identity, &receipt.content_id)
.unwrap()
.0,
receipt
);
assert_eq!(fixture.ready(&receipt).await, ready_id);
assert!(std::fs::read_dir(fixture.root.path().join(STORE))
.unwrap()
.all(|entry| !entry
.unwrap()
.file_name()
.to_string_lossy()
.starts_with("lease-")));
let mut bytes = std::fs::read(&index_path).unwrap();
*bytes.last_mut().unwrap() ^= 1;
std::fs::write(&index_path, bytes).unwrap();
assert!(crate::rental_chunk_index::Index::load(
&index_path,
&chunk_binding(&record).unwrap(),
&fixture.identity
)
.is_err());
crate::rental_readiness::shared()
.forget_for_restart(&readiness_key(fixture.root.path(), &record).unwrap());
let replacement = fixture.ready(&receipt).await;
assert_ne!(replacement, ready_id);
assert_eq!(crate::rental_chunk_index::scans(&receipt.content_id), 2);
assert!(crate::rental_chunk_index::Index::load(
&index_path,
&chunk_binding(&record).unwrap(),
&fixture.identity
)
.unwrap()
.is_some());
assert!(std::fs::read_dir(fixture.root.path().join(STORE))
.unwrap()
.all(|entry| !entry
@@ -1271,11 +1740,5 @@ mod tests {
.file_name()
.to_string_lossy()
.starts_with("lease-")));
let mut cache: VerifiedSnapshot = read(&path).unwrap().unwrap();
cache.sha256 = "ff".repeat(32);
std::fs::write(&path, serde_json::to_vec(&cache).unwrap()).unwrap();
assert!(
registered_terms(fixture.root.path(), &fixture.identity, &receipt.content_id).is_err()
);
}
}
+319
View File
@@ -0,0 +1,319 @@
//! Node-signed byte indexes, created only by a scan matching the producer's
//! original signed full SHA. Metadata alone never authorizes streamed bytes.
use crate::identity::NodeIdentity;
use anyhow::{Context, Result};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use std::{
fs::{File, OpenOptions},
io::{Read, Seek, SeekFrom, Write},
os::unix::fs::OpenOptionsExt,
path::Path,
};
pub(crate) const CHUNK_BYTES: usize = 64 * 1024;
const MAX_BYTES: u64 = 16 * 1024 * 1024 * 1024;
#[derive(Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(deny_unknown_fields)]
pub(crate) struct Binding {
pub content_id: String,
pub receipt_sha256: String,
pub full_sha256: String,
pub size: u64,
}
impl Binding {
fn validate(&self) -> Result<()> {
anyhow::ensure!(
self.content_id
.strip_prefix("registered_")
.and_then(|id| uuid::Uuid::parse_str(id).ok())
.map(|id| format!("registered_{id}") == self.content_id)
.unwrap_or(false)
&& self.size > 0
&& self.size <= MAX_BYTES,
"Invalid byte index binding"
);
for value in [&self.receipt_sha256, &self.full_sha256] {
anyhow::ensure!(
value.len() == 64
&& value
.bytes()
.all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b)),
"Invalid byte index digest"
);
}
Ok(())
}
pub fn index_bytes(&self) -> Result<u64> {
self.validate()?;
Ok(self.size.div_ceil(CHUNK_BYTES as u64) * 32)
}
}
#[derive(Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct Header {
version: u32,
binding: Binding,
chunks_sha256: String,
signature: String,
}
fn preimage(binding: &Binding, chunks_sha256: &str) -> Result<Vec<u8>> {
Ok(serde_json::to_vec(&(
"archipelago-rental-chunk-index-v1",
CHUNK_BYTES,
&binding.content_id,
&binding.receipt_sha256,
&binding.full_sha256,
binding.size,
chunks_sha256,
))?)
}
pub(crate) struct Index {
pub binding: Binding,
chunks: Vec<[u8; 32]>,
pub commitment: String,
}
#[cfg(test)]
fn scan_counts() -> &'static std::sync::Mutex<std::collections::HashMap<String, usize>> {
static COUNTS: std::sync::OnceLock<std::sync::Mutex<std::collections::HashMap<String, usize>>> =
std::sync::OnceLock::new();
COUNTS.get_or_init(Default::default)
}
#[cfg(test)]
pub(crate) fn scans(content_id: &str) -> usize {
*scan_counts().lock().unwrap().get(content_id).unwrap_or(&0)
}
impl Index {
/// Exactly one whole-file scan. Caller runs it as a bounded background job;
/// progress/cancellation cannot create or alter a purchase lease.
pub fn scan(
file: &mut File,
binding: Binding,
mut progress: impl FnMut(u64) -> Result<()>,
) -> Result<Self> {
let bytes = binding.index_bytes()?;
#[cfg(test)]
{
*scan_counts()
.lock()
.unwrap()
.entry(binding.content_id.clone())
.or_default() += 1;
}
anyhow::ensure!(
file.metadata()?.is_file() && file.metadata()?.len() == binding.size,
"Snapshot size changed"
);
file.seek(SeekFrom::Start(0))?;
let mut chunks = Vec::with_capacity(bytes as usize / 32);
let mut full = Sha256::new();
let mut buffer = [0u8; CHUNK_BYTES];
let mut read = 0;
while read < binding.size {
progress(read)?;
let count = (binding.size - read).min(CHUNK_BYTES as u64) as usize;
file.read_exact(&mut buffer[..count])?;
full.update(&buffer[..count]);
chunks.push(Sha256::digest(&buffer[..count]).into());
read += count as u64;
}
anyhow::ensure!(
file.metadata()?.len() == binding.size
&& hex::encode(full.finalize()) == binding.full_sha256,
"Registered snapshot content changed"
);
progress(read)?;
let mut digest = Sha256::new();
for chunk in &chunks {
digest.update(chunk);
}
let commitment = hex::encode(digest.finalize());
file.seek(SeekFrom::Start(0))?;
Ok(Self {
binding,
chunks,
commitment,
})
}
/// Persist under the registration's existing verification flock. Sync rename
/// is the commit point; no mutation is queued after a canceled future drops.
pub fn save(&self, path: &Path, identity: &NodeIdentity) -> Result<()> {
let header = Header {
version: 1,
binding: self.binding.clone(),
chunks_sha256: self.commitment.clone(),
signature: identity.sign(&preimage(&self.binding, &self.commitment)?),
};
let encoded = serde_json::to_vec(&header)?;
anyhow::ensure!(encoded.len() <= 4096, "Byte index header too large");
let parent = path.parent().context("Byte index directory missing")?;
let temporary = parent.join(format!(".chunk-index-{}.tmp", uuid::Uuid::new_v4()));
let result = (|| {
let mut file = OpenOptions::new()
.create_new(true)
.write(true)
.mode(0o600)
.custom_flags(libc::O_NOFOLLOW | libc::O_CLOEXEC)
.open(&temporary)?;
file.write_all(&(encoded.len() as u32).to_be_bytes())?;
file.write_all(&encoded)?;
for chunk in &self.chunks {
file.write_all(chunk)?;
}
file.sync_all()?;
std::fs::rename(&temporary, path)?;
File::open(parent)?.sync_all()?;
Ok(())
})();
if result.is_err() {
let _ = std::fs::remove_file(temporary);
}
result
}
/// A signed index survives restart without treating an editable metadata
/// cache as proof. Every delivered chunk is independently checked below.
pub fn load(path: &Path, binding: &Binding, identity: &NodeIdentity) -> Result<Option<Self>> {
let expected = binding.index_bytes()?;
let mut file = match OpenOptions::new()
.read(true)
.custom_flags(libc::O_NOFOLLOW | libc::O_CLOEXEC | libc::O_NONBLOCK)
.open(path)
{
Ok(file) => file,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
Err(error) => return Err(error.into()),
};
let metadata = file.metadata()?;
anyhow::ensure!(
metadata.is_file() && metadata.len() <= expected + 4100,
"Invalid byte index file"
);
let mut len = [0u8; 4];
file.read_exact(&mut len)?;
let len = u32::from_be_bytes(len) as usize;
anyhow::ensure!(
len <= 4096 && metadata.len() == 4 + len as u64 + expected,
"Invalid byte index length"
);
let mut header = vec![0; len];
file.read_exact(&mut header)?;
let header: Header = serde_json::from_slice(&header)?;
anyhow::ensure!(
header.version == 1
&& &header.binding == binding
&& NodeIdentity::verify(
&identity.pubkey_hex(),
&preimage(binding, &header.chunks_sha256)?,
&header.signature
)?,
"Byte index signature or immutable binding changed"
);
let mut chunks = Vec::with_capacity(expected as usize / 32);
let mut digest = Sha256::new();
for _ in 0..expected / 32 {
let mut chunk = [0; 32];
file.read_exact(&mut chunk)?;
digest.update(chunk);
chunks.push(chunk);
}
anyhow::ensure!(
hex::encode(digest.finalize()) == header.chunks_sha256,
"Byte index was altered"
);
Ok(Some(Self {
binding: binding.clone(),
chunks,
commitment: header.chunks_sha256,
}))
}
/// Return only a slice of a completely verified aligned chunk. This bounds
/// seek/range work to64KiB and detects same-inode/same-stamp byte changes.
pub fn read_slice(&self, file: &mut File, start: u64, requested: usize) -> Result<Vec<u8>> {
anyhow::ensure!(
start < self.binding.size
&& requested > 0
&& file.metadata()?.len() == self.binding.size,
"Snapshot range or size changed"
);
let chunk = start / CHUNK_BYTES as u64;
let aligned = chunk * CHUNK_BYTES as u64;
let size = (self.binding.size - aligned).min(CHUNK_BYTES as u64) as usize;
let mut bytes = vec![0; size];
file.seek(SeekFrom::Start(aligned))?;
file.read_exact(&mut bytes)?;
let actual: [u8; 32] = Sha256::digest(&bytes).into();
anyhow::ensure!(
self.chunks.get(chunk as usize) == Some(&actual),
"Registered snapshot chunk changed; recover original entitlement"
);
let offset = (start - aligned) as usize;
Ok(bytes[offset..offset + requested.min(size - offset)].to_vec())
}
}
#[cfg(test)]
mod tests {
use super::*;
fn binding(bytes: &[u8]) -> Binding {
Binding {
content_id: format!("registered_{}", uuid::Uuid::new_v4()),
receipt_sha256: "ab".repeat(32),
full_sha256: hex::encode(Sha256::digest(bytes)),
size: bytes.len() as u64,
}
}
#[tokio::test]
async fn signed_index_reopens_without_whole_scan_and_seek_slices_reject_mutation() {
let root = tempfile::tempdir().unwrap();
let identity = NodeIdentity::load_or_create(&root.path().join("identity"))
.await
.unwrap();
let bytes: Vec<u8> = (0..CHUNK_BYTES * 5 + 19).map(|n| (n % 251) as u8).collect();
let media = root.path().join("media");
std::fs::write(&media, &bytes).unwrap();
let mut file = File::open(&media).unwrap();
let mut progress = Vec::new();
let original = binding(&bytes);
let index = Index::scan(&mut file, original.clone(), |n| {
progress.push(n);
Ok(())
})
.unwrap();
assert_eq!(progress.last(), Some(&(bytes.len() as u64)));
index.save(&root.path().join("index"), &identity).unwrap();
drop(index);
let loaded = Index::load(&root.path().join("index"), &original, &identity)
.unwrap()
.unwrap();
for start in [0, CHUNK_BYTES + 7, CHUNK_BYTES * 4, bytes.len() - 1] {
let got = loaded.read_slice(&mut file, start as u64, 37).unwrap();
assert_eq!(got, bytes[start..start + got.len()]);
}
let mut corrupt = bytes.clone();
corrupt[CHUNK_BYTES + 9] ^= 1;
std::fs::write(&media, corrupt).unwrap();
assert!(loaded
.read_slice(&mut file, (CHUNK_BYTES + 7) as u64, 1)
.is_err());
assert_eq!(loaded.read_slice(&mut file, 0, 5).unwrap(), &bytes[..5]);
let mut saved = std::fs::read(root.path().join("index")).unwrap();
*saved.last_mut().unwrap() ^= 1;
std::fs::write(root.path().join("index"), saved).unwrap();
assert!(Index::load(&root.path().join("index"), &original, &identity).is_err());
}
#[test]
fn canceled_or_wrong_hash_scan_never_produces_index() {
let root = tempfile::tempdir().unwrap();
let bytes = vec![7; CHUNK_BYTES * 2];
let media = root.path().join("media");
std::fs::write(&media, &bytes).unwrap();
let mut file = File::open(media).unwrap();
assert!(Index::scan(&mut file, binding(&bytes), |n| {
anyhow::ensure!(n == 0, "canceled");
Ok(())
})
.is_err());
let mut wrong = binding(&bytes);
wrong.full_sha256 = "cd".repeat(32);
assert!(Index::scan(&mut file, wrong, |_| Ok(())).is_err());
}
}
+315
View File
@@ -0,0 +1,315 @@
//! Bounded background verification. Preparing/ready never creates a lease.
use crate::rental_chunk_index::{Binding, Index};
use anyhow::{Context, Result};
use serde::{Deserialize, Serialize};
use std::{
collections::HashMap,
sync::{
atomic::{AtomicU64, Ordering},
Arc, Mutex, OnceLock,
},
time::Instant,
};
use tokio::sync::Semaphore;
const MAX_JOBS: usize = 8;
const MAX_INDEX_MEMORY: u64 = 64 * 1024 * 1024;
#[derive(Serialize, Deserialize, Debug)]
#[serde(tag = "state", rename_all = "snake_case")]
pub(crate) enum Status {
Preparing {
completed_bytes: u64,
total_bytes: u64,
},
Ready {
ready_id: String,
total_bytes: u64,
},
Unavailable,
}
struct Reservation {
used: Arc<AtomicU64>,
bytes: u64,
}
impl Drop for Reservation {
fn drop(&mut self) {
self.used.fetch_sub(self.bytes, Ordering::SeqCst);
}
}
pub(crate) struct Ready {
pub index: Index,
pub id: String,
_memory: Reservation,
}
#[cfg(test)]
impl Ready {
pub(crate) fn fixture(index: Index) -> Arc<Self> {
Arc::new(Self {
index,
id: uuid::Uuid::new_v4().to_string(),
_memory: Reservation {
used: Arc::new(AtomicU64::new(0)),
bytes: 0,
},
})
}
}
enum State {
Preparing,
Ready(Arc<Ready>),
Failed,
}
struct Job {
binding: Binding,
state: Mutex<State>,
progress: Arc<AtomicU64>,
touched: Instant,
}
pub(crate) struct Manager {
jobs: Mutex<HashMap<String, Arc<Job>>>,
workers: Arc<Semaphore>,
memory: Arc<AtomicU64>,
}
impl Default for Manager {
fn default() -> Self {
Self {
jobs: Mutex::new(HashMap::new()),
workers: Arc::new(Semaphore::new(2)),
memory: Arc::new(AtomicU64::new(0)),
}
}
}
pub(crate) fn shared() -> &'static Manager {
static INSTANCE: OnceLock<Manager> = OnceLock::new();
INSTANCE.get_or_init(Manager::default)
}
impl Manager {
#[cfg(test)]
pub(crate) fn forget_for_restart(&self, key: &str) {
self.jobs.lock().unwrap().remove(key);
}
fn status(job: &Job) -> Result<Status> {
Ok(
match &*job
.state
.lock()
.map_err(|_| anyhow::anyhow!("Media readiness unavailable"))?
{
State::Preparing => Status::Preparing {
completed_bytes: job.progress.load(Ordering::SeqCst),
total_bytes: job.binding.size,
},
State::Ready(ready) => Status::Ready {
ready_id: ready.id.clone(),
total_bytes: job.binding.size,
},
State::Failed => Status::Unavailable,
},
)
}
/// Work closure opens only the already-authorized immutable snapshot. At most
/// two closures execute, eight jobs exist, and64MiB is reserved for indexes.
/// Stream-held Arcs retain their reservation even after cache eviction.
pub fn prepare(
&self,
key: String,
binding: Binding,
work: impl FnOnce(Arc<AtomicU64>) -> Result<Index> + Send + 'static,
) -> Result<Status> {
let mut jobs = self
.jobs
.lock()
.map_err(|_| anyhow::anyhow!("Media readiness unavailable"))?;
if let Some(job) = jobs.get(&key) {
anyhow::ensure!(job.binding == binding, "Readiness immutable terms changed");
return Self::status(job);
}
let bytes = binding
.index_bytes()?
.checked_mul(2)
.and_then(|v| v.checked_add(128 * 1024))
.context("Index budget overflow")?;
loop {
let used = self.memory.load(Ordering::SeqCst);
if jobs.len() < MAX_JOBS
&& used
.checked_add(bytes)
.is_some_and(|n| n <= MAX_INDEX_MEMORY)
{
break;
}
let victim = jobs
.iter()
.filter(|(_, job)| {
job.state
.lock()
.map(|state| !matches!(*state, State::Preparing))
.unwrap_or(false)
})
.min_by_key(|(_, job)| job.touched)
.map(|(key, _)| key.clone());
if let Some(victim) = victim {
jobs.remove(&victim);
} else {
return Ok(Status::Preparing {
completed_bytes: 0,
total_bytes: binding.size,
});
}
}
let permit = match self.workers.clone().try_acquire_owned() {
Ok(permit) => permit,
Err(_) => {
return Ok(Status::Preparing {
completed_bytes: 0,
total_bytes: binding.size,
})
}
};
// Admissions are serialized by jobs; concurrent drops can only lower use.
self.memory.fetch_add(bytes, Ordering::SeqCst);
let reservation = Reservation {
used: self.memory.clone(),
bytes,
};
let progress = Arc::new(AtomicU64::new(0));
let job = Arc::new(Job {
binding: binding.clone(),
state: Mutex::new(State::Preparing),
progress: progress.clone(),
touched: Instant::now(),
});
jobs.insert(key, job.clone());
tokio::task::spawn_blocking(move || {
let _permit = permit;
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| work(progress)))
.unwrap_or_else(|_| Err(anyhow::anyhow!("Media verification worker failed")));
let mut state = match job.state.lock() {
Ok(state) => state,
Err(_) => return,
};
*state = match result {
Ok(index) if index.binding == binding => {
job.progress.store(binding.size, Ordering::SeqCst);
State::Ready(Arc::new(Ready {
index,
id: uuid::Uuid::new_v4().to_string(),
_memory: reservation,
}))
}
_ => State::Failed,
};
});
Ok(Status::Preparing {
completed_bytes: 0,
total_bytes: binding.size,
})
}
/// Explicit retry clears only a failed preparation, never a running job or
/// a lease. The caller has already reauthenticated the original purchase.
pub fn retry_failed(&self, key: &str, binding: &Binding) -> Result<()> {
let mut jobs = self
.jobs
.lock()
.map_err(|_| anyhow::anyhow!("Media readiness unavailable"))?;
if let Some(job) = jobs.get(key) {
anyhow::ensure!(&job.binding == binding, "Readiness immutable terms changed");
let failed = matches!(
*job.state
.lock()
.map_err(|_| anyhow::anyhow!("Media readiness unavailable"))?,
State::Failed
);
if failed {
jobs.remove(key);
}
}
Ok(())
}
pub fn ready(
&self,
key: &str,
binding: &Binding,
ready_id: Option<&str>,
) -> Result<Arc<Ready>> {
let jobs = self
.jobs
.lock()
.map_err(|_| anyhow::anyhow!("Media readiness unavailable"))?;
let job = jobs
.get(key)
.context("Media is preparing; no rental was started")?;
anyhow::ensure!(&job.binding == binding, "Readiness immutable terms changed");
let state = job
.state
.lock()
.map_err(|_| anyhow::anyhow!("Media readiness unavailable"))?;
let State::Ready(ready) = &*state else {
anyhow::bail!("Media is not ready; no rental was started");
};
anyhow::ensure!(
ready_id.is_none_or(|id| id == ready.id),
"Media readiness changed; prepare again without paying"
);
Ok(ready.clone())
}
}
#[cfg(test)]
mod tests {
use super::*;
use sha2::{Digest, Sha256};
#[tokio::test]
async fn slow_verification_is_nonblocking_deduplicated_and_has_no_start_side_effect() {
let manager = Manager::default();
let root = tempfile::tempdir().unwrap();
let media = root.path().join("media");
let bytes = vec![9; 1024 * 1024];
std::fs::write(&media, &bytes).unwrap();
let binding = Binding {
content_id: format!("registered_{}", uuid::Uuid::new_v4()),
receipt_sha256: "ab".repeat(32),
full_sha256: hex::encode(Sha256::digest(&bytes)),
size: bytes.len() as u64,
};
let (release, wait) = std::sync::mpsc::channel();
let clone = binding.clone();
assert!(matches!(
manager
.prepare("key".into(), binding.clone(), move |progress| {
wait.recv().unwrap();
let mut file = std::fs::File::open(media)?;
Index::scan(&mut file, clone, |n| {
progress.store(n, Ordering::SeqCst);
Ok(())
})
})
.unwrap(),
Status::Preparing { .. }
));
assert!(manager.ready("key", &binding, None).is_err());
assert!(matches!(
manager
.prepare("key".into(), binding.clone(), |_| panic!(
"Duplicate full scan"
))
.unwrap(),
Status::Preparing { .. }
));
release.send(()).unwrap();
let ready = tokio::time::timeout(std::time::Duration::from_secs(5), async {
loop {
if let Ok(ready) = manager.ready("key", &binding, None) {
break ready;
}
tokio::task::yield_now().await;
}
})
.await
.unwrap();
assert!(manager.ready("key", &binding, Some("other-ready")).is_err());
assert!(manager.ready("key", &binding, Some(&ready.id)).is_ok());
let restarted = Manager::default();
assert!(restarted.ready("key", &binding, Some(&ready.id)).is_err());
}
}
@@ -2214,6 +2214,10 @@ impl crate::content_purchase_caller::PurchaseTransport for PurchaseTestTransport
fn seller_onion(&self) -> &str {
"aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa.onion"
}
async fn prepare_offer(&self, content_id: &str) -> anyhow::Result<Option<(u64, u64)>> {
// A registered movie still being hashed, without asking the fake mint.
Ok(content_id.starts_with("registered_").then_some((3, 16)))
}
async fn offer(
&self,
id: &str,
@@ -2339,6 +2343,36 @@ async fn full_caller_recovers_lost_acceptance_and_settlement_without_another_pay
lose_settle: AtomicBool::new(true),
offers: Default::default(),
};
let preparing_id = format!("registered_{}", uuid::Uuid::new_v4());
for _ in 0..3 {
assert!(matches!(
purchase(
buyer.path(),
&buyer_did,
&preparing_id,
None,
8,
None,
&transport
)
.await
.unwrap(),
ReadyPurchase::Preparing {
completed_bytes: 3,
total_bytes: 16
}
));
}
assert!(transport.offers.lock().unwrap().is_empty());
assert!(mint.requests.lock().unwrap().is_empty());
assert_eq!(load_wallet(buyer.path()).await.unwrap().balance(), 8);
assert!(crate::content_purchase::Journal::open(buyer.path())
.await
.unwrap()
.find_buyers(&buyer_did, &transport.template.seller_did, &preparing_id)
.await
.unwrap()
.is_empty());
assert!(purchase(
buyer.path(),
&buyer_did,