Integrate two-phase on-chain purchase recovery

This commit is contained in:
archipelago
2026-10-07 07:49:02 -04:00
parent 558f097fd6
commit 6547ae05fa
18 changed files with 6084 additions and 113 deletions
@@ -0,0 +1,712 @@
use super::{build_response, ApiHandler};
use crate::{content_lightning::Binding, content_onchain_seller::Journal};
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/onchain/v1/operation";
#[derive(Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct Operation {
pub binding: Binding,
pub action: String,
}
// Load wallet credentials only after authenticated request validation reaches a
// wallet operation. Tests inject the same typed boundary without live services.
struct NativeSellerWallet<'a>(&'a crate::api::rpc::RpcHandler);
impl crate::content_onchain_seller::Wallet for NativeSellerWallet<'_> {
async fn network(&self) -> Result<crate::content_onchain::ChainNetwork> {
self.0.onchain_purchase_wallet().await?.network().await
}
async fn preflight(&self, network: crate::content_onchain::ChainNetwork) -> Result<()> {
self.0
.onchain_purchase_wallet()
.await?
.preflight(network)
.await
}
async fn allocate(&self) -> Result<String> {
self.0.onchain_purchase_wallet().await?.allocate().await
}
async fn received(&self, address: &str, amount: u64) -> Result<bool> {
self.0
.onchain_purchase_wallet()
.await?
.received(address, amount)
.await
}
}
impl ApiHandler {
pub(super) async fn handle_onchain_purchase(
&self,
request: Request<Body>,
) -> Result<Response<Body>> {
self.handle_onchain_purchase_with_wallet(request, &NativeSellerWallet(&self.rpc_handler))
.await
}
async fn handle_onchain_purchase_with_wallet<W: crate::content_onchain_seller::Wallet>(
&self,
mut request: Request<Body>,
wallet: &W,
) -> Result<Response<Body>> {
anyhow::ensure!(
request.method() == Method::POST && request.uri().path() == ROUTE,
"Invalid on-chain purchase 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,
"On-chain purchase request too large"
);
bytes.extend_from_slice(&chunk)
}
Ok::<_, anyhow::Error>(bytes)
})
.await
.context("On-chain purchase 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,
"On-chain purchase peer identity mismatch"
);
anyhow::ensure!(
matches!(
operation.action.as_str(),
"create" | "offer" | "allocate" | "status" | "download" | "cancel"
),
"Invalid on-chain purchase action"
);
let binding = &operation.binding;
let journal = Journal::open(&self.config.data_dir).await?;
let retired = if operation.action == "cancel" {
Some(journal.retire_unallocated(binding)?)
} else {
journal.retirement(binding)?
};
if let Some(ack) = retired {
return Ok(build_response(
StatusCode::OK,
"application/json",
Body::from(serde_json::to_vec(&ack)?),
));
}
let mut saved = journal.load(binding)?;
if saved.is_none() {
anyhow::ensure!(
matches!(operation.action.as_str(), "create" | "offer"),
"Unknown original on-chain purchase 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, "onchain"),
"On-chain purchase 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 on-chain purchase 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 on-chain purchase"
);
// Source metadata is private and committed before address allocation.
let record = crate::content_server::publish_snapshot_onchain(
&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(),
},
wallet.network().await?,
)
.await?;
saved = Some(record);
}
saved.context("Missing original on-chain operation")?;
let status = if operation.action == "allocate" {
crate::content_server::allocate_onchain_offer(
&self.config.data_dir,
&journal,
binding,
wallet,
)
.await?
} else {
crate::content_onchain_seller::drive(&journal, binding, false, wallet).await?
};
if operation.action == "download" {
anyhow::ensure!(status.paid, "Original on-chain purchase has not settled");
let source = &status.source;
let data = self.config.data_dir.clone();
let id = binding.content_id.clone();
let retained = source.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 on-chain purchase 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)?),
))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::content_onchain_seller::{Allocation, UnallocatedAck};
use hyper::service::{make_service_fn, service_fn};
use std::{convert::Infallible, sync::Arc};
#[derive(Default)]
struct MockWallet {
allocations: std::sync::atomic::AtomicUsize,
lose_reply: std::sync::atomic::AtomicBool,
}
impl crate::content_onchain_seller::Wallet for MockWallet {
async fn network(&self) -> Result<crate::content_onchain::ChainNetwork> {
Ok(crate::content_onchain::ChainNetwork::Regtest)
}
async fn preflight(&self, _: crate::content_onchain::ChainNetwork) -> Result<()> {
Ok(())
}
async fn allocate(&self) -> Result<String> {
self.allocations
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
anyhow::ensure!(
!self
.lose_reply
.swap(false, std::sync::atomic::Ordering::SeqCst),
"Simulated lost allocation response"
);
let mut bytes = vec![0, 20];
bytes.extend([17u8; 20]);
Ok(bitcoin::Address::from_script(
&bitcoin::ScriptBuf::from_bytes(bytes),
bitcoin::Network::Regtest,
)?
.to_string())
}
async fn received(&self, _: &str, _: u64) -> Result<bool> {
Ok(false)
}
}
struct HttpFixture {
wallet: Arc<MockWallet>,
data: tempfile::TempDir,
_buyer_data: tempfile::TempDir,
buyer: crate::identity::NodeIdentity,
seller: String,
url: String,
task: tokio::task::JoinHandle<()>,
}
impl Drop for HttpFixture {
fn drop(&mut self) {
self.task.abort();
}
}
async fn fixture() -> HttpFixture {
let data = tempfile::tempdir().unwrap();
let buyer_data = tempfile::tempdir().unwrap();
let buyer = crate::identity::NodeIdentity::load_or_create(buyer_data.path())
.await
.unwrap();
let mut config = crate::config::Config::default();
config.data_dir = data.path().to_path_buf();
let handler = Arc::new(
ApiHandler::new(
config,
Arc::new(crate::state::StateManager::new()),
Arc::new(crate::monitoring::MetricsStore::new()),
None,
None,
)
.await
.unwrap(),
);
let seller = crate::identity::did_key_from_pubkey_hex(&handler.self_pubkey_hex).unwrap();
let wallet = Arc::new(MockWallet::default());
let server_wallet = wallet.clone();
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
listener.set_nonblocking(true).unwrap();
let url = format!("http://{}", listener.local_addr().unwrap());
let server = hyper::Server::from_tcp(listener)
.unwrap()
.serve(make_service_fn(move |_| {
let handler = handler.clone();
let wallet = server_wallet.clone();
async move {
Ok::<_, Infallible>(service_fn(move |request| {
let handler = handler.clone();
let wallet = wallet.clone();
async move {
Ok::<_, Infallible>(
handler
.handle_onchain_purchase_with_wallet(request, wallet.as_ref())
.await
.unwrap_or_else(|_| {
build_response(
StatusCode::BAD_REQUEST,
"application/json",
Body::from("{\"error\":\"rejected\"}"),
)
}),
)
}
}))
}
}));
let task = tokio::spawn(async move {
server.await.unwrap();
});
HttpFixture {
wallet,
data,
_buyer_data: buyer_data,
buyer,
seller,
url,
task,
}
}
impl HttpFixture {
fn binding(&self) -> Binding {
Binding {
id: uuid::Uuid::new_v4().to_string(),
buyer_did: self.buyer.did_key().unwrap(),
seller_did: self.seller.clone(),
content_id: "file".into(),
price_sats: 546,
}
}
async fn send(
&self,
body: &[u8],
signed_body: Option<&[u8]>,
audience: Option<&str>,
) -> reqwest::Response {
let mut request = reqwest::Client::new()
.post(format!("{}{}", self.url, ROUTE))
.header("content-type", "application/json")
.body(body.to_vec());
if let Some(signed) = signed_body {
let proof = crate::content_auth::sign_request(
&self.buyer,
audience.unwrap_or(&self.seller),
&Method::POST,
ROUTE,
signed,
chrono::Utc::now().timestamp(),
)
.unwrap();
request = request.header(crate::content_auth::REQUEST_HEADER, proof);
}
request.send().await.unwrap()
}
async fn operation(&self, binding: &Binding, action: &str) -> reqwest::Response {
let body = serde_json::to_vec(&Operation {
binding: binding.clone(),
action: action.into(),
})
.unwrap();
self.send(&body, Some(&body), None).await
}
}
#[tokio::test]
async fn authenticated_cancel_roundtrip_lost_reply_and_delayed_create_return_same_retirement() {
let server = fixture().await;
let binding = server.binding();
// Drop the original reply after headers: terminal state must already be durable.
let first = server.operation(&binding, "cancel").await;
assert_eq!(first.status(), reqwest::StatusCode::OK);
drop(first);
let replay = server.operation(&binding, "cancel").await;
assert_eq!(replay.status(), reqwest::StatusCode::OK);
let ack: UnallocatedAck = replay.json().await.unwrap();
ack.validate(&binding).unwrap();
let delayed = server.operation(&binding, "create").await;
assert_eq!(delayed.status(), reqwest::StatusCode::OK);
assert_eq!(delayed.json::<UnallocatedAck>().await.unwrap(), ack);
let journal = Journal::open(server.data.path()).await.unwrap();
assert_eq!(journal.retirement(&binding).unwrap(), Some(ack));
assert!(journal.load(&binding).unwrap().is_none());
assert!(!server.data.path().join("content-snapshots").exists());
}
#[tokio::test]
async fn cancellation_http_rejects_missing_proof_body_tamper_and_wrong_seller_without_tombstone(
) {
let server = fixture().await;
let binding = server.binding();
let body = serde_json::to_vec(&Operation {
binding: binding.clone(),
action: "cancel".into(),
})
.unwrap();
assert!(!server.send(&body, None, None).await.status().is_success());
let mut changed = binding.clone();
changed.price_sats += 1;
let changed = serde_json::to_vec(&Operation {
binding: changed,
action: "cancel".into(),
})
.unwrap();
assert!(!server
.send(&changed, Some(&body), None)
.await
.status()
.is_success());
let wrong = crate::identity::did_key_from_pubkey_hex(&hex::encode([8; 32])).unwrap();
assert!(!server
.send(&body, Some(&body), Some(&wrong))
.await
.status()
.is_success());
let journal = Journal::open(server.data.path()).await.unwrap();
assert!(journal.retirement(&binding).unwrap().is_none());
}
#[tokio::test]
async fn authenticated_cancel_cannot_retire_dispatched_or_issued_address() {
let server = fixture().await;
let mut script = vec![0, 20];
script.extend([1; 20]);
let address = bitcoin::Address::from_script(
&bitcoin::ScriptBuf::from_bytes(script),
bitcoin::Network::Regtest,
)
.unwrap()
.to_string();
for allocation in [Allocation::Dispatched, Allocation::Ready { address }] {
let binding = server.binding();
let journal = Journal::open(server.data.path()).await.unwrap();
let mut record = journal
.prepare(
binding.clone(),
crate::content_lightning::RetainedFile {
sha256: "a".repeat(64),
size: 4,
filename: "original.txt".into(),
mime_type: "text/plain".into(),
},
crate::content_onchain::ChainNetwork::Regtest,
)
.unwrap();
record.allocation = allocation.clone();
journal.save(&record).unwrap();
drop(journal);
assert!(!server
.operation(&binding, "cancel")
.await
.status()
.is_success());
let journal = Journal::open(server.data.path()).await.unwrap();
assert!(journal.retirement(&binding).unwrap().is_none());
assert_eq!(
journal.load(&binding).unwrap().unwrap().allocation,
allocation
);
}
}
async fn seed_unallocated_offer(server: &HttpFixture) -> Binding {
let binding = server.binding();
crate::content_server::save_catalog(
server.data.path(),
&crate::content_server::ContentCatalog {
items: vec![crate::content_server::ContentItem {
id: binding.content_id.clone(),
filename: "original.txt".into(),
mime_type: "text/plain".into(),
size_bytes: 4,
description: String::new(),
added_at: String::new(),
availability: crate::content_server::Availability::AllPeers,
access: crate::content_server::AccessControl::Paid {
price_sats: 546,
accepted: vec!["onchain".into()],
},
}],
},
)
.await
.unwrap();
let root = server.data.path().join("content/files");
std::fs::create_dir_all(&root).unwrap();
std::fs::write(root.join("original.txt"), b"test").unwrap();
let cancelled = std::sync::atomic::AtomicBool::new(false);
let snapshot = crate::content_snapshot::prepare(
server.data.path(),
&root,
&binding.content_id,
std::path::Path::new("original.txt"),
&crate::media_registration::Limits {
max_bytes: 1024,
cancelled: &cancelled,
},
1024 * 1024,
0,
|_| Ok(()),
)
.unwrap();
let journal = Journal::open(server.data.path()).await.unwrap();
journal
.prepare(
binding.clone(),
crate::content_lightning::RetainedFile {
sha256: snapshot.sha256,
size: 4,
filename: "original.txt".into(),
mime_type: "text/plain".into(),
},
crate::content_onchain::ChainNetwork::Regtest,
)
.unwrap();
binding
}
#[tokio::test]
async fn authenticated_offer_never_allocates_or_returns_a_receive_address() {
let server = fixture().await;
let binding = seed_unallocated_offer(&server).await;
let result = server.operation(&binding, "offer").await;
assert_eq!(result.status(), reqwest::StatusCode::OK);
let body: serde_json::Value = result.json().await.unwrap();
assert_eq!(body["allocation"]["state"], "prepared");
assert!(body["allocation"].get("address").is_none());
assert!(body.get("address").is_none());
let journal = Journal::open(server.data.path()).await.unwrap();
assert_eq!(
journal.load(&binding).unwrap().unwrap().allocation,
Allocation::Prepared
);
}
#[tokio::test]
async fn reviewed_offer_can_cancel_and_delayed_explicit_allocate_cannot_revive_it() {
let server = fixture().await;
let binding = seed_unallocated_offer(&server).await;
assert_eq!(
server.operation(&binding, "offer").await.status(),
reqwest::StatusCode::OK
);
let retired: UnallocatedAck = server
.operation(&binding, "cancel")
.await
.json()
.await
.unwrap();
retired.validate(&binding).unwrap();
// Represents a delayed Pay request from the old modal after cancellation.
let late = server.operation(&binding, "allocate").await;
assert_eq!(late.status(), reqwest::StatusCode::OK);
assert_eq!(late.json::<UnallocatedAck>().await.unwrap(), retired);
let journal = Journal::open(server.data.path()).await.unwrap();
assert_eq!(
journal.load(&binding).unwrap().unwrap().allocation,
Allocation::Prepared
);
assert_eq!(journal.retirement(&binding).unwrap(), Some(retired));
}
#[tokio::test]
async fn changing_authenticated_offer_body_to_allocate_cannot_dispatch_an_address() {
let server = fixture().await;
let binding = seed_unallocated_offer(&server).await;
let reviewed = serde_json::to_vec(&Operation {
binding: binding.clone(),
action: "offer".into(),
})
.unwrap();
let changed = serde_json::to_vec(&Operation {
binding: binding.clone(),
action: "allocate".into(),
})
.unwrap();
assert!(!server
.send(&changed, Some(&reviewed), None)
.await
.status()
.is_success());
let journal = Journal::open(server.data.path()).await.unwrap();
assert_eq!(
journal.load(&binding).unwrap().unwrap().allocation,
Allocation::Prepared
);
assert!(journal.retirement(&binding).unwrap().is_none());
}
#[tokio::test]
async fn explicit_allocation_reuses_original_address_after_lost_http_reply() {
let server = fixture().await;
let binding = seed_unallocated_offer(&server).await;
assert!(server
.operation(&binding, "offer")
.await
.status()
.is_success());
assert_eq!(
server
.wallet
.allocations
.load(std::sync::atomic::Ordering::SeqCst),
0
);
// Caller loses the response after seller durability; recovery returns the same record.
drop(server.operation(&binding, "allocate").await);
let recovered: crate::content_onchain_seller::Record = server
.operation(&binding, "allocate")
.await
.json()
.await
.unwrap();
assert!(recovered.quote().unwrap().is_some());
let repeated: crate::content_onchain_seller::Record = server
.operation(&binding, "allocate")
.await
.json()
.await
.unwrap();
assert_eq!(recovered, repeated);
assert_eq!(
server
.wallet
.allocations
.load(std::sync::atomic::Ordering::SeqCst),
1
);
assert!(!server
.operation(&binding, "cancel")
.await
.status()
.is_success());
}
#[tokio::test]
async fn lost_wallet_allocation_reply_never_allocates_a_second_address() {
let server = fixture().await;
let binding = seed_unallocated_offer(&server).await;
server
.wallet
.lose_reply
.store(true, std::sync::atomic::Ordering::SeqCst);
assert!(!server
.operation(&binding, "allocate")
.await
.status()
.is_success());
let recovered: crate::content_onchain_seller::Record = server
.operation(&binding, "allocate")
.await
.json()
.await
.unwrap();
assert_eq!(recovered.allocation, Allocation::Dispatched);
assert!(recovered.quote().unwrap().is_none());
assert_eq!(
server
.wallet
.allocations
.load(std::sync::atomic::Ordering::SeqCst),
1
);
assert!(!server
.operation(&binding, "cancel")
.await
.status()
.is_success());
}
}
@@ -337,6 +337,14 @@ 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.onchain-cancel" => self.handle_onchain_operation(params, "cancel").await,
"content.onchain-attempt" => self.handle_onchain_operation(params, "lookup").await,
"content.onchain-create" => self.handle_onchain_operation(params, "create").await,
"content.onchain-expose" => self.handle_onchain_operation(params, "expose").await,
"content.onchain-prepare" => self.handle_onchain_operation(params, "prepare").await,
"content.onchain-pay" => self.handle_onchain_operation(params, "pay").await,
"content.onchain-recover" => self.handle_onchain_operation(params, "status").await,
"content.onchain-download" => self.handle_onchain_operation(params, "download").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,
+1
View File
@@ -4,6 +4,7 @@ mod fee_bump;
mod fee_policy;
mod info;
mod macaroons;
pub(super) mod onchain_purchase;
mod payments;
mod seed_backup;
mod wallet;
File diff suppressed because it is too large Load Diff
+1 -1
View File
@@ -1347,7 +1347,7 @@ fn psbt_key_origin_report(psbt_base64: &str) -> Result<PsbtKeyOriginReport> {
/// LND's transaction `amount` is the wallet-wide net amount, not the value
/// paid to a purchase address. Attribute only confirmed output values, once
/// per outpoint. Missing/malformed evidence is unknown, never proof of payment.
fn confirmed_address_sats(body: &serde_json::Value, address: &str) -> Result<u64> {
pub(super) fn confirmed_address_sats(body: &serde_json::Value, address: &str) -> Result<u64> {
use std::collections::{HashMap, HashSet};
const MAX_SATS: u64 = 21_000_000 * 100_000_000;
fn integer(value: &serde_json::Value) -> Result<u64> {
+18
View File
@@ -18,6 +18,7 @@ mod handshake;
mod identity;
mod interfaces;
mod lightning_purchase;
mod onchain_purchase;
pub(crate) mod lnd;
mod marketplace;
mod media_registration;
@@ -110,6 +111,15 @@ fn native_consent_origin_allowed(method: &str, headers: &hyper::HeaderMap, dev_m
| "media.registration.context"
| "media.registration.resolve"
| "content.rental-purchase"
| "content.onchain-cancel"
| "content.onchain-attempt"
| "content.onchain-create"
| "content.onchain-expose"
| "content.onchain-prepare"
| "content.onchain-pay"
| "content.onchain-recover"
| "content.onchain-download"
| "content.invoice-pay"
| "content.invoice-download"
| "content.invoice-attempt"
@@ -837,6 +847,14 @@ mod nostr_signing_origin_tests {
"media.registration.context",
"media.registration.resolve",
"content.rental-purchase",
"content.onchain-cancel",
"content.onchain-attempt",
"content.onchain-create",
"content.onchain-expose",
"content.onchain-prepare",
"content.onchain-pay",
"content.onchain-recover",
"content.onchain-download",
"content.purchase",
"content.cancel-purchase",
"content.playback-handle",
@@ -0,0 +1,668 @@
//! Owner-only original-operation on-chain flow. No generic sendcoins fallback.
use super::RpcHandler;
use crate::{
content_lightning::Binding,
content_onchain::{self as engine, Journal, Phase, Record},
};
use anyhow::{Context, Result};
use serde::Deserialize;
use serde_json::{json, Value};
use sha2::{Digest, Sha256};
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct Params {
onion: String,
content_id: String,
operation_id: Option<String>,
price_sats: Option<u64>,
max_fee_sats: Option<u64>,
sat_per_vbyte: Option<u64>,
template_sha256: Option<String>,
plan_sha256: Option<String>,
}
fn public(record: &Record) -> Result<Value> {
let fee = record
.template
.as_ref()
.map(|t| engine::validate_funded(record, t))
.transpose()?;
let template_sha256 = record
.template
.as_ref()
.map(|t| hex::encode(Sha256::digest(t.psbt_base64.as_bytes())));
Ok(
json!({"operation_id":record.binding.id,"price_sats":record.binding.price_sats,"phase":record.phase,
"network":record.network(),"external_exposure":record.externally_exposed,
"address":if record.externally_exposed {record.quote.as_ref().map(|q|q.address.as_str())}else{None},
"fee_sats":fee.or_else(|| record.plan.as_ref().map(|p|p.fee_sats)),
"max_fee_sats":record.policy.as_ref().map(|p|p.max_fee_sats).or_else(||record.plan.as_ref().map(|p|p.max_fee_sats)),
"template_sha256":template_sha256,"plan_sha256":record.plan.as_ref().map(|p|p.hash()).transpose()?,
"txid":record.signed.as_ref().map(|s|s.txid.as_str()),"paid":record.settled,
"change_allocation_ambiguous":matches!(record.change_address,Some(engine::ChangeAddress::Dispatched)),
"can_switch_method":record.retirement.is_some(),"retired_unallocated":record.retirement.is_some()}),
)
}
impl RpcHandler {
async fn request_onchain_allocation(
&self,
record: &Record,
fips: &str,
) -> Result<crate::content_onchain_seller::Record> {
let operation = crate::api::handler::onchain_purchase::Operation {
binding: record.binding.clone(),
action: "allocate".into(),
};
let (mut response, _) = crate::fips::dial::PeerRequest::new(
Some(fips),
&record.seller_onion,
crate::api::handler::onchain_purchase::ROUTE,
)
.require_fips()
.single_delivery()
.timeout(std::time::Duration::from_secs(45))
.send_content_json(
&self.config.data_dir,
&record.binding.seller_did,
&operation,
)
.await
.context("Original seller allocation reply unavailable; recover the same operation")?;
anyhow::ensure!(
response.status().is_success(),
"Original seller allocation remains unresolved"
);
let mut bytes = Vec::new();
while let Some(chunk) = response.chunk().await? {
anyhow::ensure!(
bytes.len() + chunk.len() <= 16384,
"Seller response too large"
);
bytes.extend_from_slice(&chunk);
}
let saved: crate::content_onchain_seller::Record = serde_json::from_slice(&bytes)?;
anyhow::ensure!(
saved.binding == record.binding,
"Seller changed original purchase"
);
if let Some(offer) = &record.offer {
anyhow::ensure!(saved.offer()? == *offer, "Seller changed original offer");
}
Ok(saved)
}
pub(super) async fn ensure_onchain_allows_other_rail(
&self,
buyer: &str,
seller: &str,
content: &str,
) -> Result<()> {
anyhow::ensure!(
Journal::find_for(&self.config.data_dir, buyer, seller, content)?.is_none(),
"An original on-chain purchase remains recoverable; do not pay again or switch methods"
);
Ok(())
}
pub(super) async fn handle_onchain_operation(
&self,
params: Option<Value>,
action: &str,
) -> Result<Value> {
let params: Params = serde_json::from_value(params.context("Missing on-chain operation")?)?;
anyhow::ensure!(
!params.content_id.starts_with("registered_"),
"Registered rentals require their native purchase contract"
);
let peer =
crate::federation::load_unique_payment_peer(&self.config.data_dir, &params.onion)
.await?;
let buyer =
crate::identity::NodeIdentity::load_existing(&self.config.data_dir.join("identity"))
.await?
.did_key()?;
anyhow::ensure!(buyer != peer.did, "Cannot buy from this same node");
let _admission = crate::content_payment_admission::lock(
&self.config.data_dir,
&buyer,
&peer.did,
&params.content_id,
)
.await?;
let original = if let Some(id) = &params.operation_id {
let journal = Journal::open(&self.config.data_dir, id).await?;
let original = journal.load()?;
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,
"Original on-chain operation belongs to another purchase"
);
}
original
} else {
Journal::find_for(&self.config.data_dir, &buyer, &peer.did, &params.content_id)?
};
if action == "lookup" {
return Ok(json!({"attempt":original.as_ref().map(public).transpose()?}));
}
if let Some(id) = &params.operation_id {
anyhow::ensure!(
original.as_ref().is_some_and(|r| &r.binding.id == id),
"Original on-chain operation changed"
);
}
let mut record = if let Some(record) = original {
record
} else {
anyhow::ensure!(
matches!(action, "create" | "expose") && params.operation_id.is_none(),
"Recover original on-chain operation first"
);
self.ensure_invoice_allows_other_rail(&buyer, &peer.did, &params.content_id)
.await?;
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 original Cashu purchase first"
);
Record::new(
Binding {
id: uuid::Uuid::new_v4().to_string(),
buyer_did: buyer,
seller_did: peer.did.clone(),
content_id: params.content_id.clone(),
price_sats: params.price_sats.context("Expected price required")?,
},
params.onion.clone(),
)?
};
anyhow::ensure!(
record.seller_onion == params.onion
&& params
.price_sats
.is_none_or(|p| p == record.binding.price_sats),
"Original payment address or price changed"
);
let journal = Journal::open(&self.config.data_dir, &record.binding.id).await?;
if journal.load()?.is_none() {
journal.save(&record)?;
}
if record.retirement.is_some() {
return public(&record);
}
if action == "cancel" {
anyhow::ensure!(params.operation_id.is_some() && record.can_retire_unallocated(),"An allocated or mutated on-chain purchase cannot be canceled; recover its original payment");
}
if matches!(action, "create" | "expose" | "prepare" | "pay") && !record.settled {
// Recheck while the same admission guard is held, including resumes
// from another window and records predating this owner flow.
self.ensure_invoice_allows_other_rail(
&record.binding.buyer_did,
&peer.did,
&params.content_id,
)
.await?;
let cashu = crate::content_purchase::Journal::open(&self.config.data_dir).await?;
anyhow::ensure!(
cashu
.find_buyers(&record.binding.buyer_did, &peer.did, &params.content_id)
.await?
.iter()
.all(|r| r.phase == crate::content_purchase::BuyerPhase::Cancelled),
"Another saved Cashu liability must be recovered before on-chain dispatch"
);
}
if action == "prepare" {
anyhow::ensure!(
params.operation_id.is_some(),
"Original operation ID required"
);
if record.template.is_none() && record.plan.is_none() {
let max_fee_sats = params
.max_fee_sats
.context("Explicit maximum fee required")?;
anyhow::ensure!(
(1..=2_100_000_000_000_000).contains(&max_fee_sats),
"Invalid maximum fee"
);
if let Some(rate) = params.sat_per_vbyte {
anyhow::ensure!((1..=5000).contains(&rate), "Invalid fee rate");
}
let wallet = self.onchain_purchase_wallet().await?;
let change = wallet.prepare_change(&journal).await?;
record = wallet
.prepare_plan(
&journal,
super::lnd::onchain_purchase::PlanRequest {
change_address: change,
max_fee_sats,
sat_per_vbyte: params.sat_per_vbyte,
},
)
.await?;
}
return public(&record);
}
if action == "pay" {
anyhow::ensure!(
params.operation_id.is_some(),
"Original operation ID required"
);
if record.settled {
return public(&record);
}
if let Some(plan) = &record.plan {
anyhow::ensure!(
params.plan_sha256.as_deref() == Some(plan.hash()?.as_str()),
"Confirm the original saved funding plan before payment"
);
} else {
let template = record
.template
.as_ref()
.context("Review the original fee first")?;
anyhow::ensure!(
params.template_sha256.as_deref()
== Some(
hex::encode(Sha256::digest(template.psbt_base64.as_bytes())).as_str()
),
"Confirm the original saved transaction before payment"
);
}
let wallet = self.onchain_purchase_wallet().await?;
if record.plan.is_some() && record.quote.is_none() {
record = engine::lease_plan(&journal, &wallet).await?;
engine::mark_address_allocation(&journal, false)?;
let status = self
.request_onchain_allocation(
&record,
peer.fips_npub
.as_deref()
.context("Seller has no authenticated mesh connection")?,
)
.await?;
let quote = status
.quote()?
.context("Original seller allocation is unresolved; recover this operation")?;
record = engine::accept_quote(&journal, quote)?;
}
if record.plan.is_some() && record.template.is_none() {
record = engine::bind_plan(&journal)?;
}
if matches!(
record.phase,
Phase::TemplatePrepared | Phase::LeaseDispatched
) {
record = engine::drive(&journal, &wallet, engine::Action::Lease, None).await?;
}
if matches!(record.phase, Phase::Funded | Phase::SigningDispatched) {
record = engine::drive(&journal, &wallet, engine::Action::Sign, None).await?;
}
if matches!(
record.phase,
Phase::Signed | Phase::BroadcastDispatched | Phase::Published
) {
record = engine::drive(&journal, &wallet, engine::Action::Publish, None).await?;
}
return public(&record);
}
anyhow::ensure!(
matches!(
action,
"create" | "status" | "expose" | "download" | "cancel"
),
"Unsupported on-chain action"
);
let fips = peer
.fips_npub
.context("Seller has no authenticated mesh connection")?;
if action == "expose" && record.offer.is_some() && record.quote.is_none() {
engine::mark_address_allocation(&journal, true)?;
let status = self.request_onchain_allocation(&record, &fips).await?;
record = engine::accept_quote(
&journal,
status
.quote()?
.context("Original seller allocation is unresolved; recover this operation")?,
)?;
}
let remote_action = if action == "create" || action == "expose" && record.offer.is_none() {
"offer"
} else if action == "expose" {
"status"
} else {
action
};
let operation = crate::api::handler::onchain_purchase::Operation {
binding: record.binding.clone(),
action: remote_action.into(),
};
let route = crate::api::handler::onchain_purchase::ROUTE;
let remote = 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 remote {
Ok(value) => value,
Err(_) => {
return Ok(
json!({"attempt":public(&record)?,"recovery_required":true,"error":"Original on-chain request is saved. Recover this operation; do not request another address or pay again."}),
)
}
};
anyhow::ensure!(
response.status().is_success(),
"Seller could not recover original on-chain purchase {}; retain it",
record.binding.id
);
if action == "download" {
let source = record
.quote
.as_ref()
.context("Recover original address first")?
.source
.clone();
anyhow::ensure!(
response.content_length() == Some(source.size),
"Original file length changed"
);
record.settled = true;
journal.save(&record)?;
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: "onchain".into(),
purchased_at: chrono::Utc::now().to_rfc3339(),
download_complete: false,
},
Box::pin(stream),
Some(source.size),
)
.await?;
return Ok(
json!({"owned":true,"owned_content_id":owned.content_id,"mime_type":owned.mime_type}),
);
}
let mut bytes = vec![];
while let Some(chunk) = response.chunk().await? {
anyhow::ensure!(
bytes.len() + chunk.len() <= 16384,
"On-chain response too large"
);
bytes.extend_from_slice(&chunk);
}
let body: Value = serde_json::from_slice(&bytes)?;
if body["state"] == "cancelled_unallocated" {
let ack: crate::content_onchain_seller::UnallocatedAck = serde_json::from_value(body)?;
record = engine::retire_unallocated(&journal, ack)?;
return public(&record);
}
anyhow::ensure!(
action != "cancel",
"Seller did not acknowledge unallocated retirement; preserve original operation"
);
let status: crate::content_onchain_seller::Record = serde_json::from_value(body)?;
anyhow::ensure!(
status.binding == record.binding,
"Seller changed original purchase"
);
if record.offer.is_none() && record.quote.is_none() {
record = engine::accept_offer(&journal, status.offer()?)?;
}
if let Some(quote) = status.quote()? {
record = engine::accept_quote(&journal, quote)?;
}
anyhow::ensure!(
!status.paid || record.quote.is_some(),
"Paid purchase lacks original address"
);
record.settled |= status.paid;
journal.save(&record)?;
if action == "expose" && !record.settled {
if record.quote.is_none() {
engine::mark_address_allocation(&journal, true)?;
let status = self.request_onchain_allocation(&record, &fips).await?;
record = engine::accept_quote(
&journal,
status.quote()?.context(
"Original seller allocation is unresolved; recover this operation",
)?,
)?;
}
engine::expose_address(&journal)?;
record = journal.load()?.context("Original record unavailable")?;
}
public(&record)
}
}
#[cfg(test)]
mod tests {
use super::*;
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: 546,
}
}
#[tokio::test]
async fn buyer_discovery_retains_unresolved_address_and_rejects_duplicate_operations() {
let data = tempfile::tempdir().unwrap();
let binding = binding();
let saved = Record::new(binding.clone(), format!("{}.onion", "a".repeat(56))).unwrap();
let j = Journal::open(data.path(), &binding.id).await.unwrap();
j.save(&saved).unwrap();
drop(j);
let found = Journal::find_for(
data.path(),
&binding.buyer_did,
&binding.seller_did,
&binding.content_id,
)
.unwrap()
.unwrap();
assert_eq!(found.binding.id, binding.id);
assert!(found.blocks_other_rails());
assert!(public(&found).unwrap()["address"].is_null());
assert!(Journal::find_for(
data.path(),
&binding.seller_did,
&binding.buyer_did,
&binding.content_id
)
.unwrap()
.is_none());
let mut second = binding.clone();
second.id = uuid::Uuid::new_v4().to_string();
let j = Journal::open(data.path(), &second.id).await.unwrap();
j.save(&Record::new(second, saved.seller_onion).unwrap())
.unwrap();
drop(j);
assert!(Journal::find_for(
data.path(),
&binding.buyer_did,
&binding.seller_did,
&binding.content_id
)
.is_err());
}
#[tokio::test]
async fn corrupted_node_record_cannot_be_treated_as_permission_to_pay_again() {
let data = tempfile::tempdir().unwrap();
let binding = binding();
let j = Journal::open(data.path(), &binding.id).await.unwrap();
j.save(&Record::new(binding.clone(), format!("{}.onion", "a".repeat(56))).unwrap())
.unwrap();
drop(j);
std::fs::write(
data.path()
.join("content-onchain")
.join(format!("{}.json", binding.id)),
b"{}",
)
.unwrap();
assert!(Journal::find_for(
data.path(),
&binding.buyer_did,
&binding.seller_did,
&binding.content_id
)
.is_err());
}
#[tokio::test]
async fn only_durable_matching_empty_ack_releases_cross_rail_and_stale_callback_cannot_revive()
{
let data = tempfile::tempdir().unwrap();
let binding = binding();
let _rail = crate::content_payment_admission::lock(
data.path(),
&binding.buyer_did,
&binding.seller_did,
&binding.content_id,
)
.await
.unwrap();
let journal = Journal::open(data.path(), &binding.id).await.unwrap();
let original = Record::new(binding.clone(), format!("{}.onion", "a".repeat(56))).unwrap();
journal.save(&original).unwrap();
let ack = crate::content_onchain_seller::UnallocatedAck {
binding: binding.clone(),
state: "cancelled_unallocated".into(),
address: serde_json::Value::Null,
allocation_dispatched: false,
can_switch_method: true,
};
for wrong in [
crate::content_onchain_seller::UnallocatedAck {
allocation_dispatched: true,
..ack.clone()
},
crate::content_onchain_seller::UnallocatedAck {
address: serde_json::json!("not-empty"),
..ack.clone()
},
crate::content_onchain_seller::UnallocatedAck {
binding: Binding {
id: uuid::Uuid::new_v4().to_string(),
..binding.clone()
},
..ack.clone()
},
] {
assert!(engine::retire_unallocated(&journal, wrong).is_err());
}
assert!(Journal::find_for(
data.path(),
&binding.buyer_did,
&binding.seller_did,
&binding.content_id
)
.unwrap()
.is_some());
let retired = engine::retire_unallocated(&journal, ack).unwrap();
assert!(!retired.blocks_other_rails());
assert!(public(&retired).unwrap()["can_switch_method"] == true);
assert!(journal.save(&original).is_err());
drop(journal);
assert!(Journal::find_for(
data.path(),
&binding.buyer_did,
&binding.seller_did,
&binding.content_id
)
.unwrap()
.is_none());
let mut replacement = binding.clone();
replacement.id = uuid::Uuid::new_v4().to_string();
let journal = Journal::open(data.path(), &replacement.id).await.unwrap();
journal
.save(&Record::new(replacement.clone(), original.seller_onion).unwrap())
.unwrap();
drop(journal);
assert_eq!(
Journal::find_for(
data.path(),
&binding.buyer_did,
&binding.seller_did,
&binding.content_id
)
.unwrap()
.unwrap()
.binding
.id,
replacement.id
);
}
#[tokio::test]
async fn owner_address_is_redacted_until_exposure_is_durable_and_cannot_then_be_retired() {
let data = tempfile::tempdir().unwrap();
let binding = binding();
let journal = Journal::open(data.path(), &binding.id).await.unwrap();
journal
.save(&Record::new(binding.clone(), format!("{}.onion", "a".repeat(56))).unwrap())
.unwrap();
let mut bytes = vec![0, 20];
bytes.extend([1; 20]);
let address = bitcoin::Address::from_script(
&bitcoin::ScriptBuf::from_bytes(bytes),
bitcoin::Network::Regtest,
)
.unwrap()
.to_string();
let saved = engine::accept_quote(
&journal,
engine::Quote {
binding: binding.clone(),
address: address.clone(),
network: engine::ChainNetwork::Regtest,
source: crate::content_lightning::RetainedFile {
sha256: "a".repeat(64),
size: 4,
filename: "original.txt".into(),
mime_type: "text/plain".into(),
},
},
)
.unwrap();
assert!(public(&saved).unwrap()["address"].is_null());
let ack = crate::content_onchain_seller::UnallocatedAck {
binding: binding.clone(),
state: "cancelled_unallocated".into(),
address: serde_json::Value::Null,
allocation_dispatched: false,
can_switch_method: true,
};
assert!(engine::retire_unallocated(&journal, ack.clone()).is_err());
assert_eq!(engine::expose_address(&journal).unwrap(), address);
drop(journal);
let journal = Journal::open(data.path(), &binding.id).await.unwrap();
let exposed = journal.load().unwrap().unwrap();
assert!(exposed.externally_exposed);
assert_eq!(public(&exposed).unwrap()["address"], address);
assert!(engine::retire_unallocated(&journal, ack).is_err());
assert!(exposed.blocks_other_rails());
}
}
+2 -2
View File
@@ -19,7 +19,7 @@ pub(crate) struct Binding {
pub price_sats: u64,
}
impl Binding {
fn validate(&self) -> Result<()> {
pub(crate) fn validate(&self) -> Result<()> {
anyhow::ensure!(
uuid::Uuid::parse_str(&self.id)?.to_string() == self.id,
"Invalid invoice operation"
@@ -61,7 +61,7 @@ pub(crate) struct RetainedFile {
pub mime_type: String,
}
impl RetainedFile {
fn validate(&self) -> Result<()> {
pub(crate) fn validate(&self) -> Result<()> {
anyhow::ensure!(
self.sha256.len() == 64
&& self.sha256.bytes().all(|b| b.is_ascii_hexdigit())
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,250 @@
//! Non-signable funding intent. No recipient address or executable PSBT exists
//! until explicit Pay has leased these exact inputs and recovered seller allocation.
use crate::{
content_lightning::{Binding, RetainedFile},
content_onchain::{ChainNetwork, FeePolicy, Funded, Lease, Quote, Record},
};
use anyhow::{Context, Result};
use base64::Engine;
use bitcoin::{
absolute::LockTime, consensus, psbt::Psbt, transaction::Version, Amount, ScriptBuf, Sequence,
Transaction, TxIn, TxOut, Witness,
};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
const MAX_SATS: u64 = 2_100_000_000_000_000;
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct Offer {
pub binding: Binding,
pub source: RetainedFile,
pub network: ChainNetwork,
/// This version accepts only a native P2WPKH recipient (22 script bytes).
pub recipient_script_type: String,
}
impl Offer {
pub fn validate(&self) -> Result<()> {
self.binding.validate()?;
self.source.validate()?;
anyhow::ensure!(
(546..=MAX_SATS).contains(&self.binding.price_sats)
&& self.recipient_script_type == "p2wpkh",
"Unsupported original on-chain offer"
);
Ok(())
}
pub fn hash(&self) -> Result<String> {
self.validate()?;
Ok(hex::encode(Sha256::digest(serde_json::to_vec(self)?)))
}
pub fn check_quote(&self, quote: &Quote) -> Result<()> {
self.validate()?;
anyhow::ensure!(
quote.binding == self.binding
&& quote.source == self.source
&& quote.network == self.network
&& quote.script()?.is_p2wpkh(),
"Allocated address changed the original offer"
);
Ok(())
}
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct PlanInput {
pub lease: Lease,
pub previous_tx_hex: String,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct FundingPlan {
pub offer_sha256: String,
pub inputs: Vec<PlanInput>,
pub change_address: String,
pub change_script: String,
pub change_sats: u64,
pub fee_sats: u64,
pub fee_rate_sat_vbyte: u64,
pub max_fee_sats: u64,
}
impl FundingPlan {
pub fn hash(&self) -> Result<String> {
Ok(hex::encode(Sha256::digest(serde_json::to_vec(self)?)))
}
pub fn validate(&self, record: &Record) -> Result<()> {
let offer = record.offer.as_ref().context("Original offer missing")?;
offer.validate()?;
anyhow::ensure!(
self.offer_sha256 == offer.hash()? && offer.binding == record.binding,
"Funding plan belongs to a different offer"
);
let change = self
.change_address
.parse::<bitcoin::Address<bitcoin::address::NetworkUnchecked>>()?
.require_network(offer.network.bitcoin())?
.script_pubkey();
anyhow::ensure!(
hex::encode(change.as_bytes()) == self.change_script
&& (change.is_p2wpkh() || change.is_p2tr()),
"Invalid original change script"
);
anyhow::ensure!(
matches!(&record.change_address,Some(crate::content_onchain::ChangeAddress::Ready{address}) if address==&self.change_address),
"Funding plan change allocation changed"
);
anyhow::ensure!(
!self.inputs.is_empty()
&& self.inputs.len() <= 32
&& self.max_fee_sats > 0
&& self.max_fee_sats <= MAX_SATS
&& (1..=5000).contains(&self.fee_rate_sat_vbyte),
"Invalid funding plan limits"
);
anyhow::ensure!(
self.change_sats == 0 || self.change_sats >= 546,
"Dust change is unsupported"
);
let mut seen = std::collections::HashSet::new();
let mut total = 0u64;
for input in &self.inputs {
input.lease.validate(&record.lock_id)?;
anyhow::ensure!(
input.lease.expires_at == 0 && seen.insert(input.lease.outpoint()?),
"Invalid or duplicate planned input"
);
anyhow::ensure!(
input.previous_tx_hex.len() <= 2 * 1024 * 1024,
"Previous transaction too large"
);
let previous: Transaction =
consensus::deserialize(&hex::decode(&input.previous_tx_hex)?)?;
anyhow::ensure!(
previous.compute_txid().to_string() == input.lease.txid,
"Original input transaction changed"
);
let output = previous
.output
.get(input.lease.vout as usize)
.context("Original input index missing")?;
anyhow::ensure!(
output.value.to_sat() == input.lease.value_sats
&& hex::encode(output.script_pubkey.as_bytes()) == input.lease.script,
"Original input metadata changed"
);
total = total
.checked_add(input.lease.value_sats)
.filter(|v| *v <= MAX_SATS)
.context("Input sum overflow")?;
}
let debit = offer
.binding
.price_sats
.checked_add(self.change_sats)
.and_then(|v| v.checked_add(self.fee_sats))
.context("Payment amount overflow")?;
anyhow::ensure!(
total == debit && self.fee_sats > 0 && self.fee_sats <= self.max_fee_sats,
"Funding plan fee or outputs changed"
);
let leases: Vec<_> = self.inputs.iter().map(|i| i.lease.clone()).collect();
let minimum = self
.fee_rate_sat_vbyte
.checked_mul(max_vsize(
&leases,
if self.change_sats == 0 {
None
} else {
Some(&change)
},
)?)
.context("Fee estimate overflow")?;
anyhow::ensure!(
self.fee_sats >= minimum,
"Funding plan fee does not cover reviewed rate"
);
Ok(())
}
/// Only now, with the real original seller address, construct a PSBT.
pub fn bind(&self, record: &Record, quote: &Quote) -> Result<(FeePolicy, Funded)> {
self.validate(record)?;
record.offer.as_ref().unwrap().check_quote(quote)?;
let recipient = quote.script()?;
let change = ScriptBuf::from_bytes(hex::decode(&self.change_script)?);
anyhow::ensure!(recipient != change, "Recipient cannot equal local change");
let mut outputs = vec![TxOut {
value: Amount::from_sat(record.binding.price_sats),
script_pubkey: recipient,
}];
if self.change_sats > 0 {
outputs.push(TxOut {
value: Amount::from_sat(self.change_sats),
script_pubkey: change,
});
}
let tx = Transaction {
version: Version::TWO,
lock_time: LockTime::ZERO,
input: self
.inputs
.iter()
.map(|i| {
Ok(TxIn {
previous_output: i.lease.outpoint()?,
script_sig: ScriptBuf::new(),
sequence: Sequence::ENABLE_RBF_NO_LOCKTIME,
witness: Witness::new(),
})
})
.collect::<Result<Vec<_>>>()?,
output: outputs,
};
let max_rate = self.fee_sats.div_ceil(tx.vsize() as u64);
let mut psbt = Psbt::from_unsigned_tx(tx)?;
for (metadata, input) in psbt.inputs.iter_mut().zip(&self.inputs) {
metadata.witness_utxo = Some(TxOut {
value: Amount::from_sat(input.lease.value_sats),
script_pubkey: ScriptBuf::from_bytes(hex::decode(&input.lease.script)?),
});
metadata.non_witness_utxo = Some(consensus::deserialize(&hex::decode(
&input.previous_tx_hex,
)?)?);
}
Ok((
FeePolicy {
change_script: self.change_script.clone(),
max_fee_sats: self.max_fee_sats,
max_fee_rate_sat_vbyte: max_rate,
},
Funded {
psbt_base64: base64::engine::general_purpose::STANDARD.encode(psbt.serialize()),
leases: self.inputs.iter().map(|i| i.lease.clone()).collect(),
},
))
}
}
/// Weight calculation only: never constructs a placeholder-address transaction.
/// Versions, sequence and locktime are fixed by this protocol; <=32 inputs/2 outputs
/// mean single-byte CompactSize counts. Native inputs have empty scriptSig.
pub(crate) fn max_vsize(inputs: &[Lease], change: Option<&ScriptBuf>) -> Result<u64> {
anyhow::ensure!(
!inputs.is_empty() && inputs.len() <= 32,
"Unsupported input count"
);
let mut stripped = 4 + 1 + 41 * (inputs.len() as u64) + 1 + 8 + 1 + 22 + 4;
if let Some(script) = change {
anyhow::ensure!(script.len() < 253, "Unsupported change script length");
stripped += 8 + 1 + script.len() as u64;
}
let mut witness = 2u64;
for input in inputs {
let script = ScriptBuf::from_bytes(hex::decode(&input.script)?);
witness += if script.is_p2wpkh() {
109
} else if script.is_p2tr() {
67
} else {
anyhow::bail!("Unsupported signing input")
};
}
Ok((stripped * 4 + witness).div_ceil(4))
}
@@ -0,0 +1,551 @@
//! Buyer-bound seller address allocation. Unknown allocation never creates a replacement.
use crate::{
content_lightning::{Binding, RetainedFile},
content_onchain::{ChainNetwork, Quote},
};
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(tag = "state", rename_all = "snake_case", deny_unknown_fields)]
pub(crate) enum Allocation {
Prepared,
Dispatched,
Ready { address: String },
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct Record {
pub binding: Binding,
pub source: RetainedFile,
pub network: ChainNetwork,
pub allocation: Allocation,
pub paid: bool,
}
impl Record {
fn validate(&self) -> Result<()> {
self.binding.validate()?;
self.source.validate()?;
anyhow::ensure!(
(546..=2_100_000_000_000_000).contains(&self.binding.price_sats),
"On-chain price below supported minimum"
);
if let Allocation::Ready { address } = &self.allocation {
self.quote()?
.context("Missing original address")?
.script()?;
anyhow::ensure!(!address.is_empty(), "Missing address");
}
anyhow::ensure!(
!self.paid || matches!(self.allocation, Allocation::Ready { .. }),
"Paid operation lacks address"
);
Ok(())
}
pub fn offer(&self) -> Result<crate::content_onchain_plan::Offer> {
let offer = crate::content_onchain_plan::Offer {
binding: self.binding.clone(),
source: self.source.clone(),
network: self.network,
recipient_script_type: "p2wpkh".into(),
};
offer.validate()?;
Ok(offer)
}
pub fn quote(&self) -> Result<Option<Quote>> {
Ok(match &self.allocation {
Allocation::Ready { address } => Some(Quote {
binding: self.binding.clone(),
address: address.clone(),
network: self.network,
source: self.source.clone(),
}),
_ => None,
})
}
}
/// Terminal proof for an operation whose receive allocation was never dispatched.
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct UnallocatedAck {
pub binding: Binding,
pub state: String,
pub address: serde_json::Value,
pub allocation_dispatched: bool,
pub can_switch_method: bool,
}
impl UnallocatedAck {
pub fn validate(&self, binding: &Binding) -> Result<()> {
binding.validate()?;
anyhow::ensure!(
&self.binding == binding
&& self.state == "cancelled_unallocated"
&& self.address.is_null()
&& !self.allocation_dispatched
&& self.can_switch_method,
"Invalid unallocated retirement acknowledgement"
);
Ok(())
}
}
#[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<Self> {
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-onchain-seller");
fs::create_dir_all(&directory)?;
anyhow::ensure!(
fs::symlink_metadata(&directory)?.is_dir(),
"On-chain seller 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 on-chain seller 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<PathBuf> {
anyhow::ensure!(
matches!(role, "seller" | "retired") && uuid::Uuid::parse_str(id)?.to_string() == id,
"Invalid on-chain seller journal key"
);
Ok(self.directory.join(format!("{role}-{id}.json")))
}
fn read<T: serde::de::DeserializeOwned>(&self, role: &str, id: &str) -> Result<Option<T>> {
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 on-chain seller record");
let mut bytes = Vec::new();
file.take(65537).read_to_end(&mut bytes)?;
anyhow::ensure!(bytes.len() <= 65536, "On-chain seller record too large");
let envelope: Envelope = serde_json::from_slice(&bytes)
.context("On-chain seller recovery damaged; do not pay again")?;
anyhow::ensure!(
hex::encode(Sha256::digest(envelope.payload.as_bytes())) == envelope.checksum,
"On-chain seller recovery checksum changed"
);
Ok(Some(serde_json::from_str(&envelope.payload)?))
}
fn write<T: Serialize>(&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, "On-chain seller 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 retirement(&self, binding: &Binding) -> Result<Option<UnallocatedAck>> {
binding.validate()?;
let saved: Option<UnallocatedAck> = self.read("retired", &binding.id)?;
if let Some(ack) = &saved {
ack.validate(binding)?;
}
Ok(saved)
}
pub fn retire_unallocated(&self, binding: &Binding) -> Result<UnallocatedAck> {
if let Some(ack) = self.retirement(binding)? {
return Ok(ack);
}
if let Some(saved) = self.load(binding)? {
anyhow::ensure!(saved.allocation==Allocation::Prepared && !saved.paid,"Original address allocation was dispatched or paid; recover it without another payment");
}
// Even an absent operation gets a durable tombstone. A delayed create
// must observe retirement rather than allocating after the acknowledgement.
let ack = UnallocatedAck {
binding: binding.clone(),
state: "cancelled_unallocated".into(),
address: serde_json::Value::Null,
allocation_dispatched: false,
can_switch_method: true,
};
ack.validate(binding)?;
self.write("retired", &binding.id, &ack)?;
Ok(ack)
}
pub fn load(&self, binding: &Binding) -> Result<Option<Record>> {
binding.validate()?;
let saved: Option<Record> = self.read("seller", &binding.id)?;
if let Some(record) = &saved {
record.validate()?;
anyhow::ensure!(
&record.binding == binding,
"Original seller operation binding changed"
);
}
Ok(saved)
}
pub fn save(&self, record: &Record) -> Result<()> {
record.validate()?;
anyhow::ensure!(
self.retirement(&record.binding)?.is_none(),
"Original address operation is retired"
);
if let Some(old) = self.load(&record.binding)? {
anyhow::ensure!(
old.source == record.source && old.network == record.network,
"Original file/network changed"
);
anyhow::ensure!(
!old.paid || record.paid,
"Confirmed purchase cannot become unpaid"
);
match old.allocation {
Allocation::Ready { .. } => anyhow::ensure!(
old.allocation == record.allocation,
"Original receive address changed"
),
Allocation::Dispatched => anyhow::ensure!(
record.allocation != Allocation::Prepared,
"Ambiguous address allocation cannot restart"
),
Allocation::Prepared => {}
}
}
self.write("seller", &record.binding.id, record)
}
pub fn prepare(
&self,
binding: Binding,
source: RetainedFile,
network: ChainNetwork,
) -> Result<Record> {
anyhow::ensure!(
self.retirement(&binding)?.is_none(),
"Original address operation is retired"
);
if let Some(old) = self.load(&binding)? {
anyhow::ensure!(
old.source == source && old.network == network,
"Original sale changed"
);
return Ok(old);
}
let record = Record {
binding,
source,
network,
allocation: Allocation::Prepared,
paid: false,
};
self.save(&record)?;
Ok(record)
}
}
pub(crate) trait Wallet {
async fn network(&self) -> Result<ChainNetwork> {
anyhow::bail!("Seller wallet network unavailable")
}
async fn preflight(&self, network: ChainNetwork) -> Result<()>;
async fn allocate(&self) -> Result<String>;
async fn received(&self, address: &str, amount: u64) -> Result<bool>;
}
pub(crate) async fn drive<W: Wallet>(
journal: &Journal,
binding: &Binding,
create: bool,
wallet: &W,
) -> Result<Record> {
anyhow::ensure!(
journal.retirement(binding)?.is_none(),
"Original address operation is retired"
);
let mut saved = journal
.load(binding)?
.context("Original on-chain sale missing")?;
if saved.allocation == Allocation::Prepared && create {
wallet.preflight(saved.network).await?;
saved.allocation = Allocation::Dispatched;
journal.save(&saved)?;
let address = wallet.allocate().await?;
saved.allocation = Allocation::Ready { address };
journal.save(&saved)?;
}
if !saved.paid {
if let Allocation::Ready { address } = &saved.allocation {
if wallet.received(address, saved.binding.price_sats).await? {
saved.paid = true;
journal.save(&saved)?;
}
}
}
Ok(saved)
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::{
atomic::{AtomicBool, AtomicUsize, Ordering},
Mutex,
};
struct Mock {
calls: AtomicUsize,
lost: AtomicBool,
paid: AtomicBool,
preflight_fails: AtomicBool,
observed: Mutex<Vec<String>>,
}
impl Wallet for Mock {
async fn preflight(&self, _: ChainNetwork) -> Result<()> {
anyhow::ensure!(
!self.preflight_fails.load(Ordering::SeqCst),
"Wallet unavailable before allocation"
);
Ok(())
}
async fn allocate(&self) -> Result<String> {
self.calls.fetch_add(1, Ordering::SeqCst);
anyhow::ensure!(
!self.lost.swap(false, Ordering::SeqCst),
"Lost address allocation reply"
);
Ok(address())
}
async fn received(&self, address: &str, _: u64) -> Result<bool> {
self.observed.lock().unwrap().push(address.into());
Ok(self.paid.load(Ordering::SeqCst))
}
}
fn address() -> String {
let mut script = vec![0, 20];
script.extend([1; 20]);
bitcoin::Address::from_script(
&bitcoin::ScriptBuf::from_bytes(script),
bitcoin::Network::Regtest,
)
.unwrap()
.to_string()
}
fn mock() -> Mock {
Mock {
calls: AtomicUsize::new(0),
lost: AtomicBool::new(false),
paid: AtomicBool::new(false),
preflight_fails: AtomicBool::new(false),
observed: Mutex::new(vec![]),
}
}
async fn fixture() -> (tempfile::TempDir, Journal, Binding) {
let data = tempfile::tempdir().unwrap();
let journal = Journal::open(data.path()).await.unwrap();
let 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: 546,
};
journal
.prepare(
binding.clone(),
RetainedFile {
sha256: "a".repeat(64),
size: 4,
filename: "original.txt".into(),
mime_type: "text/plain".into(),
},
ChainNetwork::Regtest,
)
.unwrap();
(data, journal, binding)
}
#[tokio::test]
async fn lost_address_response_never_allocates_again_even_after_restart() {
let (data, journal, binding) = fixture().await;
let wallet = mock();
wallet.lost.store(true, Ordering::SeqCst);
assert!(drive(&journal, &binding, true, &wallet).await.is_err());
drop(journal);
let journal = Journal::open(data.path()).await.unwrap();
let recovered = drive(&journal, &binding, true, &wallet).await.unwrap();
assert_eq!(recovered.allocation, Allocation::Dispatched);
assert!(!recovered.paid);
assert_eq!(wallet.calls.load(Ordering::SeqCst), 1);
assert!(wallet.observed.lock().unwrap().is_empty());
}
#[tokio::test]
async fn preflight_retry_and_read_only_status_preserve_single_allocation() {
let (_data, journal, binding) = fixture().await;
let wallet = mock();
drive(&journal, &binding, false, &wallet).await.unwrap();
assert_eq!(wallet.calls.load(Ordering::SeqCst), 0);
wallet.preflight_fails.store(true, Ordering::SeqCst);
assert!(drive(&journal, &binding, true, &wallet).await.is_err());
assert_eq!(
journal.load(&binding).unwrap().unwrap().allocation,
Allocation::Prepared
);
wallet.preflight_fails.store(false, Ordering::SeqCst);
drive(&journal, &binding, true, &wallet).await.unwrap();
drive(&journal, &binding, true, &wallet).await.unwrap();
assert_eq!(wallet.calls.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn paid_source_survives_reload_and_stale_unpaid_callback_cannot_regress() {
let (data, journal, binding) = fixture().await;
let wallet = mock();
let unpaid = drive(&journal, &binding, true, &wallet).await.unwrap();
wallet.paid.store(true, Ordering::SeqCst);
let paid = drive(&journal, &binding, false, &wallet).await.unwrap();
assert!(paid.paid);
assert!(journal.save(&unpaid).is_err());
drop(journal);
let journal = Journal::open(data.path()).await.unwrap();
wallet.paid.store(false, Ordering::SeqCst);
assert_eq!(
drive(&journal, &binding, false, &wallet).await.unwrap(),
paid
);
assert_eq!(paid.source.filename, "original.txt");
assert_eq!(wallet.calls.load(Ordering::SeqCst), 1);
let mut other = binding.clone();
other.buyer_did = binding.seller_did.clone();
assert!(journal.load(&other).is_err());
}
#[tokio::test]
async fn cancel_before_allocation_survives_lost_ack_and_blocks_delayed_create() {
let (data, journal, binding) = fixture().await;
let wallet = mock();
wallet.preflight_fails.store(true, Ordering::SeqCst);
assert!(drive(&journal, &binding, true, &wallet).await.is_err());
let original = journal.load(&binding).unwrap().unwrap();
let ack = journal.retire_unallocated(&binding).unwrap();
ack.validate(&binding).unwrap();
assert!(ack.address.is_null());
assert!(!ack.allocation_dispatched);
drop(journal);
let journal = Journal::open(data.path()).await.unwrap();
assert_eq!(journal.retire_unallocated(&binding).unwrap(), ack);
wallet.preflight_fails.store(false, Ordering::SeqCst);
assert!(drive(&journal, &binding, true, &wallet).await.is_err());
assert!(journal
.prepare(binding.clone(), original.source, original.network)
.is_err());
assert_eq!(wallet.calls.load(Ordering::SeqCst), 0);
}
#[tokio::test]
async fn absent_operation_retirement_is_a_tombstone_not_absence_evidence() {
let (_data, journal, binding) = fixture().await;
let mut unknown = binding.clone();
unknown.id = uuid::Uuid::new_v4().to_string();
let ack = journal.retire_unallocated(&unknown).unwrap();
assert!(journal.load(&unknown).unwrap().is_none());
let original = journal.load(&binding).unwrap().unwrap();
assert!(journal
.prepare(unknown.clone(), original.source, original.network)
.is_err());
assert_eq!(journal.retirement(&unknown).unwrap(), Some(ack));
}
#[tokio::test]
async fn dispatched_unknown_and_issued_addresses_cannot_be_retired() {
let (_data, journal, binding) = fixture().await;
let wallet = mock();
wallet.lost.store(true, Ordering::SeqCst);
assert!(drive(&journal, &binding, true, &wallet).await.is_err());
assert!(journal.retire_unallocated(&binding).is_err());
assert!(journal.retirement(&binding).unwrap().is_none());
drop(journal);
let (_data, journal, binding) = fixture().await;
let wallet = mock();
drive(&journal, &binding, true, &wallet).await.unwrap();
assert!(journal.retire_unallocated(&binding).is_err());
assert_eq!(wallet.calls.load(Ordering::SeqCst), 1);
}
#[test]
fn retirement_ack_requires_explicit_null_address_and_all_terminal_fields() {
let 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: 546,
};
let complete = serde_json::json!({"binding":binding,"state":"cancelled_unallocated","address":null,"allocation_dispatched":false,"can_switch_method":true});
for key in [
"address",
"state",
"allocation_dispatched",
"can_switch_method",
] {
let mut partial = complete.clone();
partial.as_object_mut().unwrap().remove(key);
assert!(
serde_json::from_value::<UnallocatedAck>(partial).is_err(),
"missing {key}"
);
}
serde_json::from_value::<UnallocatedAck>(complete)
.unwrap()
.validate(&binding)
.unwrap();
}
}
+85
View File
@@ -1961,3 +1961,88 @@ pub(crate) async fn publish_snapshot_invoice(
anyhow::ensure!(visible, "Item is not shared with this invoice buyer");
journal.prepare_seller_source(binding, Some(retained))
}
pub(crate) async fn publish_snapshot_onchain(
data_dir: &Path,
original: &ContentItem,
journal: &crate::content_onchain_seller::Journal,
binding: crate::content_lightning::Binding,
retained: crate::content_lightning::RetainedFile,
network: crate::content_onchain::ChainNetwork,
) -> Result<crate::content_onchain_seller::Record> {
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, "onchain"),
"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(binding, retained, network)
}
/// Serialize the first allocation with catalog changes. Previously allocated
/// operations recover their original terms even if sharing changes afterwards.
pub(crate) async fn allocate_onchain_offer<W: crate::content_onchain_seller::Wallet>(
data_dir: &Path,
journal: &crate::content_onchain_seller::Journal,
binding: &crate::content_lightning::Binding,
wallet: &W,
) -> Result<crate::content_onchain_seller::Record> {
let saved = journal.load(binding)?.context("Original offer missing")?;
if saved.allocation != crate::content_onchain_seller::Allocation::Prepared {
return crate::content_onchain_seller::drive(journal, binding, true, wallet).await;
}
let source_data = data_dir.to_path_buf();
let source_id = binding.content_id.clone();
let source = saved.source.clone();
tokio::task::spawn_blocking(move || {
crate::content_snapshot::open_matching(
&source_data,
&source_id,
&source.sha256,
source.size,
)
})
.await??;
let _held = CATALOG_WRITES.lock().await;
let catalog = load_catalog(data_dir).await?;
let item = catalog
.items
.iter()
.find(|item| item.id == binding.content_id)
.context("Original offer is no longer shared; cancel before allocation")?;
let visible = match &item.availability {
Availability::Nobody => false,
Availability::AllPeers => true,
Availability::Specific { peers } => peers.contains(&binding.buyer_did),
};
anyhow::ensure!(
visible
&& item.filename == saved.source.filename
&& item.mime_type == saved.source.mime_type
&& item.size_bytes == saved.source.size
&& matches!(&item.access,AccessControl::Paid {price_sats,..} if *price_sats==binding.price_sats)
&& method_accepted(&item.access, "onchain"),
"Original offer terms changed; no seller address allocated"
);
crate::content_onchain_seller::drive(journal, binding, true, wallet).await
}
+3
View File
@@ -45,6 +45,9 @@ mod content_hash;
mod content_indeehub;
mod content_invoice;
mod content_lightning;
mod content_onchain;
mod content_onchain_plan;
mod content_onchain_seller;
mod content_payment_admission;
mod content_owned;
mod content_purchase;