Integrate recoverable native purchases, registered rentals and explicit payment consent

This commit is contained in:
archipelago
2026-10-06 22:44:06 -04:00
parent e4eae71314
commit 49703d7e88
63 changed files with 8028 additions and 134 deletions
@@ -0,0 +1,142 @@
//! Permanent Cloud snapshot delivery. Current source path/share state cannot
//! revoke a settled immutable snapshot; no rental clock is started here.
use super::{build_response, ApiHandler};
use crate::{
content_purchase::{Journal, SellerPhase},
content_server::ByteRange,
};
use anyhow::{Context, Result};
use hyper::{Body, HeaderMap, Response, StatusCode};
use tokio::io::{AsyncReadExt, AsyncSeekExt};
impl ApiHandler {
pub(super) async fn handle_cloud_purchase(
&self,
path: &str,
headers: &HeaderMap,
) -> Result<Response<Body>> {
let (content_id, purchase_id) = path
.strip_prefix("/content/")
.and_then(|value| value.split_once("/purchase/"))
.context("Invalid purchase delivery route")?;
anyhow::ensure!(
!content_id.contains('/')
&& !content_id.starts_with("registered_")
&& !purchase_id.contains('/'),
"Invalid Cloud delivery route"
);
let audience = crate::identity::did_key_from_pubkey_hex(&self.self_pubkey_hex)?;
let buyer = crate::content_auth::incoming(
headers,
&audience,
path,
chrono::Utc::now().timestamp(),
)?
.context("Authenticated peer proof is required")?;
let mut values = headers.get_all("x-content-capability").iter();
let capability = values
.next()
.context("Delivery capability is required")?
.to_str()?;
anyhow::ensure!(values.next().is_none(), "Duplicate delivery capability");
let (contract, mime) = {
let journal = Journal::open(&self.config.data_dir).await?;
let record = journal
.seller(purchase_id)
.await?
.context("Purchase settlement not found")?;
let receipt = match record.phase {
SellerPhase::ReceiptSaved(receipt) => receipt,
_ => anyhow::bail!("Purchase settlement is not durable"),
};
anyhow::ensure!(
record.contract.buyer_did == buyer
&& record.contract.seller_did == audience
&& record.contract.content_id == content_id
&& receipt.capability == capability,
"Purchase delivery binding changed"
);
let envelope = journal
.protocol_envelope("seller", purchase_id)
.await?
.context("Original delivery metadata is missing")?;
anyhow::ensure!(
envelope.contract()? == record.contract,
"Delivery metadata binding changed"
);
(record.contract, envelope.offer.mime_type)
};
let range = headers
.get("range")
.map(|value| -> Result<_> {
crate::content_server::parse_range_header(value.to_str()?).context("Invalid range")
})
.transpose()?;
let total = contract.content_size;
let (start, end, partial) = match range {
None => (0, total - 1, false),
Some(ByteRange::From { start, end }) => {
(start, end.unwrap_or(total - 1).min(total - 1), true)
}
Some(ByteRange::Suffix(count)) if count > 0 => {
(total.saturating_sub(count), total - 1, true)
}
_ => anyhow::bail!("Invalid range"),
};
if start > end || start >= total {
let mut response = build_response(
StatusCode::RANGE_NOT_SATISFIABLE,
"text/plain",
Body::empty(),
);
response
.headers_mut()
.insert("content-range", format!("bytes */{total}").parse()?);
return Ok(response);
}
let data = self.config.data_dir.clone();
let file = tokio::task::spawn_blocking(move || {
crate::content_snapshot::open_matching(
&data,
&contract.content_id,
&contract.content_sha256,
contract.content_size,
)
})
.await??;
let mut file = tokio::fs::File::from_std(file.file);
file.seek(std::io::SeekFrom::Start(start)).await?;
let length = end - start + 1;
let chunks =
futures_util::stream::try_unfold((file, length), |(mut file, left)| async move {
if left == 0 {
return Ok::<_, std::io::Error>(None);
}
let mut bytes = vec![0; left.min(65536) as usize];
let count = file.read(&mut bytes).await?;
if count == 0 {
return Err(std::io::Error::new(
std::io::ErrorKind::UnexpectedEof,
"Purchase snapshot ended early",
));
}
bytes.truncate(count);
Ok(Some((bytes, (file, left - count as u64))))
});
let mut response = Response::builder()
.status(if partial {
StatusCode::PARTIAL_CONTENT
} else {
StatusCode::OK
})
.header("content-type", mime)
.header("content-length", length)
.header("accept-ranges", "bytes")
.header("cache-control", "private, no-store")
.header("x-content-type-options", "nosniff")
.header("content-security-policy", "sandbox; default-src 'none'");
if partial {
response = response.header("content-range", format!("bytes {start}-{end}/{total}"));
}
Ok(response.body(Body::wrap_stream(chunks))?)
}
}
+6 -6
View File
@@ -416,7 +416,7 @@ impl ApiHandler {
path: &str,
) -> Result<Response<hyper::Body>> {
Ok(invoice_status_response(path, |hash, id| async move {
self.rpc_handler.settle_content_invoice(&hash, &id).await
self.rpc_handler.content_invoice_lifecycle(&hash, &id).await
})
.await)
}
@@ -663,7 +663,7 @@ impl ApiHandler {
async fn invoice_status_response<F, Fut>(path: &str, settle: F) -> Response<hyper::Body>
where
F: FnOnce(String, String) -> Fut,
Fut: std::future::Future<Output = Result<bool>>,
Fut: std::future::Future<Output = Result<serde_json::Value>>,
{
let parsed = path
.strip_prefix("/content/")
@@ -682,10 +682,10 @@ where
);
};
match settle(hash.to_ascii_lowercase(), id.to_owned()).await {
Ok(paid) => build_response(
Ok(body) => build_response(
StatusCode::OK,
"application/json",
hyper::Body::from(serde_json::json!({"paid": paid}).to_string()),
hyper::Body::from(body.to_string()),
),
Err(_) => {
tracing::warn!("Peer-file payment status verification is temporarily unavailable");
@@ -721,7 +721,7 @@ mod invoice_status_tests {
let response = invoice_status_response(path, |_, _| async {
panic!("Invalid request reached wallet");
#[allow(unreachable_code)]
Ok(false)
Ok(serde_json::json!({"paid":false}))
})
.await;
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
@@ -741,7 +741,7 @@ mod invoice_status_tests {
let response = invoice_status_response(&path, |hash, id| async move {
assert_eq!(hash, "ab".repeat(32));
assert_eq!(id, "file");
Ok(paid)
Ok(serde_json::json!({"paid":paid}))
})
.await;
assert_eq!(response.status(), StatusCode::OK);
+18
View File
@@ -1,7 +1,10 @@
mod purchase;
mod cloud_purchase;
mod blob;
mod cdp;
mod content;
mod registered_media;
mod rental_playback;
mod dwn;
mod model_proxy;
mod node_message;
@@ -385,6 +388,10 @@ impl ApiHandler {
let path = req.uri().path().to_string();
let method = req.method().clone();
if path.starts_with("/api/rental-playback/") {
return self.handle_local_rental_request(&method, &path, req.headers()).await;
}
// Handle CORS preflight for all routes
if method == Method::OPTIONS {
let mut builder = Response::builder()
@@ -445,6 +452,14 @@ impl ApiHandler {
.await;
}
// 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) {
return self.handle_purchase_request(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();
@@ -587,6 +602,9 @@ impl ApiHandler {
// 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/") && p.contains("/purchase/") => {
self.handle_cloud_purchase(p, &headers).await
}
(Method::GET, p) if p.starts_with("/content/registered_") && p.contains("/rental/") => {
self.handle_registered_rental(p, &headers).await
}
@@ -0,0 +1,170 @@
//! Add as api/handler/purchase.rs; dispatch only exact supported POST routes.
use super::{build_response, ApiHandler};
use crate::{
content_purchase::Journal, content_purchase_protocol as protocol, identity::NodeIdentity,
};
use anyhow::{Context, Result};
use hyper::{body::HttpBody, Body, Method, Request, Response, StatusCode};
use serde::Deserialize;
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct OfferRequest {
id: String,
content_id: String,
}
impl ApiHandler {
pub(super) async fn handle_purchase_request(
&self,
mut request: Request<Body>,
) -> Result<Response<Body>> {
let path = request.uri().path().to_owned();
anyhow::ensure!(
request.method() == Method::POST
&& matches!(
path.as_str(),
protocol::OFFER_ROUTE
| protocol::ACCEPT_ROUTE
| protocol::SETTLE_ROUTE
| protocol::STATUS_ROUTE
| protocol::CANCEL_ROUTE
),
"Unsupported purchase route"
);
let audience = crate::identity::did_key_from_pubkey_hex(&self.self_pubkey_hex)?;
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 <= 1024 * 1024),
"Purchase body too large"
);
bytes.extend_from_slice(&chunk);
}
Ok::<_, anyhow::Error>(bytes)
})
.await
.context("Purchase body timed out")??;
let buyer = crate::content_auth::authenticate_request(
request.headers(),
&audience,
&Method::POST,
&path,
&bytes,
chrono::Utc::now().timestamp(),
)?;
let data_dir = &self.config.data_dir;
let result = match path.as_str() {
protocol::OFFER_ROUTE => {
let body: OfferRequest = serde_json::from_slice(&bytes)?;
anyhow::ensure!(
uuid::Uuid::parse_str(&body.id)?.to_string() == body.id,
"Invalid operation identifier"
);
let saved = {
Journal::open(data_dir)
.await?
.protocol_offer(&body.id)
.await?
};
let offer = if let Some(saved) = saved {
anyhow::ensure!(
saved.buyer_did == buyer && saved.content_id == body.content_id,
"Original offer binding changed"
);
saved
} else {
// Registration pins and immutable snapshot are node-owned;
// no content hash/price/path is accepted from the request.
if !body.content_id.starts_with("registered_") {
let wallet = crate::wallet::ecash::load_wallet(data_dir).await?;
let offer = crate::content_cloud_offer::offer(
data_dir,
&body.id,
&body.content_id,
&buyer,
&audience,
crate::wallet::ecash::load_network(data_dir).await?,
wallet.mint_url.trim_end_matches('/'),
crate::content_cloud_offer::SnapshotPolicy {
max_file_bytes: 64 * 1024 * 1024 * 1024,
max_total_bytes: 64 * 1024 * 1024 * 1024,
minimum_free_bytes: 512 * 1024 * 1024,
},
)
.await?;
return Ok(build_response(
StatusCode::OK,
"application/json",
Body::from(serde_json::to_vec(&offer)?),
));
}
let identity = NodeIdentity::load_existing(&data_dir.join("identity")).await?;
anyhow::ensure!(identity.did_key()? == audience, "Node identity changed");
let selected = body.content_id.clone();
let root = data_dir.clone();
let (receipt, terms) = tokio::task::spawn_blocking(move || {
crate::registered_media::registered_terms(&root, &identity, &selected)
})
.await??;
anyhow::ensure!(
receipt
.payment_methods
.iter()
.any(|method| method == "cashu"),
"Content does not accept Cashu"
);
let now = chrono::Utc::now().timestamp();
let deadline = now.checked_add(120).context("Offer clock overflow")?;
let wallet = crate::wallet::ecash::load_wallet(data_dir).await?;
let offer = protocol::Offer {
id: body.id,
buyer_did: buyer.clone(),
seller_did: audience.clone(),
filename: receipt.content_id.clone(),
mime_type: "application/octet-stream".into(),
content_id: receipt.content_id,
content_sha256: receipt.sha256,
content_size: receipt.size_bytes.parse()?,
viewing_seconds: Some(receipt.viewing_seconds),
terms_sha256: terms,
network: crate::wallet::ecash::load_network(data_dir).await?,
mint_url: wallet.mint_url.trim_end_matches('/').to_owned(),
seller_net_sats: receipt.price_sats,
offered_at: now,
expires_at: deadline,
};
protocol::save_offer(data_dir, &offer, &buyer, now).await?
};
protocol::ensure_seller_mint_policy(data_dir, offer.network, &offer.mint_url)
.await?;
serde_json::to_value(offer)?
}
protocol::ACCEPT_ROUTE => serde_json::to_value(
protocol::accept(data_dir, &serde_json::from_slice(&bytes)?, &buyer, || {
chrono::Utc::now().timestamp()
})
.await?,
)?,
protocol::SETTLE_ROUTE => serde_json::to_value(
protocol::settle(data_dir, &serde_json::from_slice(&bytes)?, &buyer).await?,
)?,
protocol::CANCEL_ROUTE => serde_json::to_value(
protocol::cancel(data_dir, &serde_json::from_slice(&bytes)?, &buyer).await?,
)?,
protocol::STATUS_ROUTE => serde_json::to_value(
protocol::status(data_dir, &serde_json::from_slice(&bytes)?, &buyer).await?,
)?,
_ => unreachable!(),
};
Ok(build_response(
StatusCode::OK,
"application/json",
Body::from(serde_json::to_vec(&result)?),
))
}
}
@@ -0,0 +1,548 @@
//! Local browser playback. Only opaque local handles cross the browser boundary.
use super::{build_response, ApiHandler};
use crate::{content_purchase::Journal, content_server::ByteRange, identity::NodeIdentity};
use anyhow::{Context, Result};
use hyper::{Body, HeaderMap, Method, Response, StatusCode};
use std::{
io,
sync::Arc,
time::{Duration, Instant},
};
fn requested_bounds(headers: &HeaderMap, total: u64) -> Result<Option<(u64, u64)>> {
anyhow::ensure!(total > 0, "Empty purchased media");
anyhow::ensure!(
headers.get_all("range").iter().count() <= 1,
"Ambiguous playback ranges"
);
let Some(header) = headers.get("range") else {
return Ok(None);
};
let range = crate::content_server::parse_range_header(header.to_str()?)
.context("Invalid playback byte range")?;
let last = total - 1;
let (start, end) = match range {
ByteRange::From { start, end } => (start, end.unwrap_or(last).min(last)),
ByteRange::Suffix(count) => {
anyhow::ensure!(count > 0, "Invalid byte range");
(total.saturating_sub(count), last)
}
};
anyhow::ensure!(start <= end && start < total, "Invalid playback byte range");
Ok(Some((start, end)))
}
fn validate_upstream(
status: u16,
headers: &HeaderMap,
total: u64,
bounds: Option<(u64, u64)>,
now: u64,
) -> Result<(u16, u64, String, u64)> {
let expected_status = if bounds.is_some() { 206 } else { 200 };
anyhow::ensure!(
status == expected_status,
"Seller returned another range status"
);
let length = bounds.map_or(total, |(start, end)| end - start + 1);
anyhow::ensure!(
headers
.get("content-length")
.and_then(|v| v.to_str().ok())
.and_then(|v| v.parse::<u64>().ok())
== Some(length),
"Seller changed purchased byte length"
);
if let Some((start, end)) = bounds {
let expected = format!("bytes {start}-{end}/{total}");
anyhow::ensure!(
headers.get("content-range").and_then(|v| v.to_str().ok()) == Some(expected.as_str()),
"Seller changed purchased byte range"
);
}
let mime = headers
.get("content-type")
.context("Missing media type")?
.to_str()?
.to_owned();
anyhow::ensure!(
mime.starts_with("video/") || mime.starts_with("audio/"),
"Unsupported rental media type"
);
let expires = headers
.get("x-rental-expires-at")
.context("Missing rental expiry")?
.to_str()?
.parse::<u64>()?;
anyhow::ensure!(expires > now, "Rental viewing window ended");
Ok((expected_status, length, mime, expires))
}
fn installed_playback_origin(
origin: &str,
expected: &str,
host: &str,
gated_tls_port: bool,
) -> bool {
let (Ok(actual), Ok(expected), Ok(request)) = (
reqwest::Url::parse(origin),
reqwest::Url::parse(expected),
reqwest::Url::parse(&format!("http://{host}")),
) else {
return false;
};
if !matches!(actual.scheme(), "http" | "https")
|| actual.origin().ascii_serialization() != origin
|| request.path() != "/"
|| !request.username().is_empty()
|| request.password().is_some()
|| request.query().is_some()
|| request.fragment().is_some()
{
return false;
}
if actual.origin() == expected.origin() {
return true;
}
let same_port = actual.port_or_known_default() == expected.port_or_known_default();
let scheme = actual.scheme() == expected.scheme()
|| (gated_tls_port && actual.scheme() == "https" && expected.scheme() == "http");
let expected_loopback = matches!(
expected.host_str(),
Some("localhost" | "127.0.0.1" | "[::1]")
);
same_port
&& scheme
&& actual.host_str() == request.host_str()
&& (expected_loopback || actual.host_str() == expected.host_str())
}
fn unix_now() -> Result<u64> {
Ok(std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)?
.as_secs())
}
fn stream_error(message: &'static str) -> io::Error {
io::Error::new(io::ErrorKind::PermissionDenied, message)
}
impl ApiHandler {
pub(super) async fn handle_local_rental_request(
&self,
method: &Method,
path: &str,
headers: &HeaderMap,
) -> Result<Response<Body>> {
// HEAD is deliberately not GET: a browser probe must never open a lease.
if method != Method::GET && method != Method::OPTIONS {
return Ok(Response::builder()
.status(StatusCode::METHOD_NOT_ALLOWED)
.header("Allow", "GET, OPTIONS")
.header("Cache-Control", "no-store")
.body(Body::empty())?);
}
let origin = headers.get("origin").map(|v| v.to_str()).transpose()?;
if let Some(origin) = origin {
let identity =
NodeIdentity::load_existing(&self.config.data_dir.join("identity")).await?;
let (state, _) = self.state_manager.get_snapshot().await;
let root = self.config.data_dir.clone();
let context = tokio::task::spawn_blocking(move || {
crate::container::registration_pin::installed_context(&root, &identity, &state)
})
.await??;
let host = headers
.get("host")
.and_then(|value| value.to_str().ok())
.unwrap_or("");
let port = reqwest::Url::parse(origin)
.ok()
.and_then(|url| url.port_or_known_default());
let ports = self.rpc_handler.app_gate.port_map().await;
let gated_tls_port = port
.and_then(|port| ports.gated(port))
.is_some_and(|gate| gate.app_id == context.app_id && gate.declared);
if !context
.app_origins
.iter()
.any(|allowed| installed_playback_origin(origin, allowed, host, gated_tls_port))
{
return Ok(build_response(
StatusCode::FORBIDDEN,
"text/plain",
Body::from("Playback origin is not the installed app"),
));
}
}
let mut response = if method == Method::OPTIONS {
Response::builder()
.status(StatusCode::NO_CONTENT)
.body(Body::empty())?
} else {
self.handle_local_rental(path, headers).await?
};
response
.headers_mut()
.insert("Cache-Control", "private, no-store".parse()?);
response.headers_mut().insert("Vary", "Origin".parse()?);
if let Some(origin) = origin {
response
.headers_mut()
.insert("Access-Control-Allow-Origin", origin.parse()?);
response
.headers_mut()
.insert("Access-Control-Allow-Credentials", "true".parse()?);
response
.headers_mut()
.insert("Access-Control-Allow-Methods", "GET, OPTIONS".parse()?);
response
.headers_mut()
.insert("Access-Control-Allow-Headers", "Range".parse()?);
response.headers_mut().insert(
"Access-Control-Expose-Headers",
"Content-Length, Content-Range, Accept-Ranges, X-Rental-Expires-At".parse()?,
);
}
Ok(response)
}
/// Dispatcher accepts GET only after normal session handling. HEAD and other
/// methods never reach upstream, so metadata probes cannot start a lease.
pub(super) async fn handle_local_rental(
&self,
path: &str,
headers: &HeaderMap,
) -> Result<Response<Body>> {
let token = match crate::session::extract_session_cookie(headers) {
Some(token) if self.session_store.validate(&token).await => token,
_ => return Ok(Self::unauthorized()),
};
let handle = path
.strip_prefix("/api/rental-playback/")
.context("Invalid playback route")?;
let identity =
Arc::new(NodeIdentity::load_existing(&self.config.data_dir.join("identity")).await?);
let (state, _) = self.state_manager.get_snapshot().await;
let root = self.config.data_dir.clone();
let key = identity.clone();
let context = tokio::task::spawn_blocking(move || {
crate::container::registration_pin::installed_context(&root, &key, &state)
})
.await??;
let binding = self
.rpc_handler
.playback_handles()
.lookup(handle, &token, &context)?;
let capability = {
let journal = Journal::open(&self.config.data_dir).await?;
let record = journal
.buyer(&binding.contract.id)
.await?
.context("Original purchase is missing")?;
anyhow::ensure!(
record.contract == binding.contract,
"Original purchase changed"
);
record
.receipt()
.context("Original purchase is not settled")?
.capability
.clone()
};
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 is unavailable")?;
let total = binding.contract.content_size;
let bounds = match requested_bounds(headers, total) {
Ok(bounds) => bounds,
Err(_) => {
return Ok(Response::builder()
.status(StatusCode::RANGE_NOT_SATISFIABLE)
.header("Content-Range", format!("bytes */{total}"))
.body(Body::empty())?)
}
};
let remote_path = format!(
"/content/{}/rental/{}",
binding.contract.content_id, binding.contract.id
);
let mut request =
crate::fips::dial::PeerRequest::new(Some(&mesh), &binding.seller_onion, &remote_path)
.require_fips()
.single_delivery()
.timeout(Duration::from_secs(24 * 60 * 60))
.header("X-Content-Capability", capability);
if let Some((start, end)) = bounds {
request = request.header("Range", format!("bytes={start}-{end}"));
}
let (response, transport) = tokio::time::timeout(
Duration::from_secs(20),
request.send_content_get(&self.config.data_dir),
)
.await
.context("Seller did not begin the original rental stream")??;
if !response.status().is_success() {
// Never forward arbitrary upstream bodies, redirects, cookies or private headers.
let status = if response.status().as_u16() == 403 {
StatusCode::FORBIDDEN
} else {
StatusCode::BAD_GATEWAY
};
return Ok(build_response(
status,
"text/plain",
Body::from(
"Original rental is unavailable; recover this purchase without paying again",
),
));
}
let (expected_status, length, mime, expires) = validate_upstream(
response.status().as_u16(),
response.headers(),
total,
bounds,
unix_now()?,
)?;
self.rpc_handler
.playback_handles()
.note_expiry(handle, &binding, expires)?;
let sessions = self.session_store.clone();
let state_manager = self.state_manager.clone();
let data_dir = self.config.data_dir.clone();
let chunks = futures_util::stream::try_unfold(
(response, length, None::<Instant>),
move |(mut response, left, mut checked)| {
let sessions = sessions.clone();
let token = token.clone();
let state_manager = state_manager.clone();
let data_dir = data_dir.clone();
let identity = identity.clone();
let context = context.clone();
async move {
if unix_now().map_err(|_| stream_error("Playback clock unavailable"))?
>= expires
{
return Err(stream_error("Rental viewing window ended"));
}
if left == 0 {
return Ok::<_, io::Error>(None);
}
let waiting_since = Instant::now();
loop {
if waiting_since.elapsed() >= Duration::from_secs(30) {
return Err(io::Error::new(
io::ErrorKind::TimedOut,
"Rental stream stalled; reopen the original purchase",
));
}
if unix_now().map_err(|_| stream_error("Playback clock unavailable"))?
>= expires
{
return Err(stream_error("Rental viewing window ended"));
}
if checked.is_none_or(|at| at.elapsed() >= Duration::from_secs(1)) {
if !sessions.validate(&token).await {
return Err(stream_error("Playback session ended"));
}
let (state, _) = state_manager.get_snapshot().await;
let root = data_dir.clone();
let key = identity.clone();
let actual = tokio::task::spawn_blocking(move || {
crate::container::registration_pin::installed_context(
&root, &key, &state,
)
})
.await
.map_err(|_| stream_error("Playback app context unavailable"))?
.map_err(|_| stream_error("Playback app context unavailable"))?;
if actual != context {
return Err(stream_error("Playback app context changed"));
}
checked = Some(Instant::now());
}
// Keep checking revocation while the peer stalls; no local media cache.
let chunk = tokio::select! {
chunk=response.chunk() => chunk.map_err(|_|io::Error::new(io::ErrorKind::ConnectionAborted,"Rental stream interrupted; reopen the original purchase"))?,
_=tokio::time::sleep(Duration::from_secs(1)) => continue,
};
let bytes = chunk.ok_or_else(|| {
io::Error::new(
io::ErrorKind::UnexpectedEof,
"Purchased media ended early",
)
})?;
if bytes.len() as u64 > left {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"Purchased media exceeded its declared length",
));
}
let remaining = left - bytes.len() as u64;
return Ok(Some((bytes, (response, remaining, checked))));
}
}
},
);
let mut result = Response::builder()
.status(expected_status)
.header("Content-Type", mime)
.header("Content-Length", length)
.header("Accept-Ranges", "bytes")
.header("Cache-Control", "private, no-store")
.header("X-Content-Type-Options", "nosniff")
.header("X-Rental-Expires-At", expires)
.header("X-Archipelago-Transport", transport.to_string());
if let Some((start, end)) = bounds {
result = result.header("Content-Range", format!("bytes {start}-{end}/{total}"));
}
Ok(result.body(Body::wrap_stream(chunks))?)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn installed_origin_maps_only_current_host_and_verified_app_port() {
assert!(installed_playback_origin(
"http://192.168.1.5:7778",
"http://127.0.0.1:7778",
"192.168.1.5",
false
));
assert!(installed_playback_origin(
"https://192.168.1.5:7778",
"http://127.0.0.1:7778",
"192.168.1.5",
true
));
assert!(installed_playback_origin(
"https://[fd00::5]:7778",
"https://[::1]:7778",
"[fd00::5]:443",
false
));
for (actual, host, tls) in [
("https://192.168.1.5:7778", "192.168.1.5", false),
("http://evil.test:7778", "192.168.1.5", true),
("http://192.168.1.5:7779", "192.168.1.5", true),
("http://192.168.1.5:7778", "evil.test", true),
("http://192.168.1.5:7778/path", "192.168.1.5", true),
("http://192.168.1.5:7778", "user@192.168.1.5", true),
] {
assert!(
!installed_playback_origin(actual, "http://127.0.0.1:7778", host, tls),
"{actual} {host}"
);
}
}
#[test]
fn range_bounds_follow_purchased_size_and_reject_ambiguous_ranges() {
let check = |range: &str| {
let mut h = HeaderMap::new();
h.insert("range", range.parse().unwrap());
requested_bounds(&h, 100)
};
assert_eq!(check("bytes=20-39").unwrap(), Some((20, 39)));
assert_eq!(check("bytes=90-").unwrap(), Some((90, 99)));
assert_eq!(check("bytes=-10").unwrap(), Some((90, 99)));
assert_eq!(check("bytes=-200").unwrap(), Some((0, 99)));
assert_eq!(check("bytes=90-500").unwrap(), Some((90, 99)));
for range in [
"bytes=100-",
"bytes=20-10",
"bytes=-0",
"bytes=0-1,4-6",
"other=0-1",
] {
assert!(check(range).is_err(), "{range}");
}
assert_eq!(requested_bounds(&HeaderMap::new(), 100).unwrap(), None);
assert!(requested_bounds(&HeaderMap::new(), 0).is_err());
}
#[test]
fn upstream_range_expiry_and_media_headers_are_bound_before_bytes_escape() {
let mut headers = HeaderMap::new();
for (name, value) in [
("content-length", "20"),
("content-range", "bytes 20-39/100"),
("content-type", "video/mp4"),
("x-rental-expires-at", "200"),
] {
headers.insert(name, value.parse().unwrap());
}
assert_eq!(
validate_upstream(206, &headers, 100, Some((20, 39)), 100).unwrap(),
(206, 20, "video/mp4".into(), 200)
);
assert!(validate_upstream(200, &headers, 100, Some((20, 39)), 100).is_err());
assert!(validate_upstream(206, &headers, 100, Some((20, 39)), 200).is_err());
for (name, bad) in [
("content-length", "21"),
("content-range", "bytes 21-40/100"),
("content-type", "text/html"),
("x-rental-expires-at", "0"),
] {
let mut changed = headers.clone();
changed.insert(name, bad.parse().unwrap());
assert!(
validate_upstream(206, &changed, 100, Some((20, 39)), 100).is_err(),
"{name}"
);
}
headers.remove("content-range");
headers.insert("content-length", "100".parse().unwrap());
assert!(validate_upstream(200, &headers, 100, None, 100).is_ok());
}
#[tokio::test]
async fn metadata_probes_and_unauthenticated_get_do_not_touch_purchase_or_identity() {
let root = tempfile::tempdir().unwrap();
let mut config = crate::config::Config::default();
config.data_dir = root.path().to_path_buf();
let handler = ApiHandler::new(
config,
Arc::new(crate::state::StateManager::new()),
Arc::new(crate::monitoring::MetricsStore::new()),
None,
None,
)
.await
.unwrap();
// ApiHandler initialization may establish its own node identity, but the
// denied route must not need installed apps, saved receipts or any peer.
for method in [Method::HEAD, Method::POST] {
let result = handler
.handle_request(
hyper::Request::builder()
.method(method)
.uri("/api/rental-playback/invalid")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(result.status(), StatusCode::METHOD_NOT_ALLOWED);
}
let result = handler
.handle_request(
hyper::Request::builder()
.uri("/api/rental-playback/invalid")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(result.status(), StatusCode::UNAUTHORIZED);
assert!(!root.path().join("content-purchases").exists());
}
}
@@ -332,7 +332,24 @@ 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.playback-handle" => self.handle_playback_handle(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,
"content.cancel-purchase" => self.handle_content_cancel_purchase(params).await,
"content.payment-status" => self.handle_content_payment_status(params).await,
"media.registration.context" => {
self.handle_media_registration_context(params.unwrap_or_default())
.await
}
"media.registration.prepare" => {
self.handle_media_registration_prepare(params.unwrap_or_default())
.await
}
"media.registration.resolve" => {
self.handle_media_registration_resolve(params.unwrap_or_default())
.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,
+125
View File
@@ -7,6 +7,57 @@ use zeroize::Zeroize;
use super::LND_REST_BASE_URL;
impl RpcHandler {
// Add inside api/rpc/lnd/wallet.rs RpcHandler impl; seller owns this lookup.
// Wire invoice-status response to this Value instead of reducing it to paid bool.
pub(crate) async fn content_invoice_lifecycle(
&self,
hash: &str,
content_id: &str,
) -> Result<serde_json::Value> {
anyhow::ensure!(
hash.len() == 64 && hash.bytes().all(|c| c.is_ascii_hexdigit()),
"Invalid payment hash"
);
let hash = hash.to_ascii_lowercase();
let existing = crate::content_invoice::lookup(&self.config.data_dir, &hash).await?;
anyhow::ensure!(
existing.as_ref().is_none_or(|(id, _)| id == content_id),
"Invoice belongs to another content item"
);
if crate::content_invoice::is_paid_for(&self.config.data_dir, &hash, content_id).await {
return Ok(
serde_json::json!({"paid":true,"state":"settled","can_switch_method":false}),
);
}
let (client, macaroon_hex) = self.lnd_client().await?;
let response = client
.get(format!("{LND_REST_BASE_URL}/v1/invoice/{hash}"))
.header("Grpc-Metadata-macaroon", &macaroon_hex)
.send()
.await?;
if response.status() == reqwest::StatusCode::NOT_FOUND {
return Ok(
serde_json::json!({"paid":false,"state":"unknown","can_switch_method":false}),
);
}
let body: serde_json::Value = response.error_for_status()?.json().await?;
let price = content_invoice_amount(&body, content_id)
.context("Invoice content binding is unavailable")?;
anyhow::ensure!(
existing
.as_ref()
.is_none_or(|(_, expected)| *expected == price),
"Invoice price binding changed"
);
crate::content_invoice::record_pending(&self.config.data_dir, &hash, content_id, price)
.await?;
let result = content_invoice_lifecycle_body(&body, price);
if result["paid"] == true {
crate::content_invoice::mark_paid(&self.config.data_dir, &hash).await?;
}
Ok(result)
}
/// Generate a new on-chain Bitcoin address.
pub(in crate::api::rpc) async fn handle_lnd_newaddress(&self) -> Result<serde_json::Value> {
let (client, macaroon_hex) = self.lnd_client().await.map_err(|e| {
@@ -1561,3 +1612,77 @@ mod peer_file_invoice_tests {
}
}
}
fn content_invoice_lifecycle_body(body: &serde_json::Value, price: u64) -> serde_json::Value {
let settled = content_invoice_fully_settled(body, price);
let cancelled = !settled
&& body["state"] == "CANCELED"
&& body.get("settled").and_then(|value| value.as_bool()) != Some(true)
&& body.get("amt_paid_sat").and_then(json_u64) == Some(0)
&& body
.get("amt_paid_msat")
.is_none_or(|value| json_u64(value) == Some(0));
let state = if settled {
"settled"
} else if cancelled {
"canceled"
} else {
match body.get("state").and_then(|value| value.as_str()) {
Some("OPEN") => "open",
Some("ACCEPTED") => "accepted",
_ => "unknown",
}
};
let expires_at = body
.get("creation_date")
.and_then(json_u64)
.zip(body.get("expiry").and_then(json_u64))
.and_then(|(created, expiry)| created.checked_add(expiry));
// Expiry is informational. Only LND's terminal canceled state releases the
// cross-method block; wall-clock passage or lookup failure never does.
serde_json::json!({"paid":settled,"state":state,"can_switch_method":cancelled,
"expires_at":expires_at,"cancel_supported":false})
}
#[cfg(test)]
mod invoice_lifecycle_tests {
use super::*;
#[test]
fn only_authoritative_zero_paid_cancel_unlocks_and_settlement_survives_expiry() {
for state in ["OPEN", "ACCEPTED", "UNKNOWN", "CANCELED"] {
for paid in [0u64, 1] {
let result = content_invoice_lifecycle_body(
&serde_json::json!({
"state":state,"settled":false,"amt_paid_sat":paid.to_string(),
"creation_date":"1","expiry":"1"}),
8,
);
assert_eq!(
result["can_switch_method"],
state == "CANCELED" && paid == 0
);
assert_eq!(result["paid"], false);
}
}
assert_eq!(
content_invoice_lifecycle_body(&serde_json::json!({"state":"CANCELED"}), 8)
["can_switch_method"],
false
);
assert_eq!(
content_invoice_lifecycle_body(
&serde_json::json!({"state":"CANCELED", "amt_paid_sat":"0", "amt_paid_msat":"1"}),
8
)["can_switch_method"],
false
);
let paid = content_invoice_lifecycle_body(
&serde_json::json!({
"state":"SETTLED","settled":true,"amt_paid_sat":"8","value":"8",
"creation_date":"1","expiry":"1"}),
8,
);
assert_eq!(paid["paid"], true);
assert_eq!(paid["can_switch_method"], false);
}
}
@@ -0,0 +1,353 @@
//! Owner-session/CSRF RPC plus producer-signed, exact media approval.
//! App callers use the native dashboard bridge; this is never an origin-only grant.
use super::RpcHandler;
use crate::media_registration::{AuthorizedSelection, Intent, Limits};
use anyhow::{Context, Result};
use nostr_sdk::prelude::{Event, Kind};
use serde::Deserialize;
use std::{
path::Path,
sync::{
atomic::{AtomicBool, Ordering},
Arc,
},
};
// Dropping the request future cancels queued locks and chunked snapshot work.
// A completed durable record is still recovered by the original operation ID.
struct RequestCancellation(Arc<AtomicBool>);
impl Drop for RequestCancellation {
fn drop(&mut self) {
self.0.store(true, Ordering::Relaxed);
}
}
fn request_cancellation() -> (RequestCancellation, Arc<AtomicBool>) {
let signal = Arc::new(AtomicBool::new(false));
(RequestCancellation(signal.clone()), signal)
}
const DOMAIN: &str = "archipelago.media-registration.approval.v1";
const KIND: u16 = 27236;
const MAX_APPROVAL: usize = 32 * 1024;
#[derive(Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
struct Params {
intent: Intent,
selection: AuthorizedSelection,
producer_event: Event,
}
const RESOLUTION_DOMAIN: &str = "archipelago.media-registration.resolution.v1";
#[derive(Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
struct ResolveParams {
intent: Intent,
producer_event: Event,
}
fn verified_resolution(input: serde_json::Value, now: u64) -> Result<ResolveParams> {
anyhow::ensure!(
serde_json::to_vec(&input)?.len() <= MAX_APPROVAL,
"Resolution approval is too large"
);
let params: ResolveParams = serde_json::from_value(input)?;
let event = &params.producer_event;
event
.verify()
.context("Resolution producer signature failed")?;
anyhow::ensure!(
event.kind == Kind::Custom(27237) && event.pubkey.to_hex() == params.intent.producer,
"Resolution signature purpose or producer changed"
);
let encoded = serde_json::to_value(event)?;
anyhow::ensure!(
encoded["tags"] == serde_json::json!([["d", RESOLUTION_DOMAIN]])
&& event.created_at.as_u64() >= params.intent.created_at.saturating_sub(30)
&& event.created_at.as_u64() <= now.saturating_add(30),
"Invalid resolution signature time or scope"
);
let content: serde_json::Value = serde_json::from_str(&event.content)?;
anyhow::ensure!(
content
== serde_json::json!({
"action":"Recover prepared video or retire this expired incomplete registration",
"scope":RESOLUTION_DOMAIN, "intent":params.intent,
}),
"Resolution signature changed the original intent"
);
Ok(params)
}
fn approval_content(intent: &Intent, selection: &AuthorizedSelection) -> Result<serde_json::Value> {
let path = selection
.relative_path
.to_str()
.context("Cloud file name is not UTF-8")?;
Ok(serde_json::json!({
"action":"Register this Cloud video for an IndeeHub project",
"scope":DOMAIN,
"intent":intent,
"selection":{"cloudFile":path,"paymentMethods":selection.payment_methods},
}))
}
fn verified_producer(params: &Params, now: u64) -> Result<String> {
let event = &params.producer_event;
anyhow::ensure!(
event.kind == Kind::Custom(KIND),
"This signature is not a media registration approval"
);
event
.verify()
.context("Producer approval signature failed")?;
let producer = event.pubkey.to_hex();
anyhow::ensure!(
producer == params.intent.producer,
"The signing identity differs from the project producer"
);
let created = event.created_at.as_u64();
anyhow::ensure!(
created >= params.intent.created_at.saturating_sub(30)
&& created < params.intent.expires_at
&& created <= now.saturating_add(30),
"Producer approval was not signed within this registration intent"
);
let encoded = serde_json::to_value(event)?;
anyhow::ensure!(
encoded["tags"] == serde_json::json!([["d", DOMAIN]]),
"Media approval signature scope changed"
);
let content: serde_json::Value = serde_json::from_str(&event.content)
.context("Producer approval is not readable registration terms")?;
anyhow::ensure!(
content == approval_content(&params.intent, &params.selection)?,
"Approved project, Cloud selection or rental terms changed"
);
Ok(producer)
}
fn parse(input: serde_json::Value, now: u64) -> Result<(Params, String)> {
anyhow::ensure!(
serde_json::to_vec(&input)?.len() <= MAX_APPROVAL,
"Registration approval is too large"
);
let params: Params = serde_json::from_value(input).context("Invalid registration approval")?;
let producer = verified_producer(&params, now)?;
Ok((params, producer))
}
fn registration_limit(_data_dir: &Path) -> u64 {
// Explicit per-file staging bound; shared immutable snapshot storage handles
// disk reservations separately before this route is enabled for live apps.
16 * 1024 * 1024 * 1024
}
impl RpcHandler {
/// Read-only public installation bindings for the native consent bridge.
/// A hidden standalone signer must compare its actual app origin with these
/// installed addresses; a caller-supplied app name is never sufficient.
pub(super) async fn handle_media_registration_context(
&self,
_input: serde_json::Value,
) -> Result<serde_json::Value> {
let (state, _) = self.state_manager.get_snapshot().await;
let identity =
crate::identity::NodeIdentity::load_existing(&self.config.data_dir.join("identity"))
.await?;
let data_dir = self.config.data_dir.clone();
let context = tokio::task::spawn_blocking(move || {
crate::container::registration_pin::installed_context(&data_dir, &identity, &state)
})
.await??;
let mut result = serde_json::to_value(context)?;
result["paymentMethods"] = serde_json::json!(["cashu"]);
Ok(result)
}
pub(super) async fn handle_media_registration_resolve(
&self,
input: serde_json::Value,
) -> Result<serde_json::Value> {
let now = u64::try_from(chrono::Utc::now().timestamp()).context("Invalid node clock")?;
let params = verified_resolution(input, now)?;
let (state, _) = self.state_manager.get_snapshot().await;
let identity =
crate::identity::NodeIdentity::load_existing(&self.config.data_dir.join("identity"))
.await?;
let data_dir = self.config.data_dir.clone();
let (_cancellation, cancelled) = request_cancellation();
tokio::task::spawn_blocking(move || {
crate::container::registration_pin::installed_context(&data_dir, &identity, &state)?;
crate::registered_media::resolve_registration(
&data_dir,
&identity,
&params.intent,
&params.producer_event.pubkey.to_hex(),
now,
&Limits {
max_bytes: registration_limit(&data_dir),
cancelled: &cancelled,
},
)
})
.await?
}
/// Standard RPC front door already requires owner session, permitted origin
/// and CSRF. Producer signature is additional exact-scope consent, not a
/// replacement for those owner checks. A replay repeats the original UUID.
pub(super) async fn handle_media_registration_prepare(
&self,
input: serde_json::Value,
) -> Result<serde_json::Value> {
let now = u64::try_from(chrono::Utc::now().timestamp()).context("Invalid node clock")?;
let (params, producer) = parse(input, now)?;
let (state, _) = self.state_manager.get_snapshot().await;
anyhow::ensure!(
state
.package_data
.get("indeedhub-api")
.is_some_and(|entry| matches!(
entry.state,
crate::data_model::PackageState::Running
)),
"The installed IndeeHub API must be running to register its media"
);
let identity =
crate::identity::NodeIdentity::load_existing(&self.config.data_dir.join("identity"))
.await?;
let data_dir = self.config.data_dir.clone();
let (_cancellation, cancelled) = request_cancellation();
let receipt = tokio::task::spawn_blocking(move || {
crate::container::registration_pin::installed_context(&data_dir, &identity, &state)?;
let project = params.intent.project_id.clone();
let limits = Limits {
max_bytes: registration_limit(&data_dir),
cancelled: &cancelled,
};
crate::registered_media::register_approved_selection(
&data_dir,
&data_dir.join("filebrowser"),
&identity,
&crate::registered_media::ApprovedSelection {
authenticated_producer: &producer,
authenticated_project: &project,
intent: &params.intent,
selection: &params.selection,
},
now,
&limits,
|_| Ok(()),
)
})
.await??;
Ok(serde_json::to_value(receipt)?)
}
}
#[cfg(test)]
mod tests {
use super::*;
use nostr_sdk::prelude::{EventBuilder, Keys, Tag, Timestamp};
fn fixture() -> Params {
let keys = Keys::parse(&"07".repeat(32)).unwrap();
let intent = Intent {
version: 1,
request_id: uuid::Uuid::new_v4().to_string(),
nonce: "ab".repeat(32),
app_audience: uuid::Uuid::new_v4().to_string(),
node_did: crate::identity::did_key_from_pubkey_hex(&hex::encode([7; 32])).unwrap(),
producer: keys.public_key().to_hex(),
project_id: "fixture-project".into(),
price_sats: 8,
viewing_seconds: 3600,
created_at: 1000,
expires_at: 1600,
};
let selection = AuthorizedSelection {
relative_path: "Movies/Film.mp4".into(),
payment_methods: vec!["cashu".into()],
};
let event = EventBuilder::new(
Kind::Custom(KIND),
serde_json::to_string_pretty(&approval_content(&intent, &selection).unwrap()).unwrap(),
)
.tag(Tag::identifier(DOMAIN))
.custom_created_at(Timestamp::from(1200))
.sign_with_keys(&keys)
.unwrap();
Params {
intent,
selection,
producer_event: event,
}
}
#[test]
fn producer_signature_binds_human_readable_exact_selection_terms_node_and_installation() {
let original = fixture();
assert_eq!(
verified_producer(&original, 1300).unwrap(),
original.intent.producer
);
// Original consent can recover an already-completed operation after
// expiry; underlying snapshot journal refuses creating a fresh one.
assert!(verified_producer(&original, 2000).is_ok());
let mut changed = fixture();
changed.intent.price_sats += 1;
assert!(verified_producer(&changed, 1300).is_err());
let mut changed = fixture();
changed.selection.relative_path = "Other.mp4".into();
assert!(verified_producer(&changed, 1300).is_err());
let mut changed = fixture();
changed.intent.app_audience = uuid::Uuid::new_v4().to_string();
assert!(verified_producer(&changed, 1300).is_err());
let mut changed = fixture();
changed.intent.producer = "cd".repeat(32);
assert!(verified_producer(&changed, 1300).is_err());
assert!(verified_producer(&fixture(), 1000).is_err());
}
#[test]
fn request_parser_rejects_unsigned_claims_unknown_fields_and_changed_signed_content() {
let original = fixture();
let mut encoded = serde_json::json!({"intent":original.intent,"selection":original.selection,"producerEvent":original.producer_event});
assert!(parse(encoded.clone(), 1300).is_ok());
encoded["producerEvent"]["content"] = serde_json::json!("approve everything");
assert!(parse(encoded, 1300).is_err());
assert!(parse(serde_json::json!({"producer":"cd".repeat(32)}), 1300).is_err());
}
#[test]
fn dropped_request_signals_blocking_copy_cancellation() {
let (guard, signal) = request_cancellation();
assert!(!signal.load(Ordering::Relaxed));
drop(guard);
assert!(signal.load(Ordering::Relaxed));
}
#[test]
fn resolution_signature_recovers_only_exact_intent_after_expiry_without_file_authority() {
let original = fixture();
let keys = Keys::parse(&"07".repeat(32)).unwrap();
let event = EventBuilder::new(
Kind::Custom(27237),
serde_json::json!({
"action":"Recover prepared video or retire this expired incomplete registration",
"scope":RESOLUTION_DOMAIN,"intent":original.intent,
})
.to_string(),
)
.tag(Tag::identifier(RESOLUTION_DOMAIN))
.custom_created_at(Timestamp::from(1700))
.sign_with_keys(&keys)
.unwrap();
let encoded = serde_json::json!({"intent":original.intent,"producerEvent":event});
assert!(verified_resolution(encoded.clone(), 1800).is_ok());
assert!(verified_resolution(encoded.clone(), 9999).is_ok());
assert!(parse(encoded.clone(), 1800).is_err());
let mut changed = encoded.clone();
changed["intent"]["priceSats"] = serde_json::json!(999);
assert!(verified_resolution(changed, 1800).is_err());
let mut changed = encoded;
changed["selection"] =
serde_json::json!({"relative_path":"film.mp4","payment_methods":["cashu"]});
assert!(verified_resolution(changed, 1800).is_err());
assert!(verified_resolution(
serde_json::json!({"intent":original.intent,"producerEvent":original.producer_event}),
1800
)
.is_err());
}
}
+54 -6
View File
@@ -19,6 +19,9 @@ mod identity;
mod interfaces;
pub(crate) mod lnd;
mod marketplace;
mod media_registration;
mod purchase;
mod playback;
// 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
@@ -97,6 +100,22 @@ fn nostr_signing_origin_allowed(headers: &hyper::HeaderMap, dev_mode: bool) -> b
matches!(url.port_or_known_default(), Some(80 | 443))
}
fn native_consent_origin_allowed(method: &str, headers: &hyper::HeaderMap, dev_mode: bool) -> bool {
!matches!(
method,
"node.nostr-sign"
| "identity.nostr-sign"
| "media.registration.prepare"
| "media.registration.context"
| "media.registration.resolve"
| "content.rental-purchase"
| "content.purchase"
| "content.cancel-purchase"
| "content.playback-handle"
| "content.playback-status"
) || nostr_signing_origin_allowed(headers, dev_mode)
}
/// Read-only authenticated methods may skip CSRF, but they must still exist in
/// the dispatcher. The tab signer uses `system.get-hostname` as its lightweight
/// session probe, so keeping the policy in one testable function protects that
@@ -148,6 +167,7 @@ pub struct RpcHandler {
pub(crate) app_gate: Arc<crate::appgate::AppGate>,
endpoint_rate_limiter: EndpointRateLimiter,
response_cache: ResponseCache,
playback_handles: crate::playback_handles::PlaybackHandles,
mesh_service: Arc<tokio::sync::RwLock<Option<crate::mesh::MeshService>>>,
/// LoRa radio firmware-flash job state, sibling to `mesh_service` — one
/// job at a time, since flashing needs exclusive access to the port.
@@ -172,6 +192,10 @@ pub struct RpcHandler {
}
impl RpcHandler {
pub(crate) fn playback_handles(&self) -> &crate::playback_handles::PlaybackHandles {
&self.playback_handles
}
pub async fn new(
config: Config,
state_manager: Arc<StateManager>,
@@ -231,6 +255,7 @@ impl RpcHandler {
app_gate,
endpoint_rate_limiter,
response_cache: ResponseCache::new(5),
playback_handles: Default::default(),
mesh_service: Arc::new(tokio::sync::RwLock::new(None)),
flash_job: crate::mesh::flash::new_job_handle(),
transport_router: Arc::new(tokio::sync::RwLock::new(None)),
@@ -333,14 +358,10 @@ impl RpcHandler {
debug!("RPC method: {}", rpc_req.method);
if matches!(
rpc_req.method.as_str(),
"node.nostr-sign" | "identity.nostr-sign"
) && !nostr_signing_origin_allowed(&parts.headers, self.config.dev_mode)
{
if !native_consent_origin_allowed(&rpc_req.method, &parts.headers, self.config.dev_mode) {
return Ok(self.error_response(
403,
"Nostr signing from app origins requires the dashboard consent bridge",
"Native signing and Cloud registration from app origins require the dashboard consent bridge",
StatusCode::FORBIDDEN,
));
}
@@ -799,6 +820,33 @@ mod nostr_signing_origin_tests {
headers
}
#[test]
fn native_registration_and_purchase_use_dashboard_origin_and_keep_authentication_and_csrf() {
for method in [
"media.registration.prepare",
"media.registration.context",
"media.registration.resolve",
"content.rental-purchase",
"content.purchase",
"content.cancel-purchase",
"content.playback-handle",
"content.playback-status",
] {
assert!(!native_consent_origin_allowed(
method,
&headers(Some("http://node.local:7778")),
false
));
assert!(native_consent_origin_allowed(
method,
&headers(Some("https://node.local")),
false
));
assert!(!UNAUTHENTICATED_METHODS.contains(&method));
assert!(!csrf_exempt_method(method));
}
}
#[test]
fn signing_accepts_dashboard_and_authenticated_non_browser_clients() {
assert!(nostr_signing_origin_allowed(&headers(None), false));
+74
View File
@@ -0,0 +1,74 @@
//! Native broker only: settled purchases become session-bound opaque media URLs.
use super::RpcHandler;
use anyhow::{Context, Result};
use serde::Deserialize;
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct Issue {
purchase_id: String,
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct Status {
handle: String,
}
impl RpcHandler {
async fn playback_context(
&self,
session: &Option<String>,
) -> Result<(
String,
crate::container::registration_pin::InstalledAppContext,
)> {
let session = session.as_ref().context("Owner session required")?;
anyhow::ensure!(
self.session_store.validate(session).await,
"Owner session expired"
);
let identity =
crate::identity::NodeIdentity::load_existing(&self.config.data_dir.join("identity"))
.await?;
let (state, _) = self.state_manager.get_snapshot().await;
let root = self.config.data_dir.clone();
let context = tokio::task::spawn_blocking(move || {
crate::container::registration_pin::installed_context(&root, &identity, &state)
})
.await??;
Ok((session.clone(), context))
}
pub(super) async fn handle_playback_handle(
&self,
params: Option<serde_json::Value>,
session: &Option<String>,
) -> Result<serde_json::Value> {
let params: Issue = serde_json::from_value(params.context("Missing purchase identifier")?)?;
let (session, context) = self.playback_context(session).await?;
let handle = self
.playback_handles()
.issue(
&self.config.data_dir,
context.clone(),
&session,
&params.purchase_id,
)
.await?;
let expires = self
.playback_handles()
.expiry(&handle, &session, &context)?;
Ok(
serde_json::json!({"playback_url":format!("/api/rental-playback/{handle}"), "expires_at":expires}),
)
}
pub(super) async fn handle_playback_status(
&self,
params: Option<serde_json::Value>,
session: &Option<String>,
) -> Result<serde_json::Value> {
let params: Status = serde_json::from_value(params.context("Missing playback handle")?)?;
let (session, context) = self.playback_context(session).await?;
let expires = self
.playback_handles()
.expiry(&params.handle, &session, &context)?;
Ok(serde_json::json!({"expires_at":expires}))
}
}
+216
View File
@@ -0,0 +1,216 @@
//! Owner RPC for permanent Cloud purchases. Registered app rentals use the
//! separate native installed-app context + opaque playback-handle dispatcher.
use super::RpcHandler;
use crate::{
content_purchase::Journal,
content_purchase_caller::{self as caller, PurchaseConsent, PurchaseTransport, ReadyPurchase},
content_purchase_transport::FipsPurchaseTransport,
};
use anyhow::{Context, Result};
use serde::Deserialize;
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct PurchaseParams {
onion: String,
content_id: String,
filename: Option<String>,
max_wallet_debit: u64,
consent: Option<PurchaseConsent>,
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct CancelParams {
onion: String,
operation_id: String,
}
impl RpcHandler {
pub(super) async fn handle_content_purchase(
&self,
params: Option<serde_json::Value>,
) -> Result<serde_json::Value> {
let params: PurchaseParams =
serde_json::from_value(params.context("Missing purchase parameters")?)?;
anyhow::ensure!(
!params.content_id.starts_with("registered_"),
"Registered rentals require the native app purchase context"
);
let transport =
FipsPurchaseTransport::load(self.config.data_dir.clone(), params.onion.clone()).await?;
let identity =
crate::identity::NodeIdentity::load_existing(&self.config.data_dir.join("identity"))
.await?;
let buyer = identity.did_key()?;
let result = caller::purchase(
&self.config.data_dir,
&buyer,
&params.content_id,
params.filename.as_deref(),
params.max_wallet_debit,
params.consent.as_ref(),
&transport,
)
.await?;
match result {
ReadyPurchase::AwaitingConfirmation {
operation_id,
envelope_sha256,
gross_token_sats,
seller_net_sats,
wallet_debit_sats,
expires_at,
network,
mint_url,
} => Ok(
serde_json::json!({"state":"confirmation_required","operation_id":operation_id,
"envelope_sha256":envelope_sha256,"gross_token_sats":gross_token_sats,
"seller_net_sats":seller_net_sats,"wallet_debit_sats":wallet_debit_sats,"expires_at":expires_at,"network":network,"mint_url":mint_url}),
),
ReadyPurchase::Cancelled { operation_id } => {
Ok(serde_json::json!({"state":"cancelled_unspent","operation_id":operation_id}))
}
ReadyPurchase::Cached { content_id, .. } => Ok(
serde_json::json!({"state":"delivered","owned":true,"owned_content_id":content_id}),
),
ReadyPurchase::Entitlement { contract, receipt } => {
let item = crate::content_purchase_download::cache(
&self.config.data_dir,
&params.onion,
&contract,
&receipt,
)
.await?;
Ok(
serde_json::json!({"state":"delivered","owned":true,"operation_id":contract.id,
"owned_content_id":item.content_id,"mime_type":item.mime_type,"size_bytes":item.size_bytes}),
)
}
}
}
pub(super) async fn handle_content_cancel_purchase(
&self,
params: Option<serde_json::Value>,
) -> Result<serde_json::Value> {
let params: CancelParams =
serde_json::from_value(params.context("Missing cancellation parameters")?)?;
let transport =
FipsPurchaseTransport::load(self.config.data_dir.clone(), params.onion).await?;
let identity =
crate::identity::NodeIdentity::load_existing(&self.config.data_dir.join("identity"))
.await?;
let (envelope, plan) = {
let journal = Journal::open(&self.config.data_dir).await?;
let record = journal
.buyer(&params.operation_id)
.await?
.context("Original purchase not found")?;
anyhow::ensure!(
record.contract.buyer_did == identity.did_key()?
&& record.contract.seller_did == transport.seller_did(),
"Cancellation purchase binding changed"
);
(
journal
.protocol_envelope("buyer", &params.operation_id)
.await?
.context("Original payment shape missing")?,
journal
.buyer_plan(&params.operation_id)
.await?
.context("Original wallet plan missing")?,
)
};
caller::cancel_purchase(&self.config.data_dir, &envelope, &plan, &transport).await?;
Ok(serde_json::json!({"state":"cancelled_unspent","operation_id":params.operation_id}))
}
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct RentalParams {
seller_did: String,
content_id: String,
expected_sha256: String,
expected_price_sats: u64,
expected_viewing_seconds: u64,
max_wallet_debit: u64,
consent: Option<PurchaseConsent>,
}
impl RpcHandler {
pub(super) async fn handle_content_rental_purchase(
&self,
input: Option<serde_json::Value>,
) -> Result<serde_json::Value> {
let params: RentalParams =
serde_json::from_value(input.context("Missing rental purchase terms")?)?;
anyhow::ensure!(
params.content_id.starts_with("registered_")
&& params.expected_price_sats > 0
&& (1..=31_536_000).contains(&params.expected_viewing_seconds)
&& params.expected_sha256.len() == 64
&& params
.expected_sha256
.bytes()
.all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b)),
"Invalid published rental terms"
);
let identity =
crate::identity::NodeIdentity::load_existing(&self.config.data_dir.join("identity"))
.await?;
let buyer = identity.did_key()?;
let (state, _) = self.state_manager.get_snapshot().await;
let data = self.config.data_dir.clone();
tokio::task::spawn_blocking(move || {
crate::container::registration_pin::installed_context(&data, &identity, &state)
})
.await??;
let onion = crate::content_purchase_transport::seller_onion_for_did(
&self.config.data_dir,
&params.seller_did,
)
.await?;
let transport =
FipsPurchaseTransport::load(self.config.data_dir.clone(), onion.clone()).await?;
let expected = caller::ExpectedRental {
seller_did: params.seller_did,
content_id: params.content_id.clone(),
sha256: params.expected_sha256,
price_sats: params.expected_price_sats,
viewing_seconds: params.expected_viewing_seconds,
};
match caller::purchase_bound(
&self.config.data_dir,
&buyer,
&params.content_id,
None,
params.max_wallet_debit,
params.consent.as_ref(),
&transport,
Some(&expected),
)
.await?
{
ReadyPurchase::AwaitingConfirmation {
operation_id,
envelope_sha256,
gross_token_sats,
seller_net_sats,
wallet_debit_sats,
expires_at,
network,
mint_url,
} => Ok(serde_json::json!({
"state":"confirmation_required","operation_id":operation_id,"envelope_sha256":envelope_sha256,
"gross_token_sats":gross_token_sats,"seller_net_sats":seller_net_sats,"wallet_debit_sats":wallet_debit_sats,
"expires_at":expires_at,"seller_onion":onion,"network":network,"mint_url":mint_url})),
ReadyPurchase::Entitlement { contract, .. } => Ok(
serde_json::json!({"state":"entitled","operation_id":contract.id,"seller_onion":onion}),
),
ReadyPurchase::Cancelled { operation_id } => Ok(
serde_json::json!({"state":"cancelled_unspent","operation_id":operation_id,"seller_onion":onion}),
),
ReadyPurchase::Cached { .. } => {
anyhow::bail!("A timed rental cannot use a permanent owned copy")
}
}
}
}