Add durable buyer-bound Lightning recovery and explicit native retry

Preserve original invoice preimages, private snapshots and exposure provenance; serialize rail admission and retire native-only failures before explicit replacement. Qualify 55 focused UI tests and vue-tsc. Expanded 17 engine cases and combined backend acceptance remain pending; six earlier engine cases passed in isolation. No live payment or publication.
This commit is contained in:
archipelago
2026-10-07 00:32:39 -04:00
parent ed96df0ac3
commit fba3273c67
18 changed files with 2525 additions and 109 deletions
+1
View File
@@ -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
@@ -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<Body>,
) -> Result<Response<Body>> {
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<std::sync::atomic::AtomicBool>);
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<std::sync::atomic::AtomicBool>);
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)?),
))
}
}
+4
View File
@@ -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!(
@@ -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,
@@ -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<u64>,
operation_id: Option<String>,
#[serde(default)]
external_exposure: bool,
}
impl RpcHandler {
pub(super) async fn handle_lightning_operation(
&self,
params: Option<serde_json::Value>,
action: &str,
) -> Result<serde_json::Value> {
let params: Params = serde_json::from_value(params.context("Missing invoice operation")?)?;
let peer =
crate::federation::load_unique_payment_peer(&self.config.data_dir, &params.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,
&params.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, &params.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) = &params.operation_id {
journal.buyer(id)?
} else {
journal.buyer_for(&buyer, &peer.did, &params.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), &params.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(())
}
}
@@ -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<u64> {
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<Option<Invoice>> {
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<Status> {
// 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<serde_json::Value> {
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<PreparedNativePayment> {
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<serde_json::Value> {
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::Prepared> {
self.0
.prepare_bound_invoice_payment(invoice, hash, amount)
.await
}
async fn lookup_payment(&self, hash: &str) -> Result<serde_json::Value> {
self.0
.handle_lnd_paymentstatus(Some(serde_json::json!({"payment_hash":hash})))
.await
}
}
+1
View File
@@ -1,4 +1,5 @@
mod channels;
pub(super) mod external_invoice;
mod fee_bump;
mod fee_policy;
mod info;
+1 -1
View File
@@ -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,
+8
View File
@@ -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"
+10
View File
@@ -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(),
&params.content_id,
)
.await?;
self.ensure_invoice_allows_other_rail(&buyer, transport.seller_did(), &params.content_id)
.await?;
let result = caller::purchase(
&self.config.data_dir,
&buyer,
File diff suppressed because it is too large Load Diff
@@ -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<Guard> {
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?
}
@@ -79,7 +79,7 @@ pub(crate) async fn cache(
Ok(owned)
}
fn verified_stream<S, E>(
pub(crate) fn verified_stream<S, E>(
stream: S,
expected_hash: String,
expected_size: u64,
+38
View File
@@ -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<crate::content_lightning::SellerRecord> {
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))
}
+2
View File
@@ -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;