diff --git a/core/Cargo.lock b/core/Cargo.lock index 845eba7c..b0d6aa53 100644 --- a/core/Cargo.lock +++ b/core/Cargo.lock @@ -141,6 +141,7 @@ dependencies = [ "iroh", "iroh-blobs", "libc", + "lightning-invoice", "lofty", "mainline", "mdns-sd", diff --git a/core/archipelago/Cargo.toml b/core/archipelago/Cargo.toml index 4a2e974a..af13a0cc 100644 --- a/core/archipelago/Cargo.toml +++ b/core/archipelago/Cargo.toml @@ -73,6 +73,7 @@ chrono = "0.4" # BIP-39 mnemonic seed generation + BIP-32 HD key derivation bip39 = { version = "2.1", features = ["rand"] } +lightning-invoice = "=0.34.1" bitcoin = { version = "=0.32.5", features = ["rand-std"] } # Configuration diff --git a/core/archipelago/src/api/handler/lightning_purchase.rs b/core/archipelago/src/api/handler/lightning_purchase.rs new file mode 100644 index 00000000..36d55e0a --- /dev/null +++ b/core/archipelago/src/api/handler/lightning_purchase.rs @@ -0,0 +1,241 @@ +use super::{build_response, ApiHandler}; +use crate::content_lightning::{Binding, Journal, Phase}; +use anyhow::{Context, Result}; +use hyper::{body::HttpBody, Body, Method, Request, Response, StatusCode}; +use serde::{Deserialize, Serialize}; +use tokio::io::AsyncReadExt; +pub(crate) const ROUTE: &str = "/content/lightning/v1/operation"; +#[derive(Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(crate) struct Operation { + pub binding: Binding, + pub action: String, +} +impl ApiHandler { + pub(super) async fn handle_lightning_purchase( + &self, + mut request: Request, + ) -> Result> { + anyhow::ensure!( + request.method() == Method::POST && request.uri().path() == ROUTE, + "Invalid invoice route" + ); + 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() + chunk.len() <= 16384, + "Invoice request too large" + ); + bytes.extend_from_slice(&chunk) + } + Ok::<_, anyhow::Error>(bytes) + }) + .await + .context("Invoice request timed out")??; + let seller = crate::identity::did_key_from_pubkey_hex(&self.self_pubkey_hex)?; + let buyer = crate::content_auth::authenticate_request( + request.headers(), + &seller, + &Method::POST, + ROUTE, + &bytes, + chrono::Utc::now().timestamp(), + )?; + let operation: Operation = serde_json::from_slice(&bytes)?; + anyhow::ensure!( + operation.binding.buyer_did == buyer && operation.binding.seller_did == seller, + "Invoice peer identity mismatch" + ); + anyhow::ensure!( + matches!( + operation.action.as_str(), + "create" | "status" | "cancel" | "download" + ), + "Invalid invoice action" + ); + let binding = &operation.binding; + let journal = Journal::open(&self.config.data_dir).await?; + let mut saved = journal.seller(binding)?; + if saved.is_none() { + anyhow::ensure!( + operation.action == "create", + "Unknown original invoice operation" + ); + anyhow::ensure!( + !binding.content_id.starts_with("registered_"), + "Registered rentals use their native purchase contract" + ); + let catalog = crate::content_server::load_catalog(&self.config.data_dir).await?; + let item = catalog + .items + .iter() + .find(|v| v.id == binding.content_id) + .context("Shared item unavailable")?; + let visible = match &item.availability { + crate::content_server::Availability::Nobody => false, + crate::content_server::Availability::AllPeers => true, + crate::content_server::Availability::Specific { peers } => peers.contains(&buyer), + }; + anyhow::ensure!(visible, "Item is not shared with this buyer"); + anyhow::ensure!( + matches!(&item.access,crate::content_server::AccessControl::Paid{price_sats,..} if *price_sats==binding.price_sats) + && crate::content_server::method_accepted(&item.access, "lightning"), + "Invoice price or accepted method changed" + ); + crate::content_server::ensure_payment_source_available(&self.config.data_dir, item) + .await?; + let source = crate::content_server::content_file_path(&self.config.data_dir, item); + let roots = [ + self.config.data_dir.join("content/files"), + self.config.data_dir.join("filebrowser"), + ]; + let (root, relative) = roots + .iter() + .find_map(|root| { + source + .strip_prefix(root) + .ok() + .map(|p| (root.clone(), p.to_path_buf())) + }) + .context("Unsupported invoice source root")?; + let data = self.config.data_dir.clone(); + let id = binding.content_id.clone(); + struct CancelCopy(std::sync::Arc); + impl Drop for CancelCopy { + fn drop(&mut self) { + self.0.store(true, std::sync::atomic::Ordering::SeqCst); + } + } + let cancel_copy = CancelCopy(std::sync::Arc::new(std::sync::atomic::AtomicBool::new( + false, + ))); + let cancelled = cancel_copy.0.clone(); + let snapshot = tokio::task::spawn_blocking(move || { + crate::content_snapshot::prepare( + &data, + &root, + &id, + &relative, + &crate::media_registration::Limits { + max_bytes: 64 * 1024 * 1024 * 1024, + cancelled: &cancelled, + }, + 64 * 1024 * 1024 * 1024, + 512 * 1024 * 1024, + |_| Ok(()), + ) + }) + .await??; + anyhow::ensure!( + snapshot.size == item.size_bytes, + "Shared file changed before invoice" + ); + // Source metadata is private and committed before AddInvoice dispatch. + let record = crate::content_server::publish_snapshot_invoice( + &self.config.data_dir, + item, + &journal, + binding.clone(), + crate::content_lightning::RetainedFile { + sha256: snapshot.sha256, + size: snapshot.size, + filename: item.filename.clone(), + mime_type: item.mime_type.clone(), + }, + ) + .await?; + saved = Some(record); + } + let mut saved = saved.context("Missing invoice operation")?; + anyhow::ensure!( + saved.source.is_some(), + "Original invoice source is not prepared; no new invoice dispatched" + ); + let status = if operation.action == "cancel" && saved.phase == Phase::Prepared { + saved.phase = Phase::CanceledUnpaid; + journal.save_seller(&saved)?; + saved.status() + } else if operation.action != "create" + && operation.action != "cancel" + && saved.phase == Phase::Prepared + { + saved.status() + } else { + self.rpc_handler + .drive_external_invoice(&journal, binding, operation.action == "cancel") + .await? + }; + // The original legacy delivery mechanism remains usable by its hash. + if status.bolt11.is_some() { + crate::content_invoice::record_pending( + &self.config.data_dir, + &status.payment_hash, + &binding.content_id, + binding.price_sats, + ) + .await?; + if status.state == Phase::Settled { + crate::content_invoice::mark_paid(&self.config.data_dir, &status.payment_hash) + .await?; + } + } + if operation.action == "download" { + anyhow::ensure!( + status.state == Phase::Settled, + "Original invoice has not settled" + ); + let source = status + .source + .as_ref() + .context("Original invoice snapshot is missing")?; + let data = self.config.data_dir.clone(); + let id = binding.content_id.clone(); + let retained = source.clone(); + struct CancelCopy(std::sync::Arc); + impl Drop for CancelCopy { + fn drop(&mut self) { + self.0.store(true, std::sync::atomic::Ordering::SeqCst); + } + } + let cancel_copy = CancelCopy(std::sync::Arc::new(std::sync::atomic::AtomicBool::new( + false, + ))); + let cancelled = cancel_copy.0.clone(); + let snapshot = tokio::task::spawn_blocking(move || { + crate::content_snapshot::open_matching(&data, &id, &retained.sha256, retained.size) + }) + .await??; + let stream = futures_util::stream::try_unfold( + (tokio::fs::File::from_std(snapshot.file), source.size), + |(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, + "Original invoice snapshot ended early", + )); + } + bytes.truncate(count); + Ok(Some((bytes, (file, left - count as u64)))) + }, + ); + return Ok(Response::builder() + .status(StatusCode::OK) + .header("Content-Type", &source.mime_type) + .header("Content-Length", source.size) + .header("Cache-Control", "private, no-store") + .body(Body::wrap_stream(stream))?); + } + Ok(build_response( + StatusCode::OK, + "application/json", + Body::from(serde_json::to_vec(&status)?), + )) + } +} diff --git a/core/archipelago/src/api/handler/mod.rs b/core/archipelago/src/api/handler/mod.rs index 77d7b8c1..cda073f5 100644 --- a/core/archipelago/src/api/handler/mod.rs +++ b/core/archipelago/src/api/handler/mod.rs @@ -3,6 +3,7 @@ mod cdp; mod cloud_purchase; mod content; mod dwn; +pub(crate) mod lightning_purchase; mod model_proxy; mod node_message; mod proxy; @@ -454,6 +455,9 @@ impl ApiHandler { .await; } + if method == Method::POST && path == lightning_purchase::ROUTE { + return self.handle_lightning_purchase(req).await; + } // Purchase routes bound the original body before the generic buffer. if method == Method::POST && matches!( diff --git a/core/archipelago/src/api/rpc/dispatcher.rs b/core/archipelago/src/api/rpc/dispatcher.rs index f7f227f5..1cb6de68 100644 --- a/core/archipelago/src/api/rpc/dispatcher.rs +++ b/core/archipelago/src/api/rpc/dispatcher.rs @@ -337,6 +337,15 @@ impl RpcHandler { "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.invoice-pay" => self.handle_lightning_operation(params, "pay").await, + "content.invoice-download" => self.handle_lightning_operation(params, "download").await, + "content.invoice-attempt" => self.handle_lightning_operation(params, "lookup").await, + "content.invoice-retry-native" => { + self.handle_lightning_operation(params, "retry").await + } + "content.invoice-create" => self.handle_lightning_operation(params, "create").await, + "content.invoice-recover" => self.handle_lightning_operation(params, "status").await, + "content.invoice-cancel" => self.handle_lightning_operation(params, "cancel").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, diff --git a/core/archipelago/src/api/rpc/lightning_purchase.rs b/core/archipelago/src/api/rpc/lightning_purchase.rs new file mode 100644 index 00000000..6d768478 --- /dev/null +++ b/core/archipelago/src/api/rpc/lightning_purchase.rs @@ -0,0 +1,356 @@ +use super::RpcHandler; +use crate::{ + api::handler::lightning_purchase::{Operation, ROUTE}, + content_lightning::{Binding, BuyerRecord, Journal, Phase, Status}, +}; +use anyhow::{Context, Result}; +use serde::Deserialize; +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct Params { + onion: String, + content_id: String, + price_sats: Option, + operation_id: Option, + #[serde(default)] + external_exposure: bool, +} +impl RpcHandler { + pub(super) async fn handle_lightning_operation( + &self, + params: Option, + action: &str, + ) -> Result { + let params: Params = serde_json::from_value(params.context("Missing invoice operation")?)?; + let peer = + crate::federation::load_unique_payment_peer(&self.config.data_dir, ¶ms.onion) + .await?; + let fips = peer + .fips_npub + .context("Seller has no authenticated mesh connection")?; + let buyer = + crate::identity::NodeIdentity::load_existing(&self.config.data_dir.join("identity")) + .await? + .did_key()?; + anyhow::ensure!(buyer != peer.did, "Cannot buy a file from this same node"); + let _admission = crate::content_payment_admission::lock( + &self.config.data_dir, + &buyer, + &peer.did, + ¶ms.content_id, + ) + .await?; + if matches!(action, "create" | "pay" | "retry") || params.external_exposure { + let cashu = crate::content_purchase::Journal::open(&self.config.data_dir).await?; + anyhow::ensure!( + cashu + .find_buyers(&buyer, &peer.did, ¶ms.content_id) + .await? + .iter() + .all(|r| r.phase == crate::content_purchase::BuyerPhase::Cancelled), + "Recover or cancel the original Cashu purchase before exposing a Lightning invoice" + ); + } + let journal = Journal::open(&self.config.data_dir).await?; + let original = if let Some(id) = ¶ms.operation_id { + journal.buyer(id)? + } else { + journal.buyer_for(&buyer, &peer.did, ¶ms.content_id)? + }; + if let Some(record) = &original { + anyhow::ensure!( + record.binding.buyer_did == buyer + && record.binding.seller_did == peer.did + && record.binding.content_id == params.content_id + && record.seller_onion == params.onion, + "Original invoice belongs to another purchase" + ); + } + if action == "lookup" { + return Ok(match original { + None => serde_json::json!({"attempt":null}), + Some(mut record) => { + let mut native_result = record.native_result.clone(); + if native_result.is_none() && record.native_dispatched { + if let Some(status) = &record.last { + if let Ok(payment) = self + .handle_lnd_paymentstatus(Some( + serde_json::json!({"payment_hash":status.payment_hash}), + )) + .await + { + if let Some(result @ ("failed" | "succeeded")) = + payment["status"].as_str() + { + native_result = Some(result.to_owned()); + record.native_result = native_result.clone(); + journal.save_buyer(&record)?; + } + } + } + } + let native_failed = + !record.external_exposure && native_result.as_deref() == Some("failed"); + let native_succeeded = native_result.as_deref() == Some("succeeded"); + if !record.external_exposure { + if let Some(status) = record.last.as_mut() { + status.bolt11 = None; + } + } + serde_json::json!({"attempt":{"operation_id":record.binding.id,"price_sats":record.binding.price_sats,"external_exposure":record.external_exposure,"native_failed":native_failed,"native_succeeded":native_succeeded,"status":record.last}}) + } + }); + } + let mut record = if let Some(record) = original { + record + } else { + anyhow::ensure!( + action == "create" && params.operation_id.is_none(), + "Original invoice operation is unavailable" + ); + BuyerRecord { + binding: Binding { + id: uuid::Uuid::new_v4().to_string(), + buyer_did: buyer.clone(), + seller_did: peer.did.clone(), + content_id: params.content_id.clone(), + price_sats: params + .price_sats + .context("Expected invoice price is required")?, + }, + seller_onion: params.onion.clone(), + external_exposure: false, + native_retired: false, + native_replacement: None, + native_dispatched: false, + native_result: None, + last: None, + } + }; + anyhow::ensure!( + record.binding.buyer_did == buyer + && record.binding.seller_did == peer.did + && record.binding.content_id == params.content_id + && record.seller_onion == params.onion + && params + .operation_id + .as_ref() + .is_none_or(|id| id == &record.binding.id), + "Original invoice operation changed" + ); + if action == "retry" { + anyhow::ensure!( + params.operation_id.is_some() + && !params.external_exposure + && params.price_sats == Some(record.binding.price_sats), + "Explicit original native retry and original price required" + ); + if record.native_result.is_none() + && record.native_dispatched + && !record.external_exposure + { + let hash = &record + .last + .as_ref() + .context("Original invoice metadata missing")? + .payment_hash; + let payment = self + .handle_lnd_paymentstatus(Some(serde_json::json!({"payment_hash":hash}))) + .await?; + if matches!(payment["status"].as_str(), Some("failed" | "succeeded")) { + record.native_result = payment["status"].as_str().map(str::to_owned); + journal.save_buyer(&record)?; + } + } + record = journal.retry_native(&record.binding.id)?; + } + anyhow::ensure!((!record.native_retired || matches!(action,"status"|"cancel"|"download")) && (!record.native_retired || !params.external_exposure),"This native invoice was retired before changing payment method; recover the replacement purchase"); + if action == "pay" { + anyhow::ensure!( + params.operation_id.is_some(), + "Original invoice operation required for native payment" + ); + return crate::content_lightning::drive_native( + &self.config.data_dir, + journal, + &record.binding.id, + &super::lnd::external_invoice::NativeNode(self), + ) + .await; + } + record.external_exposure |= params.external_exposure; + journal.save_buyer(&record)?; + drop(journal); + // Native-only local FAILED never cancels an externally exposed invoice. + // The seller terminal state is authoritative regardless of UI receipt loss. + let operation = Operation { + binding: record.binding.clone(), + action: if action == "retry" { + "create".into() + } else { + action.into() + }, + }; + let response = crate::fips::dial::PeerRequest::new(Some(&fips), ¶ms.onion, ROUTE) + .require_fips() + .single_delivery() + .timeout(std::time::Duration::from_secs(if action == "download" { + 900 + } else { + 45 + })) + .send_content_json(&self.config.data_dir, &peer.did, &operation) + .await; + let mut response = match response { + Ok((r, _)) => r, + Err(_) => { + return Ok( + serde_json::json!({"state":"unknown","operation_id":record.binding.id,"recovery_required":true,"error":"The original invoice request is saved on this node. Recover it; no replacement invoice was requested."}), + ) + } + }; + anyhow::ensure!( + response.status().is_success(), + "Seller could not resolve original invoice; recover operation {}", + record.binding.id + ); + let journal = Journal::open(&self.config.data_dir).await?; + if action == "download" { + let mut paid = record + .last + .clone() + .context("Original invoice metadata missing; recover it first")?; + let source = paid + .source + .clone() + .context("Original invoice snapshot missing")?; + anyhow::ensure!( + response.content_length() == Some(source.size), + "Original invoice file length changed" + ); + paid.state = Phase::Settled; + paid.can_switch_method = false; + record.last = Some(paid); + journal.save_buyer(&record)?; + drop(journal); + let stream = crate::content_purchase_download::verified_stream( + response.bytes_stream(), + source.sha256, + source.size, + ); + let owned = crate::content_owned::record_purchase_stream( + &self.config.data_dir, + crate::content_owned::OwnedItem { + onion: params.onion, + content_id: params.content_id, + filename: source.filename, + mime_type: source.mime_type, + size_bytes: source.size, + paid_sats: record.binding.price_sats, + ecash_backend: "lightning".into(), + purchased_at: chrono::Utc::now().to_rfc3339(), + download_complete: false, + }, + Box::pin(stream), + Some(source.size), + ) + .await?; + return Ok( + serde_json::json!({"owned":true,"owned_content_id":owned.content_id,"mime_type":owned.mime_type}), + ); + } + let mut bytes = Vec::new(); + while let Some(chunk) = response.chunk().await? { + anyhow::ensure!( + bytes.len() + chunk.len() <= 16384, + "Invoice response too large" + ); + bytes.extend_from_slice(&chunk) + } + let status: Status = serde_json::from_slice(&bytes)?; + anyhow::ensure!( + status.binding == record.binding + && status.source.is_some() + && status.payment_hash.len() == 64 + && status.payment_hash.bytes().all(|b| b.is_ascii_hexdigit()) + && status.can_switch_method == (status.state == Phase::CanceledUnpaid), + "Seller invoice binding changed" + ); + if let Some(bolt11) = &status.bolt11 { + let invoice: lightning_invoice::Bolt11Invoice = + bolt11.parse().context("Seller invoice is invalid")?; + invoice.check_signature()?; + anyhow::ensure!( + invoice.payment_hash().to_string() == status.payment_hash + && invoice.amount_milli_satoshis() + == record.binding.price_sats.checked_mul(1000), + "Invoice hash or amount differs from saved purchase" + ); + } + if let Some(previous) = &record.last { + anyhow::ensure!( + previous.payment_hash == status.payment_hash + && previous.source == status.source + && previous + .bolt11 + .as_ref() + .is_none_or(|v| status.bolt11.as_ref() == Some(v)), + "Original invoice replaced" + ); + } + record.last = Some(status.clone()); + journal.save_buyer(&record)?; + Ok( + serde_json::json!({"operation_id":record.binding.id,"price_sats":record.binding.price_sats,"payment_hash":status.payment_hash,"bolt11":if record.external_exposure{status.bolt11}else{None},"state":match status.state{Phase::Settled=>"settled",Phase::CanceledUnpaid=>"canceled",Phase::Issued=>"open",Phase::Prepared=>"prepared",Phase::Dispatched=>"unknown",Phase::CancelRequested=>"cancel_requested"},"paid":status.state==Phase::Settled,"can_switch_method":status.can_switch_method,"cancel_supported":true,"external_exposure":record.external_exposure}), + ) + } +} + +impl RpcHandler { + /// Caller holds content_payment_admission before entering any rail journal. + pub(super) async fn ensure_invoice_allows_other_rail( + &self, + buyer: &str, + seller: &str, + content: &str, + ) -> Result<()> { + let journal = Journal::open(&self.config.data_dir).await?; + if let Some(mut record) = journal.buyer_for(buyer, seller, content)? { + anyhow::ensure!(record.native_replacement.is_none(), + "An explicit native retry is being recovered; recover its replacement operation first"); + anyhow::ensure!( + !record.external_exposure, + "An externally payable invoice remains unresolved; cancel or recover it first" + ); + let status = record + .last + .clone() + .context("Original invoice creation is unresolved; recover it first")?; + anyhow::ensure!( + status.state != Phase::Settled, + "Original Lightning purchase is paid; recover its file" + ); + // A local terminal failure can release only a never-exposed native + // attempt. This check runs under the same outer lock as QR exposure. + anyhow::ensure!( + record.native_dispatched, + "Original invoice has not been canceled; cancel it before replacing the method" + ); + if record.native_result.as_deref() != Some("failed") { + let payment = self + .handle_lnd_paymentstatus(Some( + serde_json::json!({"payment_hash":status.payment_hash}), + )) + .await?; + anyhow::ensure!( + payment["status"] == "failed", + "Original native Lightning attempt remains unresolved" + ); + } + record.native_result = Some("failed".into()); + record.native_retired = true; + journal.save_buyer(&record)?; + } + Ok(()) + } +} diff --git a/core/archipelago/src/api/rpc/lnd/external_invoice.rs b/core/archipelago/src/api/rpc/lnd/external_invoice.rs new file mode 100644 index 00000000..b142a19a --- /dev/null +++ b/core/archipelago/src/api/rpc/lnd/external_invoice.rs @@ -0,0 +1,201 @@ +use super::LND_REST_BASE_URL; +use crate::{ + api::rpc::RpcHandler, + content_lightning::{Binding, Invoice, InvoiceNode, Journal, Status}, +}; +use anyhow::{Context, Result}; +use base64::Engine; +struct Node { + client: reqwest::Client, + macaroon: String, +} +fn number(v: &serde_json::Value) -> Option { + v.as_u64().or_else(|| v.as_str()?.parse().ok()) +} +impl InvoiceNode for Node { + async fn prepare_creation(&self) -> Result<()> { + let info: serde_json::Value = self + .client + .get(format!("{LND_REST_BASE_URL}/v1/getinfo")) + .header("Grpc-Metadata-macaroon", &self.macaroon) + .send() + .await? + .error_for_status()? + .json() + .await?; + anyhow::ensure!( + info["identity_pubkey"] + .as_str() + .is_some_and(|key| !key.is_empty()), + "LND invoice service is not ready; original preparation retained" + ); + Ok(()) + } + async fn lookup(&self, hash: &str) -> Result> { + let response = self + .client + .get(format!("{LND_REST_BASE_URL}/v1/invoice/{hash}")) + .header("Grpc-Metadata-macaroon", &self.macaroon) + .send() + .await?; + if response.status() == reqwest::StatusCode::NOT_FOUND { + return Ok(None); + } + let body: serde_json::Value = response.error_for_status()?.json().await?; + let raw = body["r_hash"].as_str().context("Invoice hash omitted")?; + let payment_hash = hex::encode(base64::engine::general_purpose::STANDARD.decode(raw)?); + Ok(Some(Invoice { + payment_hash, + bolt11: body["payment_request"] + .as_str() + .context("Invoice payment request omitted")? + .into(), + price_sats: number(&body["value"]).context("Invoice amount omitted")?, + state: body["state"] + .as_str() + .context("Invoice state omitted")? + .into(), + paid_sats: number(&body["amt_paid_sat"]), + paid_msats: number(&body["amt_paid_msat"]), + })) + } + async fn add(&self, binding: &Binding, preimage_hex: &str) -> Result<()> { + let preimage = base64::engine::general_purpose::STANDARD.encode(hex::decode(preimage_hex)?); + self.client.post(format!("{LND_REST_BASE_URL}/v1/invoices")).header("Grpc-Metadata-macaroon",&self.macaroon) + .json(&serde_json::json!({"memo":format!("Archipelago peer file {}",binding.content_id),"value":binding.price_sats.to_string(),"r_preimage":preimage,"private":true,"expiry":"3600"})).send().await?.error_for_status()?; + Ok(()) + } + async fn cancel(&self, hash: &str) -> Result<()> { + self.client.post(format!("{LND_REST_BASE_URL}/v2/invoices/cancel")).header("Grpc-Metadata-macaroon",&self.macaroon) + .json(&serde_json::json!({"payment_hash":base64::engine::general_purpose::STANDARD.encode(hex::decode(hash)?)})).send().await?.error_for_status()?; + Ok(()) + } +} +impl RpcHandler { + pub(crate) async fn drive_external_invoice( + &self, + journal: &Journal, + binding: &Binding, + cancel: bool, + ) -> Result { + // Configuration/auth preflight before the engine persists dispatch. + let (client, macaroon) = self.lnd_client().await?; + crate::content_lightning::drive(journal, binding, &Node { client, macaroon }, cancel).await + } +} + +/// Prepared before the durable native-dispatch marker. Once execute is called, +/// every transport error is ambiguous and only original-hash lookup may follow. +pub(crate) struct PreparedNativePayment { + client: reqwest::Client, + request: reqwest::Request, + hash: String, + amount: u64, +} +impl PreparedNativePayment { + pub(crate) async fn execute(self) -> Result { + let response = self + .client + .execute(self.request) + .await + .context("Native payment response is unknown; recover the original operation")?; + let status = response.status(); + let body: serde_json::Value = response + .json() + .await + .context("Native payment response is unknown")?; + anyhow::ensure!( + status.is_success(), + "LND did not confirm the original payment; recover its status" + ); + let payment = body.get("result").unwrap_or(&body); + anyhow::ensure!( + payment + .get("payment_hash") + .and_then(|v| v.as_str()) + .is_none_or(|hash| hash == self.hash), + "LND payment hash changed" + ); + Ok(super::payments::router_payment_outcome( + payment, + &self.hash, + self.amount as i64, + )) + } +} +impl RpcHandler { + pub(crate) async fn prepare_bound_invoice_payment( + &self, + bolt11: &str, + hash: &str, + amount: u64, + ) -> Result { + let invoice: lightning_invoice::Bolt11Invoice = + bolt11.parse().context("Invalid original invoice")?; + invoice.check_signature()?; + anyhow::ensure!( + invoice.payment_hash().to_string() == hash + && invoice.amount_milli_satoshis() == amount.checked_mul(1000), + "Original invoice amount/hash changed" + ); + anyhow::ensure!( + !invoice.is_expired(), + "Original invoice expired; cancel or recover it before choosing another method" + ); + let (client, macaroon) = self.lnd_client().await?; + let info: serde_json::Value = client + .get(format!("{LND_REST_BASE_URL}/v1/getinfo")) + .header("Grpc-Metadata-macaroon", &macaroon) + .send() + .await? + .error_for_status()? + .json() + .await?; + let network = match invoice.currency() { + lightning_invoice::Currency::Bitcoin => "mainnet", + lightning_invoice::Currency::BitcoinTestnet => "testnet", + lightning_invoice::Currency::Regtest => "regtest", + lightning_invoice::Currency::Signet => "signet", + lightning_invoice::Currency::Simnet => "simnet", + }; + anyhow::ensure!( + info["chains"].as_array().is_some_and(|chains| chains + .iter() + .any(|chain| chain["chain"] == "bitcoin" && chain["network"] == network)), + "Original invoice belongs to another Bitcoin network" + ); + let client = reqwest::Client::builder() + .no_proxy() + .connect_timeout(std::time::Duration::from_secs(10)) + .timeout(std::time::Duration::from_secs(8)) + .danger_accept_invalid_certs(true) + .build()?; + let request=client.post(format!("{LND_REST_BASE_URL}/v2/router/send")).header("Grpc-Metadata-macaroon",macaroon).json(&serde_json::json!({"payment_request":bolt11,"no_inflight_updates":true,"timeout_seconds":120,"fee_limit_sat":amount})).build()?; + Ok(PreparedNativePayment { + client, + request, + hash: hash.into(), + amount, + }) + } +} + +impl crate::content_lightning::PreparedPayment for PreparedNativePayment { + async fn execute(self) -> Result { + PreparedNativePayment::execute(self).await + } +} +pub(crate) struct NativeNode<'a>(pub &'a RpcHandler); +impl crate::content_lightning::NativeInvoiceNode for NativeNode<'_> { + type Prepared = PreparedNativePayment; + async fn prepare(&self, invoice: &str, hash: &str, amount: u64) -> Result { + self.0 + .prepare_bound_invoice_payment(invoice, hash, amount) + .await + } + async fn lookup_payment(&self, hash: &str) -> Result { + self.0 + .handle_lnd_paymentstatus(Some(serde_json::json!({"payment_hash":hash}))) + .await + } +} diff --git a/core/archipelago/src/api/rpc/lnd/mod.rs b/core/archipelago/src/api/rpc/lnd/mod.rs index 9e26822b..a7c6ac21 100644 --- a/core/archipelago/src/api/rpc/lnd/mod.rs +++ b/core/archipelago/src/api/rpc/lnd/mod.rs @@ -1,4 +1,5 @@ mod channels; +pub(super) mod external_invoice; mod fee_bump; mod fee_policy; mod info; diff --git a/core/archipelago/src/api/rpc/lnd/payments.rs b/core/archipelago/src/api/rpc/lnd/payments.rs index 3a23765e..c3768a73 100644 --- a/core/archipelago/src/api/rpc/lnd/payments.rs +++ b/core/archipelago/src/api/rpc/lnd/payments.rs @@ -36,7 +36,7 @@ fn payment_failure_reason(reason: &str) -> &'static str { /// Preserve terminal LND state as structured data. An RPC exception is an /// ambiguous outcome to callers and must not hide a verified unpaid failure. -fn router_payment_outcome( +pub(super) fn router_payment_outcome( payment: &serde_json::Value, hash: &str, decoded_amt: i64, diff --git a/core/archipelago/src/api/rpc/mod.rs b/core/archipelago/src/api/rpc/mod.rs index e5c04f17..c8b735a1 100644 --- a/core/archipelago/src/api/rpc/mod.rs +++ b/core/archipelago/src/api/rpc/mod.rs @@ -17,6 +17,7 @@ mod fips; mod handshake; mod identity; mod interfaces; +mod lightning_purchase; pub(crate) mod lnd; mod marketplace; mod media_registration; @@ -109,6 +110,13 @@ fn native_consent_origin_allowed(method: &str, headers: &hyper::HeaderMap, dev_m | "media.registration.context" | "media.registration.resolve" | "content.rental-purchase" + | "content.invoice-pay" + | "content.invoice-download" + | "content.invoice-attempt" + | "content.invoice-retry-native" + | "content.invoice-create" + | "content.invoice-recover" + | "content.invoice-cancel" | "content.purchase" | "content.cancel-purchase" | "content.playback-handle" diff --git a/core/archipelago/src/api/rpc/purchase.rs b/core/archipelago/src/api/rpc/purchase.rs index da51aea9..83585d88 100644 --- a/core/archipelago/src/api/rpc/purchase.rs +++ b/core/archipelago/src/api/rpc/purchase.rs @@ -40,6 +40,16 @@ impl RpcHandler { crate::identity::NodeIdentity::load_existing(&self.config.data_dir.join("identity")) .await?; let buyer = identity.did_key()?; + let _rail = crate::content_payment_admission::lock( + &self.config.data_dir, + &buyer, + transport.seller_did(), + ¶ms.content_id, + ) + .await?; + self.ensure_invoice_allows_other_rail(&buyer, transport.seller_did(), ¶ms.content_id) + .await?; + let result = caller::purchase( &self.config.data_dir, &buyer, diff --git a/core/archipelago/src/content_lightning.rs b/core/archipelago/src/content_lightning.rs new file mode 100644 index 00000000..2e41938e --- /dev/null +++ b/core/archipelago/src/content_lightning.rs @@ -0,0 +1,1233 @@ +//! Durable external-invoice operations. Browser storage is supplemental only. +//! An ambiguous AddInvoice is never replayed: lookup the saved hash instead. +use anyhow::{Context, Result}; +use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; +use std::{ + fs, + io::{Read, Write}, + path::{Path, PathBuf}, +}; + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(crate) struct Binding { + pub id: String, + pub buyer_did: String, + pub seller_did: String, + pub content_id: String, + pub price_sats: u64, +} +impl Binding { + fn validate(&self) -> Result<()> { + anyhow::ensure!( + uuid::Uuid::parse_str(&self.id)?.to_string() == self.id, + "Invalid invoice operation" + ); + crate::identity::pubkey_bytes_from_did_key(&self.buyer_did)?; + crate::identity::pubkey_bytes_from_did_key(&self.seller_did)?; + anyhow::ensure!( + !self.content_id.is_empty() + && self.content_id.len() <= 128 + && self + .content_id + .bytes() + .all(|b| b.is_ascii_alphanumeric() || b == b'_' || b == b'-'), + "Invalid invoice content" + ); + anyhow::ensure!( + self.price_sats > 0 && self.price_sats <= i64::MAX as u64, + "Invalid invoice price" + ); + Ok(()) + } +} +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub(crate) enum Phase { + Prepared, + Dispatched, + Issued, + CancelRequested, + CanceledUnpaid, + Settled, +} +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(crate) struct RetainedFile { + pub sha256: String, + pub size: u64, + pub filename: String, + pub mime_type: String, +} +impl RetainedFile { + fn validate(&self) -> Result<()> { + anyhow::ensure!( + self.sha256.len() == 64 + && self.sha256.bytes().all(|b| b.is_ascii_hexdigit()) + && self.size > 0 + && !self.filename.is_empty() + && self.filename.len() <= 4096 + && !self.filename.chars().any(char::is_control) + && !self.mime_type.is_empty() + && self.mime_type.len() <= 256, + "Invalid retained invoice source metadata" + ); + hyper::header::HeaderValue::from_str(&self.mime_type)?; + Ok(()) + } +} +#[derive(Clone, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(crate) struct SellerRecord { + pub binding: Binding, + preimage: String, + pub payment_hash: String, + pub phase: Phase, + pub source: Option, + pub bolt11: Option, +} +#[derive(Clone, Debug, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(crate) struct BuyerRecord { + pub binding: Binding, + pub seller_onion: String, + pub external_exposure: bool, + #[serde(default)] + pub native_retired: bool, + #[serde(default)] + pub native_replacement: Option, + #[serde(default)] + pub native_dispatched: bool, + #[serde(default)] + pub native_result: Option, + pub last: Option, +} +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(crate) struct Status { + pub binding: Binding, + pub payment_hash: String, + pub bolt11: Option, + pub state: Phase, + pub can_switch_method: bool, + pub source: Option, +} +#[derive(Clone)] +pub(crate) struct Invoice { + pub payment_hash: String, + pub bolt11: String, + pub price_sats: u64, + pub state: String, + pub paid_sats: Option, + pub paid_msats: Option, +} +pub(crate) trait InvoiceNode { + async fn prepare_creation(&self) -> Result<()>; + async fn lookup(&self, hash: &str) -> Result>; + async fn add(&self, binding: &Binding, preimage_hex: &str) -> Result<()>; + async fn cancel(&self, hash: &str) -> Result<()>; +} +#[derive(Serialize, Deserialize)] +struct Envelope { + payload: String, + checksum: String, +} +pub(crate) struct Journal { + directory: PathBuf, + _lock: fs::File, +} +impl Journal { + pub async fn open(data_dir: &Path) -> Result { + let data = data_dir.to_path_buf(); + tokio::task::spawn_blocking(move || { + use std::os::{ + fd::AsRawFd, + unix::fs::{OpenOptionsExt, PermissionsExt}, + }; + fs::create_dir_all(&data)?; + let data = fs::canonicalize(data)?; + let directory = data.join("content-lightning"); + fs::create_dir_all(&directory)?; + anyhow::ensure!( + fs::symlink_metadata(&directory)?.is_dir(), + "Invoice journal is not a directory" + ); + fs::set_permissions(&directory, fs::Permissions::from_mode(0o700))?; + fs::File::open(&data)?.sync_all()?; + let lock = fs::OpenOptions::new() + .read(true) + .write(true) + .create(true) + .mode(0o600) + .custom_flags(libc::O_NOFOLLOW | libc::O_NONBLOCK) + .open(directory.join(".lock"))?; + anyhow::ensure!(lock.metadata()?.is_file(), "Invalid invoice lock"); + loop { + if unsafe { libc::flock(lock.as_raw_fd(), libc::LOCK_EX) } == 0 { + break; + } + let e = std::io::Error::last_os_error(); + if e.kind() != std::io::ErrorKind::Interrupted { + return Err(e.into()); + } + } + Ok(Self { + directory, + _lock: lock, + }) + }) + .await? + } + fn path(&self, role: &str, id: &str) -> Result { + anyhow::ensure!( + matches!(role, "buyer" | "seller") && uuid::Uuid::parse_str(id)?.to_string() == id, + "Invalid invoice journal key" + ); + Ok(self.directory.join(format!("{role}-{id}.json"))) + } + fn read(&self, role: &str, id: &str) -> Result> { + use std::os::unix::fs::OpenOptionsExt; + let file = match fs::OpenOptions::new() + .read(true) + .custom_flags(libc::O_NOFOLLOW | libc::O_NONBLOCK) + .open(self.path(role, id)?) + { + Ok(v) => v, + Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None), + Err(e) => return Err(e.into()), + }; + anyhow::ensure!(file.metadata()?.is_file(), "Invalid invoice record"); + let mut bytes = Vec::new(); + file.take(65537).read_to_end(&mut bytes)?; + anyhow::ensure!(bytes.len() <= 65536, "Invoice record too large"); + let envelope: Envelope = + serde_json::from_slice(&bytes).context("Invoice recovery damaged; do not pay again")?; + anyhow::ensure!( + hex::encode(Sha256::digest(envelope.payload.as_bytes())) == envelope.checksum, + "Invoice recovery checksum changed" + ); + Ok(Some(serde_json::from_str(&envelope.payload)?)) + } + fn write(&self, role: &str, id: &str, value: &T) -> Result<()> { + use std::os::unix::fs::OpenOptionsExt; + let payload = serde_json::to_string(value)?; + let bytes = serde_json::to_vec(&Envelope { + checksum: hex::encode(Sha256::digest(payload.as_bytes())), + payload, + })?; + anyhow::ensure!(bytes.len() <= 65536, "Invoice record too large"); + let temporary = self + .directory + .join(format!(".{}.tmp", uuid::Uuid::new_v4())); + let result = (|| -> Result<()> { + let mut f = fs::OpenOptions::new() + .write(true) + .create_new(true) + .mode(0o600) + .open(&temporary)?; + f.write_all(&bytes)?; + f.sync_all()?; + fs::rename(&temporary, self.path(role, id)?)?; + fs::File::open(&self.directory)?.sync_all()?; + Ok(()) + })(); + if result.is_err() { + let _ = fs::remove_file(temporary); + } + result + } + pub fn seller(&self, binding: &Binding) -> Result> { + binding.validate()?; + let record: Option = self.read("seller", &binding.id)?; + if let Some(v) = &record { + anyhow::ensure!(&v.binding == binding, "Invoice operation binding changed"); + let secret = hex::decode(&v.preimage)?; + anyhow::ensure!( + secret.len() == 32 && hex::encode(Sha256::digest(secret)) == v.payment_hash, + "Invoice preimage binding damaged" + ); + } + Ok(record) + } + pub fn prepare_seller(&self, binding: Binding) -> Result { + self.prepare_seller_source(binding, None) + } + pub fn prepare_seller_source( + &self, + binding: Binding, + source: Option, + ) -> Result { + if let Some(source) = &source { + source.validate()?; + } + if let Some(saved) = self.seller(&binding)? { + anyhow::ensure!(saved.source == source, "Original invoice source changed"); + return Ok(saved); + } + use rand::RngCore; + let mut secret = [0u8; 32]; + rand::rngs::OsRng.fill_bytes(&mut secret); + let record = SellerRecord { + payment_hash: hex::encode(Sha256::digest(secret)), + preimage: hex::encode(secret), + binding, + phase: Phase::Prepared, + source, + bolt11: None, + }; + self.save_seller(&record)?; + Ok(record) + } + pub fn save_seller(&self, record: &SellerRecord) -> Result<()> { + self.write("seller", &record.binding.id, record) + } + pub fn buyer(&self, id: &str) -> Result> { + let record: Option = self.read("buyer", id)?; + if let Some(v) = &record { + v.binding.validate()?; + anyhow::ensure!(v.binding.id == id, "Invoice operation changed"); + } + Ok(record) + } + pub fn buyer_for( + &self, + buyer: &str, + seller: &str, + content: &str, + ) -> Result> { + let mut found = None; + for entry in fs::read_dir(&self.directory)? { + let name = entry? + .file_name() + .into_string() + .map_err(|_| anyhow::anyhow!("Invalid invoice record name"))?; + let Some(id) = name + .strip_prefix("buyer-") + .and_then(|s| s.strip_suffix(".json")) + else { + continue; + }; + let record: BuyerRecord = self + .read("buyer", id)? + .context("Invoice record disappeared")?; + record.binding.validate()?; + let unfinished_replacement = if record.native_retired { + if let Some(id) = &record.native_replacement { + self.buyer(id)?.is_none() + } else { + false + } + } else { + false + }; + if record.binding.buyer_did == buyer + && record.binding.seller_did == seller + && record.binding.content_id == content + && (!record.native_retired || unfinished_replacement) + && (unfinished_replacement + || record.last.as_ref().is_none_or(|s| !s.can_switch_method)) + { + anyhow::ensure!( + found.is_none(), + "Multiple unresolved invoice operations; recover them first" + ); + found = Some(record) + } + } + Ok(found) + } + /// Explicit retry only: retire a proven native-only failure and retain its + /// replacement UUID before creating anything. Recovery reuses that UUID. + pub fn retry_native(&self, id: &str) -> Result { + let mut old = self.buyer(id)?.context("Original native invoice missing")?; + anyhow::ensure!( + !old.external_exposure + && old.native_dispatched + && old.native_result.as_deref() == Some("failed") + && old.last.as_ref().is_some_and(|s| s.state != Phase::Settled), + "Only a confirmed native-only failure can be retried" + ); + anyhow::ensure!( + !old.native_retired || old.native_replacement.is_some(), + "Original invoice was retired for another payment method" + ); + let replacement = old + .native_replacement + .clone() + .unwrap_or_else(|| uuid::Uuid::new_v4().to_string()); + let mut binding = old.binding.clone(); + binding.id = replacement.clone(); + if let Some(saved) = self.buyer(&replacement)? { + anyhow::ensure!( + saved.binding == binding && saved.seller_onion == old.seller_onion, + "Native replacement binding changed" + ); + return Ok(saved); + } + let current = self + .buyer_for( + &old.binding.buyer_did, + &old.binding.seller_did, + &old.binding.content_id, + )? + .context("Original native operation no longer owns this purchase")?; + anyhow::ensure!( + current.binding.id == old.binding.id, + "Another operation owns this purchase" + ); + old.native_retired = true; + old.native_replacement = Some(replacement); + self.save_buyer(&old)?; + let new = BuyerRecord { + binding, + seller_onion: old.seller_onion, + external_exposure: false, + native_retired: false, + native_replacement: None, + native_dispatched: false, + native_result: None, + last: None, + }; + self.save_buyer(&new)?; + Ok(new) + } + pub fn save_buyer(&self, record: &BuyerRecord) -> Result<()> { + record.binding.validate()?; + anyhow::ensure!( + matches!( + record.native_result.as_deref(), + None | Some("failed" | "succeeded") + ), + "Invalid native invoice outcome" + ); + anyhow::ensure!( + !record.native_retired + || (!record.external_exposure && record.native_result.as_deref() == Some("failed")), + "Retired native invoice cannot be exposed" + ); + if let Some(id) = &record.native_replacement { + anyhow::ensure!( + record.native_retired + && uuid::Uuid::parse_str(id)?.to_string() == *id + && *id != record.binding.id, + "Invalid native replacement identity" + ); + } + if let Some(old) = self.read::("buyer", &record.binding.id)? { + anyhow::ensure!( + old.binding == record.binding + && old.seller_onion == record.seller_onion + && (!old.external_exposure || record.external_exposure) + && (!old.native_retired || record.native_retired) + && old + .native_replacement + .as_ref() + .is_none_or(|id| record.native_replacement.as_ref() == Some(id)) + && (!old.native_dispatched || record.native_dispatched) + && (old.native_result.as_deref() != Some("succeeded") + || record.native_result.as_deref() == Some("succeeded")), + "Invoice buyer binding changed" + ); + if old.last.as_ref().is_some_and(|s| s.state == Phase::Settled) { + anyhow::ensure!( + record + .last + .as_ref() + .is_some_and(|s| s.state == Phase::Settled), + "Settled invoice cannot regress" + ); + } + } + if let Some(status) = &record.last { + if let Some(source) = &status.source { + source.validate()?; + } + anyhow::ensure!( + status.binding == record.binding + && !(record.native_result.as_deref() == Some("succeeded") + && status.can_switch_method) + && status.can_switch_method == (status.state == Phase::CanceledUnpaid), + "Invoice status binding changed" + ); + } + self.write("buyer", &record.binding.id, record) + } +} +impl SellerRecord { + pub fn status(&self) -> Status { + Status { + binding: self.binding.clone(), + payment_hash: self.payment_hash.clone(), + bolt11: self.bolt11.clone(), + state: self.phase.clone(), + can_switch_method: self.phase == Phase::CanceledUnpaid, + source: self.source.clone(), + } + } + fn observe(&mut self, invoice: Invoice) -> Result<()> { + anyhow::ensure!( + invoice.payment_hash == self.payment_hash + && invoice.price_sats == self.binding.price_sats + && !invoice.bolt11.is_empty(), + "LND invoice binding changed" + ); + if let Some(original) = &self.bolt11 { + anyhow::ensure!(original == &invoice.bolt11, "Original invoice changed"); + } + self.bolt11 = Some(invoice.bolt11); + let settled = invoice.state == "SETTLED" + && invoice + .paid_sats + .is_some_and(|v| v >= self.binding.price_sats) + && invoice.paid_msats.is_none_or(|v| { + self.binding + .price_sats + .checked_mul(1000) + .is_some_and(|required| v >= required) + }); + if settled { + self.phase = Phase::Settled + } else if self.phase != Phase::Settled + && invoice.state == "CANCELED" + && invoice.paid_sats == Some(0) + && invoice.paid_msats.is_none_or(|v| v == 0) + { + self.phase = Phase::CanceledUnpaid + } else if self.phase != Phase::Settled + && self.phase != Phase::CanceledUnpaid + && self.phase != Phase::CancelRequested + { + self.phase = Phase::Issued + } + Ok(()) + } +} +/// Journal lock is retained through network calls. Persist dispatch BEFORE await. +/// Cancellation of this future leaves a recoverable record, never permission to +/// issue another invoice. A prepared operation can be canceled before dispatch. +pub(crate) async fn drive( + journal: &Journal, + binding: &Binding, + node: &N, + cancel: bool, +) -> Result { + let mut record = journal + .seller(binding)? + .context("Unknown invoice operation")?; + if matches!(record.phase, Phase::Settled | Phase::CanceledUnpaid) { + return Ok(record.status()); + } + if cancel && record.phase == Phase::Prepared { + record.phase = Phase::CanceledUnpaid; + journal.save_seller(&record)?; + return Ok(record.status()); + } + if cancel { + record.phase = Phase::CancelRequested; + journal.save_seller(&record)?; + } + if record.phase == Phase::Prepared { + node.prepare_creation().await?; + record.phase = Phase::Dispatched; + journal.save_seller(&record)?; + // Exactly one AddInvoice attempt. A failed response may still have created + // it; subsequent operations only look up the saved hash. + let _ = node.add(binding, &record.preimage).await; + } + let Some(invoice) = node.lookup(&record.payment_hash).await? else { + return Ok(record.status()); + }; + record.observe(invoice)?; + journal.save_seller(&record)?; + if cancel && record.phase != Phase::Settled && record.phase != Phase::CanceledUnpaid { + let _ = node.cancel(&record.payment_hash).await; + if let Some(invoice) = node.lookup(&record.payment_hash).await? { + record.observe(invoice)?; + journal.save_seller(&record)?; + } + } + Ok(record.status()) +} + +/// Runtime adapters prepare a request without dispatching it, then consume it once. +pub(crate) trait PreparedPayment { + async fn execute(self) -> Result; +} +pub(crate) trait NativeInvoiceNode { + type Prepared: PreparedPayment; + async fn prepare(&self, invoice: &str, hash: &str, amount: u64) -> Result; + async fn lookup_payment(&self, hash: &str) -> Result; +} +/// Caller retains the per-buyer/seller/item admission lock across this operation. +/// The journal lock is released during actual payment, so unrelated invoices can recover. +pub(crate) async fn drive_native( + data_dir: &Path, + journal: Journal, + id: &str, + node: &N, +) -> Result { + let mut record = journal + .buyer(id)? + .context("Original invoice operation missing")?; + anyhow::ensure!( + !record.native_retired, + "Original native invoice was retired before changing methods" + ); + let status = record + .last + .as_ref() + .context("Original invoice has not been created")?; + anyhow::ensure!(!status.can_switch_method, "Original invoice was canceled"); + if status.state == Phase::Settled || record.native_result.as_deref() == Some("succeeded") { + return Ok(serde_json::json!({"status":"succeeded","payment_hash":status.payment_hash})); + } + anyhow::ensure!( + status.source.is_some(), + "Original invoice snapshot is not confirmed; no payment dispatched" + ); + if record.native_dispatched { + if record.native_result.as_deref() == Some("failed") { + return Ok(serde_json::json!({"status":"failed","payment_hash":status.payment_hash})); + } + let payment = node.lookup_payment(&status.payment_hash).await?; + if matches!(payment["status"].as_str(), Some("succeeded" | "failed")) { + record.native_result = payment["status"].as_str().map(str::to_owned); + journal.save_buyer(&record)?; + } + return Ok(payment); + } + let invoice = status + .bolt11 + .as_ref() + .context("Original invoice is unavailable")?; + // Preparation validates identity/amount/network/expiry and builds the request; + // deterministic preparation failures leave this same operation undispatched. + let prepared = node + .prepare(invoice, &status.payment_hash, record.binding.price_sats) + .await?; + record.native_dispatched = true; + journal.save_buyer(&record)?; + drop(journal); + let payment = prepared.execute().await?; + if matches!(payment["status"].as_str(), Some("succeeded" | "failed")) { + record.native_result = payment["status"].as_str().map(str::to_owned); + Journal::open(data_dir).await?.save_buyer(&record)?; + } + Ok(payment) +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::{ + atomic::{AtomicUsize, Ordering}, + Mutex, + }; + struct Node { + invoice: Mutex>, + adds: AtomicUsize, + lost_reply: bool, + settle_on_cancel: bool, + reject_preflight: std::sync::atomic::AtomicBool, + } + impl Node { + fn new() -> Self { + Self { + invoice: Mutex::new(None), + adds: AtomicUsize::new(0), + lost_reply: true, + settle_on_cancel: false, + reject_preflight: std::sync::atomic::AtomicBool::new(false), + } + } + } + impl InvoiceNode for Node { + async fn prepare_creation(&self) -> Result<()> { + anyhow::ensure!( + !self.reject_preflight.load(Ordering::SeqCst), + "LND unavailable before invoice dispatch" + ); + Ok(()) + } + async fn lookup(&self, _: &str) -> Result> { + Ok(self.invoice.lock().unwrap().clone()) + } + async fn add(&self, b: &Binding, p: &str) -> Result<()> { + self.adds.fetch_add(1, Ordering::SeqCst); + *self.invoice.lock().unwrap() = Some(Invoice { + payment_hash: hex::encode(Sha256::digest(hex::decode(p)?)), + bolt11: "ln-original".into(), + price_sats: b.price_sats, + state: "OPEN".into(), + paid_sats: Some(0), + paid_msats: Some(0), + }); + if self.lost_reply { + anyhow::bail!("reply lost after LND stored invoice") + } + Ok(()) + } + async fn cancel(&self, _: &str) -> Result<()> { + let mut guard = self.invoice.lock().unwrap(); + let v = guard.as_mut().unwrap(); + if self.settle_on_cancel { + v.state = "SETTLED".into(); + v.paid_sats = Some(v.price_sats); + v.paid_msats = Some(v.price_sats * 1000) + } else { + v.state = "CANCELED".into() + }; + anyhow::bail!("cancel reply lost") + } + } + fn binding() -> Binding { + Binding { + id: uuid::Uuid::new_v4().to_string(), + buyer_did: crate::identity::did_key_from_pubkey_hex(&hex::encode([7; 32])).unwrap(), + seller_did: crate::identity::did_key_from_pubkey_hex(&hex::encode([8; 32])).unwrap(), + content_id: "file".into(), + price_sats: 8, + } + } + #[tokio::test] + async fn lost_add_reply_and_process_restart_recover_original_invoice_once() { + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let node = Node::new(); + let j = Journal::open(root.path()).await.unwrap(); + j.prepare_seller(b.clone()).unwrap(); + let first = drive(&j, &b, &node, false).await.unwrap(); + drop(j); + let j = Journal::open(root.path()).await.unwrap(); + let second = drive(&j, &b, &node, false).await.unwrap(); + assert_eq!(first, second); + assert_eq!(node.adds.load(Ordering::SeqCst), 1); + assert_eq!(second.state, Phase::Issued); + } + #[tokio::test] + async fn ambiguous_dispatch_missing_lookup_cannot_reissue_or_cancel_as_unpaid() { + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let node = Node::new(); + let j = Journal::open(root.path()).await.unwrap(); + let mut r = j.prepare_seller(b.clone()).unwrap(); + r.phase = Phase::Dispatched; + j.save_seller(&r).unwrap(); + drop(j); + let j = Journal::open(root.path()).await.unwrap(); + let unknown = drive(&j, &b, &node, true).await.unwrap(); + assert_eq!(unknown.state, Phase::CancelRequested); + assert!(!unknown.can_switch_method); + assert_eq!(node.adds.load(Ordering::SeqCst), 0); + // Original delayed AddInvoice arrives after the first cancel lookup. + node.add(&b, &r.preimage).await.unwrap_err(); + let resolved = drive(&j, &b, &node, true).await.unwrap(); + assert_eq!(resolved.state, Phase::CanceledUnpaid); + assert!(resolved.can_switch_method); + assert_eq!(node.adds.load(Ordering::SeqCst), 1); + } + #[tokio::test] + async fn prepared_cancel_has_no_remote_creation_and_cannot_be_reopened() { + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let node = Node::new(); + let j = Journal::open(root.path()).await.unwrap(); + j.prepare_seller(b.clone()).unwrap(); + assert!(drive(&j, &b, &node, true).await.unwrap().can_switch_method); + assert!(drive(&j, &b, &node, false).await.unwrap().can_switch_method); + assert_eq!(node.adds.load(Ordering::SeqCst), 0); + } + #[tokio::test] + async fn settlement_wins_lost_cancel_response_and_survives_missing_lnd_record() { + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let mut node = Node::new(); + node.settle_on_cancel = true; + let j = Journal::open(root.path()).await.unwrap(); + j.prepare_seller(b.clone()).unwrap(); + drive(&j, &b, &node, false).await.unwrap(); + let paid = drive(&j, &b, &node, true).await.unwrap(); + assert_eq!(paid.state, Phase::Settled); + assert!(!paid.can_switch_method); + *node.invoice.lock().unwrap() = None; + assert_eq!(drive(&j, &b, &node, true).await.unwrap(), paid); + } + #[tokio::test] + async fn browser_loss_finds_original_buyer_operation_and_exposure_is_monotonic() { + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let j = Journal::open(root.path()).await.unwrap(); + let mut record = BuyerRecord { + binding: b.clone(), + seller_onion: "original.onion".into(), + external_exposure: true, + native_retired: false, + native_replacement: None, + native_dispatched: false, + native_result: None, + last: None, + }; + j.save_buyer(&record).unwrap(); + drop(j); + let j = Journal::open(root.path()).await.unwrap(); + assert_eq!( + j.buyer_for(&b.buyer_did, &b.seller_did, &b.content_id) + .unwrap() + .unwrap() + .binding, + b + ); + record.external_exposure = false; + assert!(j.save_buyer(&record).is_err()); + let mut changed = b.clone(); + changed.price_sats += 1; + j.prepare_seller(b).unwrap(); + assert!(j.prepare_seller(changed).is_err()); + } + #[tokio::test] + async fn checksum_damage_blocks_replacement() { + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let j = Journal::open(root.path()).await.unwrap(); + j.prepare_seller(b.clone()).unwrap(); + let p = j.path("seller", &b.id).unwrap(); + fs::write(p, b"{}").unwrap(); + assert!(j.prepare_seller(b).is_err()); + } + #[tokio::test] + async fn native_success_cannot_be_replaced_by_contradictory_canceled_status() { + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let j = Journal::open(root.path()).await.unwrap(); + let mut last = j.prepare_seller(b.clone()).unwrap().status(); + last.state = Phase::Issued; + let mut record = BuyerRecord { + binding: b.clone(), + seller_onion: "original.onion".into(), + external_exposure: false, + native_retired: false, + native_replacement: None, + native_dispatched: true, + native_result: Some("succeeded".into()), + last: Some(last), + }; + j.save_buyer(&record).unwrap(); + record.last.as_mut().unwrap().state = Phase::CanceledUnpaid; + record.last.as_mut().unwrap().can_switch_method = true; + assert!(j.save_buyer(&record).is_err()); + drop(j); + let j = Journal::open(root.path()).await.unwrap(); + assert_eq!( + j.buyer_for(&b.buyer_did, &b.seller_did, &b.content_id) + .unwrap() + .unwrap() + .native_result + .as_deref(), + Some("succeeded") + ); + } + #[tokio::test] + async fn retired_native_failure_stays_retired_and_cannot_later_expose_invoice() { + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let j = Journal::open(root.path()).await.unwrap(); + let mut record = BuyerRecord { + binding: b.clone(), + seller_onion: "original.onion".into(), + external_exposure: false, + native_retired: false, + native_replacement: None, + native_dispatched: true, + native_result: Some("failed".into()), + last: None, + }; + j.save_buyer(&record).unwrap(); + record.native_retired = true; + j.save_buyer(&record).unwrap(); + drop(j); + let j = Journal::open(root.path()).await.unwrap(); + assert!(j + .buyer_for(&b.buyer_did, &b.seller_did, &b.content_id) + .unwrap() + .is_none()); + assert!(j.buyer(&b.id).unwrap().unwrap().native_retired); + record.external_exposure = true; + assert!(j.save_buyer(&record).is_err()); + record.external_exposure = false; + record.native_retired = false; + assert!(j.save_buyer(&record).is_err()); + } + #[tokio::test] + async fn inconsistent_paid_units_cannot_create_settlement_or_cancellation() { + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let node = Node::new(); + let j = Journal::open(root.path()).await.unwrap(); + j.prepare_seller(b.clone()).unwrap(); + drive(&j, &b, &node, false).await.unwrap(); + { + let mut held = node.invoice.lock().unwrap(); + let invoice = held.as_mut().unwrap(); + invoice.state = "SETTLED".into(); + invoice.paid_sats = Some(b.price_sats); + invoice.paid_msats = Some(0); + } + let result = drive(&j, &b, &node, false).await.unwrap(); + assert_ne!(result.state, Phase::Settled); + assert!(!result.can_switch_method); + { + let mut held = node.invoice.lock().unwrap(); + let invoice = held.as_mut().unwrap(); + invoice.state = "CANCELED".into(); + invoice.paid_sats = Some(0); + invoice.paid_msats = Some(1); + } + assert!(!drive(&j, &b, &node, false).await.unwrap().can_switch_method); + } + #[tokio::test] + async fn snapshot_commit_rechecks_changed_terms_and_unshare_before_invoice_exists() { + use crate::content_server::{ + self as catalog, AccessControl, Availability, ContentCatalog, ContentItem, + }; + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let original = ContentItem { + id: b.content_id.clone(), + filename: "file.txt".into(), + mime_type: "text/plain".into(), + size_bytes: 4, + description: String::new(), + access: AccessControl::Paid { + price_sats: b.price_sats, + accepted: vec!["lightning".into()], + }, + availability: Availability::AllPeers, + added_at: String::new(), + }; + let retained = RetainedFile { + sha256: "ab".repeat(32), + size: 4, + filename: original.filename.clone(), + mime_type: original.mime_type.clone(), + }; + let j = Journal::open(root.path()).await.unwrap(); + let mut changed = original.clone(); + changed.access = AccessControl::Paid { + price_sats: b.price_sats + 1, + accepted: vec![], + }; + catalog::save_catalog( + root.path(), + &ContentCatalog { + items: vec![changed], + }, + ) + .await + .unwrap(); + assert!(catalog::publish_snapshot_invoice( + root.path(), + &original, + &j, + b.clone(), + retained.clone() + ) + .await + .is_err()); + assert!(j.seller(&b).unwrap().is_none()); + catalog::save_catalog(root.path(), &ContentCatalog { items: vec![] }) + .await + .unwrap(); + assert!(catalog::publish_snapshot_invoice( + root.path(), + &original, + &j, + b.clone(), + retained.clone() + ) + .await + .is_err()); + assert!(j.seller(&b).unwrap().is_none()); + catalog::save_catalog( + root.path(), + &ContentCatalog { + items: vec![original.clone()], + }, + ) + .await + .unwrap(); + let prepared = catalog::publish_snapshot_invoice( + root.path(), + &original, + &j, + b.clone(), + retained.clone(), + ) + .await + .unwrap(); + assert_eq!(prepared.phase, Phase::Prepared); + assert_eq!(prepared.source, Some(retained)); + } + struct Native { + prepares: AtomicUsize, + executions: std::sync::Arc, + lookups: AtomicUsize, + reject_preflight: std::sync::atomic::AtomicBool, + lose_reply: bool, + } + struct PreparedNative { + executions: std::sync::Arc, + hash: String, + lose_reply: bool, + } + impl PreparedPayment for PreparedNative { + async fn execute(self) -> Result { + self.executions.fetch_add(1, Ordering::SeqCst); + anyhow::ensure!(!self.lose_reply, "Response lost after dispatch"); + Ok(serde_json::json!({"status":"succeeded","payment_hash":self.hash})) + } + } + impl NativeInvoiceNode for Native { + type Prepared = PreparedNative; + async fn prepare(&self, _: &str, hash: &str, _: u64) -> Result { + self.prepares.fetch_add(1, Ordering::SeqCst); + anyhow::ensure!( + !self.reject_preflight.load(Ordering::SeqCst), + "Wrong configured network before dispatch" + ); + Ok(PreparedNative { + executions: self.executions.clone(), + hash: hash.into(), + lose_reply: self.lose_reply, + }) + } + async fn lookup_payment(&self, hash: &str) -> Result { + self.lookups.fetch_add(1, Ordering::SeqCst); + Ok(serde_json::json!({"status":"succeeded","payment_hash":hash})) + } + } + impl Native { + fn new(lose_reply: bool) -> Self { + Self { + prepares: AtomicUsize::new(0), + executions: std::sync::Arc::new(AtomicUsize::new(0)), + lookups: AtomicUsize::new(0), + reject_preflight: std::sync::atomic::AtomicBool::new(false), + lose_reply, + } + } + } + fn prepared_native_buyer(journal: &Journal, b: &Binding) -> BuyerRecord { + let mut seller = journal + .prepare_seller_source( + b.clone(), + Some(RetainedFile { + sha256: "ab".repeat(32), + size: 4, + filename: "file.txt".into(), + mime_type: "text/plain".into(), + }), + ) + .unwrap(); + seller.phase = Phase::Issued; + seller.bolt11 = Some("ln-fixture".into()); + journal.save_seller(&seller).unwrap(); + let record = BuyerRecord { + binding: b.clone(), + seller_onion: "original.onion".into(), + external_exposure: false, + native_retired: false, + native_replacement: None, + native_dispatched: false, + native_result: None, + last: Some(seller.status()), + }; + journal.save_buyer(&record).unwrap(); + record + } + #[tokio::test] + async fn native_preflight_failure_can_retry_original_operation_before_single_dispatch() { + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let node = Native::new(false); + let journal = Journal::open(root.path()).await.unwrap(); + prepared_native_buyer(&journal, &b); + node.reject_preflight.store(true, Ordering::SeqCst); + assert!(drive_native(root.path(), journal, &b.id, &node) + .await + .is_err()); + let journal = Journal::open(root.path()).await.unwrap(); + assert!(!journal.buyer(&b.id).unwrap().unwrap().native_dispatched); + node.reject_preflight.store(false, Ordering::SeqCst); + assert_eq!( + drive_native(root.path(), journal, &b.id, &node) + .await + .unwrap()["status"], + "succeeded" + ); + let journal = Journal::open(root.path()).await.unwrap(); + assert_eq!( + drive_native(root.path(), journal, &b.id, &node) + .await + .unwrap()["status"], + "succeeded" + ); + assert_eq!(node.prepares.load(Ordering::SeqCst), 2); + assert_eq!(node.executions.load(Ordering::SeqCst), 1); + assert_eq!(node.lookups.load(Ordering::SeqCst), 0); + } + #[tokio::test] + async fn native_lost_reply_restarts_with_original_hash_lookup_and_never_dispatches_twice() { + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let node = Native::new(true); + let journal = Journal::open(root.path()).await.unwrap(); + let original = prepared_native_buyer(&journal, &b); + assert!(drive_native(root.path(), journal, &b.id, &node) + .await + .is_err()); + let journal = Journal::open(root.path()).await.unwrap(); + assert!(journal.buyer(&b.id).unwrap().unwrap().native_dispatched); + let recovered = drive_native(root.path(), journal, &b.id, &node) + .await + .unwrap(); + assert_eq!( + recovered["payment_hash"], + original.last.unwrap().payment_hash + ); + assert_eq!(recovered["status"], "succeeded"); + let journal = Journal::open(root.path()).await.unwrap(); + assert_eq!( + journal + .buyer(&b.id) + .unwrap() + .unwrap() + .native_result + .as_deref(), + Some("succeeded") + ); + drive_native(root.path(), journal, &b.id, &node) + .await + .unwrap(); + assert_eq!(node.prepares.load(Ordering::SeqCst), 1); + assert_eq!(node.executions.load(Ordering::SeqCst), 1); + assert_eq!(node.lookups.load(Ordering::SeqCst), 1); + } + #[tokio::test] + async fn retired_invoice_rejects_delayed_native_callback_without_preflight_or_payment() { + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let node = Native::new(false); + let journal = Journal::open(root.path()).await.unwrap(); + let mut original = prepared_native_buyer(&journal, &b); + original.native_dispatched = true; + original.native_result = Some("failed".into()); + original.native_retired = true; + journal.save_buyer(&original).unwrap(); + assert!(drive_native(root.path(), journal, &b.id, &node) + .await + .is_err()); + assert_eq!(node.prepares.load(Ordering::SeqCst), 0); + assert_eq!(node.executions.load(Ordering::SeqCst), 0); + assert_eq!(node.lookups.load(Ordering::SeqCst), 0); + } + #[tokio::test] + async fn explicit_retry_links_one_fresh_uuid_and_rejects_old_callbacks() { + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let node = Native::new(false); + let j = Journal::open(root.path()).await.unwrap(); + let mut old = prepared_native_buyer(&j, &b); + old.native_dispatched = true; + old.native_result = Some("failed".into()); + j.save_buyer(&old).unwrap(); + let new = j.retry_native(&b.id).unwrap(); + assert_ne!(new.binding.id, b.id); + assert!(!new.native_dispatched); + assert!(!new.external_exposure); + assert_eq!(j.retry_native(&b.id).unwrap().binding, new.binding); + assert_eq!( + j.buyer_for(&b.buyer_did, &b.seller_did, &b.content_id) + .unwrap() + .unwrap() + .binding, + new.binding + ); + assert!(j.save_buyer(&old).is_err()); + assert!(drive_native(root.path(), j, &b.id, &node).await.is_err()); + assert_eq!(node.executions.load(Ordering::SeqCst), 0); + let j = Journal::open(root.path()).await.unwrap(); + assert_eq!( + j.buyer(&new.binding.id).unwrap().unwrap().binding, + new.binding + ); + } + #[tokio::test] + async fn interrupted_retry_retirement_blocks_other_rails_and_recovers_same_uuid() { + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let replacement = uuid::Uuid::new_v4().to_string(); + let j = Journal::open(root.path()).await.unwrap(); + let mut old = prepared_native_buyer(&j, &b); + old.native_dispatched = true; + old.native_result = Some("failed".into()); + old.native_retired = true; + old.native_replacement = Some(replacement.clone()); + j.save_buyer(&old).unwrap(); + drop(j); + let j = Journal::open(root.path()).await.unwrap(); + assert_eq!( + j.buyer_for(&b.buyer_did, &b.seller_did, &b.content_id) + .unwrap() + .unwrap() + .native_replacement, + Some(replacement.clone()) + ); + old.last.as_mut().unwrap().state = Phase::CanceledUnpaid; + old.last.as_mut().unwrap().can_switch_method = true; + j.save_buyer(&old).unwrap(); + assert!(j + .buyer_for(&b.buyer_did, &b.seller_did, &b.content_id) + .unwrap() + .is_some()); + assert_eq!(j.retry_native(&b.id).unwrap().binding.id, replacement); + assert_eq!(j.retry_native(&b.id).unwrap().binding.id, replacement); + } + #[tokio::test] + async fn explicit_retry_never_replaces_success_pending_or_externally_exposed_invoice() { + for (outcome, exposed) in [ + (Some("succeeded"), false), + (None, false), + (Some("failed"), true), + ] { + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let j = Journal::open(root.path()).await.unwrap(); + let mut old = prepared_native_buyer(&j, &b); + old.native_dispatched = true; + old.native_result = outcome.map(str::to_owned); + old.external_exposure = exposed; + j.save_buyer(&old).unwrap(); + assert!(j.retry_native(&b.id).is_err()); + assert!(!j.buyer(&b.id).unwrap().unwrap().native_retired); + } + } + #[tokio::test] + async fn unavailable_seller_preflight_preserves_prepared_operation_for_retry() { + let root = tempfile::tempdir().unwrap(); + let b = binding(); + let node = Node::new(); + let j = Journal::open(root.path()).await.unwrap(); + let original = j.prepare_seller(b.clone()).unwrap(); + node.reject_preflight.store(true, Ordering::SeqCst); + assert!(drive(&j, &b, &node, false).await.is_err()); + assert_eq!(j.seller(&b).unwrap().unwrap().phase, Phase::Prepared); + assert_eq!(node.adds.load(Ordering::SeqCst), 0); + node.reject_preflight.store(false, Ordering::SeqCst); + assert_eq!( + drive(&j, &b, &node, false).await.unwrap().payment_hash, + original.payment_hash + ); + assert_eq!(node.adds.load(Ordering::SeqCst), 1); + } +} diff --git a/core/archipelago/src/content_payment_admission.rs b/core/archipelago/src/content_payment_admission.rs new file mode 100644 index 00000000..131b966b --- /dev/null +++ b/core/archipelago/src/content_payment_admission.rs @@ -0,0 +1,44 @@ +//! Outermost cross-rail admission, before wallet or payment journals. +use anyhow::Result; +use sha2::{Digest, Sha256}; +use std::{fs, path::Path}; +pub(crate) struct Guard { + _file: fs::File, +} +pub(crate) async fn lock(data: &Path, buyer: &str, seller: &str, content: &str) -> Result { + let root = data.join("content-payment-admission"); + let key = hex::encode(Sha256::digest(serde_json::to_vec(&( + buyer, seller, content, + ))?)); + tokio::task::spawn_blocking(move || { + use std::os::{ + fd::AsRawFd, + unix::fs::{OpenOptionsExt, PermissionsExt}, + }; + fs::create_dir_all(&root)?; + anyhow::ensure!( + fs::symlink_metadata(&root)?.is_dir(), + "Invalid payment admission directory" + ); + fs::set_permissions(&root, fs::Permissions::from_mode(0o700))?; + let file = fs::OpenOptions::new() + .read(true) + .write(true) + .create(true) + .mode(0o600) + .custom_flags(libc::O_NOFOLLOW | libc::O_NONBLOCK) + .open(root.join(key))?; + anyhow::ensure!(file.metadata()?.is_file(), "Invalid payment admission lock"); + loop { + if unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX) } == 0 { + break; + } + let error = std::io::Error::last_os_error(); + if error.kind() != std::io::ErrorKind::Interrupted { + return Err(error.into()); + } + } + Ok(Guard { _file: file }) + }) + .await? +} diff --git a/core/archipelago/src/content_purchase_download.rs b/core/archipelago/src/content_purchase_download.rs index 15afe3c2..dd0179b3 100644 --- a/core/archipelago/src/content_purchase_download.rs +++ b/core/archipelago/src/content_purchase_download.rs @@ -79,7 +79,7 @@ pub(crate) async fn cache( Ok(owned) } -fn verified_stream( +pub(crate) fn verified_stream( stream: S, expected_hash: String, expected_size: u64, diff --git a/core/archipelago/src/content_server.rs b/core/archipelago/src/content_server.rs index 42c65ebe..f25d4cf7 100644 --- a/core/archipelago/src/content_server.rs +++ b/core/archipelago/src/content_server.rs @@ -1923,3 +1923,41 @@ pub(crate) async fn publish_snapshot_offer( ) .await } + +/// Commit a new invoice only while its exact selected share remains current. +/// Caller holds the invoice journal; catalog writers do not acquire that journal. +pub(crate) async fn publish_snapshot_invoice( + data_dir: &Path, + original: &ContentItem, + journal: &crate::content_lightning::Journal, + binding: crate::content_lightning::Binding, + retained: crate::content_lightning::RetainedFile, +) -> Result { + let _held = CATALOG_WRITES.lock().await; + let catalog = load_catalog(data_dir).await?; + let current = catalog + .items + .iter() + .find(|item| item.id == original.id) + .context("Content was unshared before invoice preparation")?; + anyhow::ensure!( + serde_json::to_value(current)? == serde_json::to_value(original)?, + "Shared content terms changed before invoice preparation" + ); + anyhow::ensure!( + binding.content_id == original.id + && retained.filename == original.filename + && retained.mime_type == original.mime_type + && retained.size == original.size_bytes + && matches!(&original.access, AccessControl::Paid { price_sats, .. } if *price_sats == binding.price_sats) + && method_accepted(&original.access, "lightning"), + "Invoice snapshot terms changed" + ); + let visible = match &original.availability { + Availability::Nobody => false, + Availability::AllPeers => true, + Availability::Specific { peers } => peers.contains(&binding.buyer_did), + }; + anyhow::ensure!(visible, "Item is not shared with this invoice buyer"); + journal.prepare_seller_source(binding, Some(retained)) +} diff --git a/core/archipelago/src/main.rs b/core/archipelago/src/main.rs index 314937c3..9eb52fe9 100644 --- a/core/archipelago/src/main.rs +++ b/core/archipelago/src/main.rs @@ -44,6 +44,8 @@ mod content_auth; mod content_hash; mod content_indeehub; mod content_invoice; +mod content_lightning; +mod content_payment_admission; mod content_owned; mod content_purchase; mod content_purchase_executor; diff --git a/neode-ui/src/views/PeerFiles.vue b/neode-ui/src/views/PeerFiles.vue index 184f7065..78c178ae 100644 --- a/neode-ui/src/views/PeerFiles.vue +++ b/neode-ui/src/views/PeerFiles.vue @@ -393,6 +393,9 @@ {{ payItem.filename.split('/').pop() }} · {{ getItemPrice(payItem.access) }} sats

+ +

Original Lightning purchase · {{ lnReceipt.price_sats }} sats

+
@@ -421,8 +424,8 @@ - {{ lnPaying ? (hasBlockingLightningReceipt ? 'Checking payment…' : 'Paying…') : (hasBlockingLightningReceipt ? 'Check payment / retry download' : 'Pay with my Lightning node') }} - {{ hasBlockingLightningReceipt ? 'Checks the saved attempt; does not send more sats' : 'Pays the seller’s invoice from your node’s Lightning wallet' }} + {{ lnPaying ? (hasBlockingLightningReceipt ? 'Checking payment…' : 'Paying…') : (canRetryNative ? `Retry Lightning · ${lnReceipt?.price_sats} sats` : hasBlockingLightningReceipt ? (lnReceipt?.operation_id ? 'Resume original Lightning purchase' : 'Check payment / retry download') : 'Pay with my Lightning node') }} + {{ canRetryNative ? 'Creates a new attempt only after the original failed' : hasBlockingLightningReceipt ? (lnReceipt?.operation_id ? 'Recovers or completes the same purchase without creating another invoice' : 'Checks the saved attempt; does not send more sats') : 'Pays the seller’s invoice from your node’s Lightning wallet' }} @@ -585,7 +588,9 @@
+

Saved request {{ invoiceOperationId.slice(0, 8) }} · recovery uses this same invoice

{{ invoiceError }}

+