diff --git a/AGENTS.md b/AGENTS.md index 6b75350c..46fbf04c 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -21,3 +21,11 @@ While its status is OPEN: This priority comes from the user's explicit instruction on 2026-09-15. It remains in effect across sessions until the documented acceptance criteria are met or the user explicitly changes it. + +## Unit tests on a live node + +Run backend unit tests through `scripts/test-backend-isolated.sh`. Do not run +unrestricted `cargo test` on a node with installed apps: older mocked-runtime +tests still reached real service commands. The runner isolates wallet data, +service buses, container storage, networking, and process IDs. Compilation with +`cargo test --no-run` is safe. Keep separately authorized live checks explicit. diff --git a/CHANGELOG.md b/CHANGELOG.md index e8c5d070..72057cc3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,24 @@ ## Unreleased +## v1.8.21-alpha (2026-09-30) + +- Fixed Bitcoin and other containers being forcibly stopped after ten seconds during managed updates and restarts. +- Existing installations now receive the same graceful shutdown allowance as new containers, without restarting apps just to apply this setting. +- Prevented unnecessary Lightning restarts when Bitcoin has stayed running; dependency restarts now require an observed Bitcoin container change. +- Includes the Cashu payment, optional Bitcoin pruning, Lightning readiness, and explorer improvements from 1.8.20. + +## v1.8.20-alpha (2026-09-29) + +- Fixed Cashu file payments rejected despite a shared mint, and preserved the payment amount when mint fees reduce change. +- Payment failures now report whether a refund actually succeeded; missing files and unsupported payment methods are rejected before charging. +- Improved saving paid files into Files and reopening purchases without paying again. +- Bitcoin Core and Knots installation offers optional pruning on larger disks, using the same settings as automatic pruning. +- Fixed false missing-port checks that unnecessarily restarted Bitcoin and LND; recovery now respects managed shutdown timeouts. +- LND explains when it is waiting for Bitcoin installation, startup, or sync, without treating normal synchronization as a restart-worthy failure. +- Bitcoin startup messages explain block-index loading without exposing raw RPC errors, and Lightning keeps known balances clearly marked during outages. +- Changed the public transaction-explorer default to mempool.space while preserving local explorers and custom choices. + ## v1.8.19-alpha (2026-09-28) - Fixed the embedded AIUI chat page painting a second background and dark scrim over Archy’s dashboard background. diff --git a/apps/bitcoin-core/manifest.yml b/apps/bitcoin-core/manifest.yml index a14a50ea..6a13168c 100644 --- a/apps/bitcoin-core/manifest.yml +++ b/apps/bitcoin-core/manifest.yml @@ -54,7 +54,7 @@ app: if [ -n "$RPC_TXRELAY_AUTH" ]; then RPC_TXRELAY_FLAGS="$RPC_TXRELAY_FLAGS -rpcauth=$RPC_TXRELAY_AUTH -rpcwhitelist=txrelay:sendrawtransaction,submitpackage,testmempoolaccept,getmempoolinfo,getrawmempool,getmempoolentry,getnetworkinfo,getblockchaininfo,getblockcount,getblockhash,getblock,getblockheader,getrawtransaction,gettxout,gettxspendingprevout,decoderawtransaction,decodescript,estimatesmartfee,uptime,ping,getconnectioncount,getpeerinfo,getindexinfo,getdeploymentinfo,getchaintips"; fi; - if [ "${DISK_GB_VALUE:-0}" -lt 1000 ]; then + if [ "${BITCOIN_PRUNE:-0}" = "1" ] || [ "${DISK_GB_VALUE:-0}" -lt 1000 ]; then exec "$BITCOIND" -datadir=/home/bitcoin/.bitcoin -conf="$RPC_CONF" -allowignoredconf=1 -printtoconsole=0 -server=1 -prune=50000 -rpcallowip=0.0.0.0/0 -rpcbind=0.0.0.0:8332 -listen=1 -bind=0.0.0.0:8333 -dbcache=1024 -par=0 -maxconnections=125 $RPC_HEADROOM $RPC_TXRELAY_FLAGS; else exec "$BITCOIND" -datadir=/home/bitcoin/.bitcoin -conf="$RPC_CONF" -allowignoredconf=1 -printtoconsole=0 -server=1 -txindex=1 -rpcallowip=0.0.0.0/0 -rpcbind=0.0.0.0:8332 -listen=1 -bind=0.0.0.0:8333 -dbcache=4096 -par=0 -maxconnections=125 $RPC_HEADROOM $RPC_TXRELAY_FLAGS; diff --git a/apps/bitcoin-knots/manifest.yml b/apps/bitcoin-knots/manifest.yml index 5545fad1..4dd080af 100644 --- a/apps/bitcoin-knots/manifest.yml +++ b/apps/bitcoin-knots/manifest.yml @@ -60,7 +60,7 @@ app: if [ -n "$RPC_TXRELAY_AUTH" ]; then RPC_TXRELAY_FLAGS="$RPC_TXRELAY_FLAGS -rpcauth=$RPC_TXRELAY_AUTH -rpcwhitelist=txrelay:sendrawtransaction,submitpackage,testmempoolaccept,getmempoolinfo,getrawmempool,getmempoolentry,getnetworkinfo,getblockchaininfo,getblockcount,getblockhash,getblock,getblockheader,getrawtransaction,gettxout,gettxspendingprevout,decoderawtransaction,decodescript,estimatesmartfee,uptime,ping,getconnectioncount,getpeerinfo,getindexinfo,getdeploymentinfo,getchaintips"; fi; - if [ "${DISK_GB_VALUE:-0}" -lt 1000 ]; then + if [ "${BITCOIN_PRUNE:-0}" = "1" ] || [ "${DISK_GB_VALUE:-0}" -lt 1000 ]; then exec "$BITCOIND" -datadir=/home/bitcoin/.bitcoin -conf="$RPC_CONF" -allowignoredconf=1 -printtoconsole=0 -server=1 -prune=50000 -rpcallowip=0.0.0.0/0 -rpcbind=0.0.0.0:8332 -listen=1 -bind=0.0.0.0:8333 -dbcache=2048 -par=0 -maxconnections=125 $RPC_HEADROOM $RPC_TXRELAY_FLAGS; else exec "$BITCOIND" -datadir=/home/bitcoin/.bitcoin -conf="$RPC_CONF" -allowignoredconf=1 -printtoconsole=0 -server=1 -txindex=1 -rpcallowip=0.0.0.0/0 -rpcbind=0.0.0.0:8332 -listen=1 -bind=0.0.0.0:8333 -dbcache=4096 -par=0 -maxconnections=125 $RPC_HEADROOM $RPC_TXRELAY_FLAGS; diff --git a/core/Cargo.lock b/core/Cargo.lock index f44ec9e0..168c15d5 100644 --- a/core/Cargo.lock +++ b/core/Cargo.lock @@ -104,7 +104,7 @@ dependencies = [ [[package]] name = "archipelago" -version = "1.8.19-alpha" +version = "1.8.21-alpha" dependencies = [ "anyhow", "archipelago-container", diff --git a/core/archipelago/Cargo.toml b/core/archipelago/Cargo.toml index d920c578..87199f6a 100644 --- a/core/archipelago/Cargo.toml +++ b/core/archipelago/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "archipelago" -version = "1.8.19-alpha" +version = "1.8.21-alpha" edition = "2021" license.workspace = true description = "Archipelago Bitcoin Node OS - Native backend" diff --git a/core/archipelago/src/api/handler/content.rs b/core/archipelago/src/api/handler/content.rs index 0fd815be..24100269 100644 --- a/core/archipelago/src/api/handler/content.rs +++ b/core/archipelago/src/api/handler/content.rs @@ -166,9 +166,14 @@ impl ApiHandler { StatusCode::SERVICE_UNAVAILABLE, "application/json", hyper::Body::from( - r#"{"error":"The seller's node can't read this file right now. No payment was taken."}"#, + r#"{"error":"The seller's node can't read this file right now. This request did not redeem an ecash payment."}"#, ), )), + Ok(content_server::ServeResult::RangeNotSatisfiable(total)) => Ok(Response::builder() + .status(StatusCode::RANGE_NOT_SATISFIABLE) + .header("Content-Range", format!("bytes */{total}")) + .body(hyper::Body::empty()) + .unwrap()), Ok(content_server::ServeResult::NotFound) => Ok(build_response( StatusCode::NOT_FOUND, "text/plain", diff --git a/core/archipelago/src/api/handler/proxy.rs b/core/archipelago/src/api/handler/proxy.rs index 7ddceb4e..cda46c47 100644 --- a/core/archipelago/src/api/handler/proxy.rs +++ b/core/archipelago/src/api/handler/proxy.rs @@ -138,6 +138,19 @@ impl ApiHandler { cors_origin: &str, ) -> Result> { let suffix = path.strip_prefix("/proxy/lnd").unwrap_or("/"); + if suffix == "/archy-status" { + return Ok(Response::builder() + .status(StatusCode::OK) + .header("Content-Type", "application/json") + .header("Cache-Control", "no-store") + .header("Access-Control-Allow-Origin", cors_origin) + .header("Access-Control-Allow-Credentials", "true") + .header("Vary", "Origin") + .body(hyper::Body::from( + rpc.handle_lnd_readiness().await.to_string(), + ))?); + } + let url = format!("{LND_REST_BASE_URL}{suffix}"); // LND REST serves a self-signed cert and requires the admin macaroon. // A bare reqwest::get() uses the default client, which rejects the diff --git a/core/archipelago/src/api/rpc/content.rs b/core/archipelago/src/api/rpc/content.rs index f7e46d2c..4ba9c107 100644 --- a/core/archipelago/src/api/rpc/content.rs +++ b/core/archipelago/src/api/rpc/content.rs @@ -22,9 +22,9 @@ const FILE_CATALOG_PROTOCOL: &str = "https://archipelago.dev/protocols/file-cata /// Best-effort reclaim of an ecash payment token that was minted but the sale /// didn't complete (seller unreachable or couldn't redeem it), so the buyer /// doesn't lose the value. For Fedimint the spender can reissue its own -/// un-redeemed notes; for Cashu the proofs are received back. Fails silently if -/// the seller already claimed the token (then the value is genuinely gone). -async fn reclaim_spent_ecash(data_dir: &std::path::Path, token: &str, backend: &str) { +/// un-redeemed notes; for Cashu the proofs are received back. Report the actual +/// recovered amount, or explicitly say when a refund could not be confirmed. +async fn reclaim_spent_ecash(data_dir: &std::path::Path, token: &str, backend: &str) -> String { let res = match backend { "fedimint" => crate::wallet::fedimint_client::reissue_into_any(data_dir, token) .await @@ -32,16 +32,120 @@ async fn reclaim_spent_ecash(data_dir: &std::path::Path, token: &str, backend: & _ => ecash::receive_token(data_dir, token).await, }; match res { - Ok(sats) => tracing::info!( - "paid download: reclaimed {sats} sats of unspent {backend} ecash after a failed sale" - ), - Err(e) => tracing::warn!( - "paid download: could not reclaim {backend} ecash (the peer may have already \ - claimed it): {e:#}" - ), + Ok(sats) => { + tracing::info!("paid download: reclaimed {sats} sats after failed sale"); + format!("Refunded {sats} sats to your wallet.") + } + Err(e) => { + tracing::warn!("paid download: refund not confirmed: {e}"); + "Your refund could not be confirmed. The seller may have received the payment. Do not pay again until this is checked.".to_string() + } } } +/// Only pass through the peer's bounded, printable explanation; refund status +/// is always determined locally and must never come from the peer's wording. +fn seller_error_message(status: reqwest::StatusCode, body: &str) -> String { + let reason = serde_json::from_str::(body) + .ok() + .and_then(|v| v.get("error").and_then(|e| e.as_str()).map(str::to_owned)); + match reason { + Some(reason) if !reason.trim().is_empty() => { + let clean: String = reason + .chars() + .filter(|c| !c.is_control()) + .take(240) + .collect(); + format!("Seller response ({status}): {clean}") + } + _ => format!("Peer returned an error ({status})."), + } +} + +/// Keep first purchases and cached repeats compatible with both existing clients. +fn paid_content_response(bytes: &[u8], mime: &str, paid_sats: u64) -> serde_json::Value { + use base64::Engine; + let data = base64::engine::general_purpose::STANDARD.encode(bytes); + serde_json::json!({ + "data": data, "data_base64": data, + "size": bytes.len(), "size_bytes": bytes.len(), + "mime_type": mime, "paid_sats": paid_sats, "owned": true, + }) +} + +/// FileBrowser owns its files through a rootless UID mapping. Use its authenticated +/// API rather than writing host paths with the backend's unrelated UID. Its +/// override=false upload atomically refuses existing names, including races. +async fn file_purchase_in_files( + client: &reqwest::Client, + base_url: &str, + token: &str, + filename: &str, + mime: &str, + bytes: &[u8], +) -> Result { + let folder = if mime.starts_with("image/") || mime.starts_with("video/") { + "Photos" + } else if mime.starts_with("audio/") { + "Music" + } else { + "Documents" + }; + let mut folder_url = reqwest::Url::parse(base_url)?; + folder_url + .path_segments_mut() + .map_err(|_| anyhow::anyhow!("Invalid Files URL"))? + .extend(["api", "resources", folder, ""]); + let response = client + .get(folder_url.clone()) + .header("X-Auth", token) + .send() + .await?; + if response.status() == reqwest::StatusCode::NOT_FOUND { + let response = client + .post(folder_url.clone()) + .header("X-Auth", token) + .send() + .await?; + if response.status() != reqwest::StatusCode::CONFLICT { + response.error_for_status()?; + } + } else { + response.error_for_status()?; + } + let base = std::path::Path::new(filename) + .file_name() + .and_then(|n| n.to_str()) + .filter(|n| !n.is_empty()) + .unwrap_or("download"); + let (stem, extension) = match base.rsplit_once('.') { + Some((stem, ext)) if !stem.is_empty() => (stem, format!(".{ext}")), + _ => (base, String::new()), + }; + for attempt in 1..=100 { + let name = if attempt == 1 { + base.to_string() + } else { + format!("{stem} ({attempt}){extension}") + }; + let mut url = folder_url.clone(); + url.path_segments_mut().unwrap().pop_if_empty().push(&name); + url.query_pairs_mut().append_pair("override", "false"); + let response = client + .post(url) + .header("X-Auth", token) + .body(bytes.to_vec()) + .send() + .await?; + if response.status() == reqwest::StatusCode::CONFLICT { + continue; + } + response.error_for_status()?; + return Ok(format!("{folder}/{name}")); + } + anyhow::bail!("Too many existing copies; purchased file remains in the purchase cache") +} + impl RpcHandler { /// List content I'm sharing. pub(super) async fn handle_content_list_mine(&self) -> Result { @@ -463,17 +567,10 @@ impl RpcHandler { crate::content_owned::read_owned(&self.config.data_dir, &o.onion, &o.content_id) .await { - use base64::Engine; - return Ok(serde_json::json!({ - "owned": true, - "already_owned": true, - "filename": o.filename, - "mime_type": mime, - "size_bytes": bytes.len(), - "paid_sats": 0, - "data_base64": - base64::engine::general_purpose::STANDARD.encode(&bytes), - })); + let mut result = paid_content_response(&bytes, &mime, 0); + result["already_owned"] = serde_json::json!(true); + result["filename"] = serde_json::json!(o.filename); + return Ok(result); } // Cache record exists but bytes are gone — fall through and // repurchase rather than stranding the user. @@ -545,34 +642,30 @@ impl RpcHandler { let path = format!("/content/{}", content_id); // Surface a real reason instead of the generic sanitized error (#30): - // the dial already tries FIPS/mesh then falls back to Tor, so a failure - // here means the peer is genuinely unreachable on both transports. - let (response, transport) = match crate::fips::dial::PeerRequest::new( - fips_npub.as_deref(), - onion, - &path, - ) - .service(crate::settings::transport::PeerService::PeerFiles) - .header("X-Federation-DID", local_did) - .header("X-Payment-Token", token_str.clone()) - // The token is a bearer instrument the seller redeems on first sight: - // a Tor replay after FIPS delivered it can only arrive spent. - .single_delivery() - .timeout(std::time::Duration::from_secs(900)) - .send_get() - .await - { - Ok(v) => v, - Err(e) => { - tracing::warn!("paid peer download dial failed for {}: {:#}", onion, e); - // The token was already minted/spent — reclaim it so the buyer - // doesn't lose the value when the seller was simply unreachable. - reclaim_spent_ecash(&self.config.data_dir, &token_str, used_backend).await; - return Ok(serde_json::json!({ - "error": "Could not reach the peer over mesh or Tor — it may be offline. Your ecash was refunded to your wallet. Please try again." - })); - } - }; + // A bearer token must not be replayed after an ambiguous delivery. + // A transport error can mean the seller received it without replying. + let (response, transport) = + match crate::fips::dial::PeerRequest::new(fips_npub.as_deref(), onion, &path) + .service(crate::settings::transport::PeerService::PeerFiles) + .header("X-Federation-DID", local_did) + .header("X-Payment-Token", token_str.clone()) + .single_delivery() + .timeout(std::time::Duration::from_secs(900)) + .send_get() + .await + { + Ok(v) => v, + Err(e) => { + tracing::warn!("paid peer download dial failed for {}: {:#}", onion, e); + // The token was already minted/spent — reclaim it so the buyer + // doesn't lose the value when the seller was simply unreachable. + let refund = + reclaim_spent_ecash(&self.config.data_dir, &token_str, used_backend).await; + return Ok(serde_json::json!({ + "error": format!("The purchase could not be completed. {refund}") + })); + } + }; // Record which transport actually reached the peer (B14). if let Err(e) = crate::federation::record_peer_transport( &self.config.data_dir, @@ -586,25 +679,17 @@ impl RpcHandler { } if response.status() == reqwest::StatusCode::PAYMENT_REQUIRED { - // Payment was rejected by the seller. Surface the most likely cause - // per backend — for ecash both sides must share a redemption network - // (a Cashu mint, or a Fedimint federation). + // A 402 can mean mint validation, network failure, underpayment, + // or an unaccepted mint. Do not invent a mint-mismatch diagnosis. let body = response.text().await.unwrap_or_default(); tracing::warn!( "paid download: seller {onion} rejected {used_backend} payment of {price_sats} sats: {body}" ); // Seller couldn't redeem the token — reclaim it so the buyer keeps // their funds (the spent-but-unredeemed-notes case the user hit). - reclaim_spent_ecash(&self.config.data_dir, &token_str, used_backend).await; - let hint = match used_backend { - "fedimint" => "the seller isn't in the same Fedimint federation as you", - _ => "the seller doesn't accept your Cashu mint", - }; + let refund = reclaim_spent_ecash(&self.config.data_dir, &token_str, used_backend).await; return Ok(serde_json::json!({ - "error": format!( - "Payment rejected by the seller — {hint}. Your ecash was refunded to \ - your wallet. Try the other ecash type, or use a shared mint/federation." - ) + "error": format!("The seller could not verify the payment. {refund}") })); } @@ -612,13 +697,9 @@ impl RpcHandler { let status = response.status(); let body = response.text().await.unwrap_or_default(); tracing::warn!("paid download: seller {onion} returned {status}: {body}"); - reclaim_spent_ecash(&self.config.data_dir, &token_str, used_backend).await; - let reason = serde_json::from_str::(&body) - .ok() - .and_then(|v| v.get("error").and_then(|e| e.as_str()).map(str::to_string)) - .unwrap_or_else(|| format!("Peer returned an error ({status}).")); + let refund = reclaim_spent_ecash(&self.config.data_dir, &token_str, used_backend).await; return Ok(serde_json::json!({ - "error": format!("{reason} Your ecash was refunded to your wallet.") + "error": format!("{} {refund}", seller_error_message(status, &body)) })); } @@ -632,10 +713,17 @@ impl RpcHandler { .filter(|s| !s.is_empty()) .unwrap_or_else(|| "application/octet-stream".to_string()); - let bytes = response - .bytes() - .await - .context("Failed to read response body")?; + let bytes = match response.bytes().await { + Ok(bytes) => bytes, + Err(error) => { + tracing::warn!("paid download: response body failed: {error}"); + let refund = + reclaim_spent_ecash(&self.config.data_dir, &token_str, used_backend).await; + return Ok(serde_json::json!({ + "error": format!("The file transfer was interrupted after payment was sent. {refund}") + })); + } + }; // Persist the purchase so it "stays unlocked" for this buyer: cache the // bytes + metadata keyed by (onion, content_id). The gallery then renders @@ -665,63 +753,41 @@ impl RpcHandler { tracing::warn!("paid download: failed to cache purchased content (non-fatal): {e:#}"); } - // Auto-file the purchase into the user's Files area (2026-07-22): - // Photos for images/video, Music for audio, Documents otherwise — - // same buckets the Cloud view uses. The in-app viewer still plays - // from the purchase cache; this makes the file ALSO show up where - // files live, on every device, without relying on a browser - // download. Best-effort: never fail a paid download over it. - { - let folder = if mime_type.starts_with("image/") || mime_type.starts_with("video/") { - "Photos" - } else if mime_type.starts_with("audio/") { - "Music" - } else { - "Documents" - }; - let base = std::path::Path::new(&filename) - .file_name() - .and_then(|n| n.to_str()) - .unwrap_or("download") - .to_string(); - let dir = self.config.data_dir.join("filebrowser").join(folder); - if let Err(e) = tokio::fs::create_dir_all(&dir).await { - tracing::warn!("paid download: cannot create {}: {e}", dir.display()); - } else { - // Don't clobber an existing file of the same name: "x.jpg" - // → "x (2).jpg" etc. - let mut target = dir.join(&base); - let (stem, ext) = match base.rsplit_once('.') { - Some((s, e)) if !s.is_empty() => (s.to_string(), format!(".{e}")), - _ => (base.clone(), String::new()), - }; - let mut n = 2; - while target.exists() { - target = dir.join(format!("{stem} ({n}){ext}")); - n += 1; - } - match tokio::fs::write(&target, &bytes).await { - Ok(()) => tracing::info!("paid download: filed into {}", target.display()), - Err(e) => tracing::warn!( - "paid download: filing into {} failed (non-fatal): {e}", - target.display() - ), - } - } + // The durable purchased-content cache above is primary. A Files copy + // remains optional: a stopped FileBrowser must not undo a paid download. + let filed = async { + let auth = self.handle_filebrowser_token().await?; + let token = auth + .get("token") + .and_then(|v| v.as_str()) + .context("FileBrowser omitted its authentication token")?; + let client = reqwest::Client::builder() + .no_proxy() + .redirect(reqwest::redirect::Policy::none()) + .timeout(std::time::Duration::from_secs(30)) + .build()?; + file_purchase_in_files( + &client, + "http://127.0.0.1:8083", + token, + &filename, + &mime_type, + &bytes, + ) + .await + } + .await; + match filed { + Ok(path) => tracing::info!("paid download: filed into Files/{path}"), + Err(error) => tracing::warn!( + "paid download: optional Files copy failed; purchase cache retained: {error}" + ), } - use base64::Engine; - let encoded = base64::engine::general_purpose::STANDARD.encode(&bytes); - tracing::info!("paid download: received {} bytes from {onion} (paid {price_sats} sats via {used_backend})", bytes.len()); - Ok(serde_json::json!({ - "data": encoded, - "size": bytes.len(), - "paid_sats": price_sats, - "ecash_backend": used_backend, - "mime_type": mime_type, - "owned": true, - })) + let mut result = paid_content_response(&bytes, &mime_type, price_sats); + result["ecash_backend"] = serde_json::json!(used_backend); + Ok(result) } /// Buyer side (#46): ask the selling node to mint a Lightning invoice for a @@ -1394,3 +1460,7 @@ impl RpcHandler { } } } + +#[cfg(test)] +#[path = "content_tests.rs"] +mod tests; diff --git a/core/archipelago/src/api/rpc/content_tests.rs b/core/archipelago/src/api/rpc/content_tests.rs new file mode 100644 index 00000000..cb8d5833 --- /dev/null +++ b/core/archipelago/src/api/rpc/content_tests.rs @@ -0,0 +1,181 @@ +use super::*; +use hyper::{ + service::{make_service_fn, service_fn}, + Body, Response, Server, +}; +use std::{ + collections::VecDeque, + convert::Infallible, + sync::{Arc, Mutex}, +}; + +struct FilesApi { + url: String, + seen: Arc)>>>, + task: tokio::task::JoinHandle<()>, +} +impl Drop for FilesApi { + fn drop(&mut self) { + self.task.abort(); + } +} +fn files_api(statuses: Vec) -> FilesApi { + let statuses = Arc::new(Mutex::new(VecDeque::from(statuses))); + let seen = Arc::new(Mutex::new(Vec::new())); + let history = seen.clone(); + let server = Server::bind(&([127, 0, 0, 1], 0).into()); + let address = server.local_addr(); + let service = make_service_fn(move |_| { + let statuses = statuses.clone(); + let seen = history.clone(); + async move { + Ok::<_, Infallible>(service_fn(move |request: hyper::Request| { + let statuses = statuses.clone(); + let seen = seen.clone(); + async move { + assert_eq!(request.headers().get("X-Auth").unwrap(), "test-session"); + let method = request.method().to_string(); + let uri = request.uri().to_string(); + let body = hyper::body::to_bytes(request.into_body()) + .await + .unwrap() + .to_vec(); + seen.lock().unwrap().push((method, uri, body)); + let status = statuses + .lock() + .unwrap() + .pop_front() + .expect("unexpected extra Files request"); + Ok::<_, Infallible>( + Response::builder() + .status(status) + .body(Body::empty()) + .unwrap(), + ) + } + })) + } + }); + FilesApi { + url: format!("http://{address}"), + seen, + task: tokio::spawn(async move { + server.serve(service).await.unwrap(); + }), + } +} + +#[test] +fn first_and_cached_paid_downloads_have_the_same_client_payload_contract() { + use base64::Engine; + for paid in [0, 1] { + let response = paid_content_response(&[0, 255, 123], "application/octet-stream", paid); + assert_eq!(response["data"], response["data_base64"]); + assert_eq!( + base64::engine::general_purpose::STANDARD + .decode(response["data"].as_str().unwrap()) + .unwrap(), + [0, 255, 123] + ); + assert_eq!(response["size"], 3); + assert_eq!(response["size_bytes"], 3); + assert_eq!(response["paid_sats"], paid); + assert_eq!(response["owned"], true); + } +} + +#[tokio::test] +async fn files_copy_uses_authenticated_api_and_preserves_existing_names() { + let api = files_api(vec![200, 409, 200]); + let client = reqwest::Client::new(); + let path = file_purchase_in_files( + &client, + &api.url, + "test-session", + "../my #file?.txt", + "text/plain", + b"paid bytes", + ) + .await + .unwrap(); + assert_eq!(path, "Documents/my #file? (2).txt"); + let seen = api.seen.lock().unwrap(); + assert_eq!(seen[0].0, "GET"); + assert_eq!(seen[0].1, "/api/resources/Documents/"); + assert_eq!(seen.len(), 3); + for (_, uri, body) in &seen[1..] { + assert!(uri.contains("override=false")); + assert!(uri.contains("%23file%3F")); + assert!(!uri.contains("../")); + assert_eq!(body, b"paid bytes"); + } +} + +#[tokio::test] +async fn files_copy_creates_missing_media_folder() { + for (mime, folder) in [ + ("image/png", "Photos"), + ("video/mp4", "Photos"), + ("audio/ogg", "Music"), + ] { + let api = files_api(vec![404, 200, 200]); + let path = file_purchase_in_files( + &reqwest::Client::new(), + &api.url, + "test-session", + "file", + mime, + b"bytes", + ) + .await + .unwrap(); + assert_eq!(path, format!("{folder}/file")); + let seen = api.seen.lock().unwrap(); + assert_eq!(seen[1].0, "POST"); + assert!(seen[1].1.ends_with('/')); + assert!(seen[1].2.is_empty()); + assert_eq!(seen[2].2, b"bytes"); + } +} + +#[tokio::test] +async fn files_copy_fails_without_overwriting_or_claiming_success_on_errors() { + for statuses in [ + vec![401], + vec![503], + vec![404, 500], + vec![200, 507], + vec![200, 403], + ] { + let expected = statuses.len(); + let api = files_api(statuses); + assert!(file_purchase_in_files( + &reqwest::Client::new(), + &api.url, + "test-session", + "file.txt", + "text/plain", + b"bytes" + ) + .await + .is_err()); + assert_eq!(api.seen.lock().unwrap().len(), expected); + } +} + +#[test] +fn seller_errors_are_bounded_printable_and_identified_as_peer_text() { + let status = reqwest::StatusCode::SERVICE_UNAVAILABLE; + let message = seller_error_message(status, r#"{"error":"Cannot read file\n\u0000"}"#); + assert!(message.starts_with("Seller response (503")); + assert!(message.ends_with("Cannot read file")); + assert!(!message.contains('\n') && !message.contains('\0')); + let body = serde_json::json!({"error": "é".repeat(1000)}).to_string(); + assert!(seller_error_message(status, &body).chars().count() < 300); + for body in ["not JSON", r#"{"error": 7}"#, r#"{"error":" "}"#] { + assert_eq!( + seller_error_message(status, body), + "Peer returned an error (503 Service Unavailable)." + ); + } +} diff --git a/core/archipelago/src/api/rpc/lnd/info.rs b/core/archipelago/src/api/rpc/lnd/info.rs index ea12108f..6d94a6c5 100644 --- a/core/archipelago/src/api/rpc/lnd/info.rs +++ b/core/archipelago/src/api/rpc/lnd/info.rs @@ -109,7 +109,50 @@ fn checked_balances( )) } +fn bitcoin_wait_state( + installed: bool, + running: bool, + fresh: bool, + ibd: Option, +) -> (&'static str, &'static str) { + if !installed { + ("waiting_install", "Waiting for Bitcoin to be installed") + } else if !running { + ("waiting_start", "Waiting for Bitcoin to start") + } else if !fresh || ibd.is_none() { + ("waiting_start", "Waiting for Bitcoin to start") + } else if ibd == Some(true) { + ("waiting_sync", "Waiting for Bitcoin to sync") + } else { + ("bitcoin_ready", "Bitcoin is ready") + } +} + impl RpcHandler { + pub(crate) async fn handle_lnd_readiness(&self) -> serde_json::Value { + let (data, _) = self.state_manager.get_snapshot().await; + if !data.server_info.status_info.containers_scanned { + return serde_json::json!({"state":"checking", "message":"Checking Bitcoin availability"}); + } + let nodes: Vec<_> = ["bitcoin-core", "bitcoin-knots", "bitcoin"] + .iter() + .filter_map(|id| data.package_data.get(*id)) + .collect(); + let installed = !nodes.is_empty(); + let running = nodes + .iter() + .any(|p| p.state == crate::data_model::PackageState::Running); + let bitcoin = crate::bitcoin_status::get_bitcoin_status().await; + let ibd = bitcoin + .blockchain_info + .as_ref() + .and_then(|v| v.get("initialblockdownload")) + .and_then(|v| v.as_bool()); + let (state, message) = + bitcoin_wait_state(installed, running, bitcoin.ok && !bitcoin.stale, ibd); + serde_json::json!({"state": state, "message": message}) + } + pub(in crate::api::rpc) async fn handle_lnd_getinfo(&self) -> Result { let macaroon_bytes = read_lnd_admin_macaroon().await?; let macaroon_hex = hex::encode(&macaroon_bytes); @@ -419,3 +462,44 @@ mod tests { assert!(!is_valid_identity_pubkey(&"g".repeat(66))); } } + +#[cfg(test)] +mod dependency_readiness_tests { + use super::bitcoin_wait_state; + #[test] + fn waiting_states_cover_install_start_sync_outage_and_recovery() { + assert_eq!( + bitcoin_wait_state(false, false, false, None).0, + "waiting_install" + ); + assert_eq!( + bitcoin_wait_state(true, false, false, None).0, + "waiting_start" + ); + assert_eq!( + bitcoin_wait_state(true, true, false, None).0, + "waiting_start" + ); + assert_eq!( + bitcoin_wait_state(true, true, true, Some(true)).0, + "waiting_sync" + ); + assert_eq!( + bitcoin_wait_state(true, true, true, Some(false)).0, + "bitcoin_ready" + ); + // Previously synced cached information must not hide a current outage. + assert_eq!( + bitcoin_wait_state(true, true, false, Some(false)).0, + "waiting_start" + ); + assert_eq!( + bitcoin_wait_state(true, true, true, None).0, + "waiting_start" + ); + assert_eq!( + bitcoin_wait_state(true, true, true, Some(false)).0, + "bitcoin_ready" + ); + } +} diff --git a/core/archipelago/src/api/rpc/lnd/mod.rs b/core/archipelago/src/api/rpc/lnd/mod.rs index c2909968..ef734ccf 100644 --- a/core/archipelago/src/api/rpc/lnd/mod.rs +++ b/core/archipelago/src/api/rpc/lnd/mod.rs @@ -133,12 +133,36 @@ async fn stream_lnd_transactions(sm: &crate::state::StateManager) -> Result<()> /// RPC-unreachable and locked-wallet states are deliberately NOT handled /// here — container-down is crash-recovery's job, and unlocking needs the /// operator. +fn bitcoin_ready_for_lnd_watchdog(status: &crate::bitcoin_status::BitcoinNodeStatus) -> bool { + status.ok + && !status.stale + && status.age_ms < 30_000 + && status + .blockchain_info + .as_ref() + .and_then(|v| v.get("initialblockdownload")) + .and_then(|v| v.as_bool()) + == Some(false) +} + pub(crate) fn spawn_lnd_health_watchdog() { tokio::spawn(async move { let mut bad_minutes: u32 = 0; let mut last_restart: Option = None; + let mut last_height: Option = None; loop { tokio::time::sleep(std::time::Duration::from_secs(60)).await; + // Initial Bitcoin sync, warmup, and outages are dependencies to + // wait for, never evidence that LND is wedged. Do not accumulate + // restart pressure during a days-long initial block download. + let bitcoin = crate::bitcoin_status::get_bitcoin_status().await; + if !bitcoin_ready_for_lnd_watchdog(&bitcoin) + || crate::app_ops::lifecycle_op_in_flight("lnd") + { + bad_minutes = 0; + last_height = None; + continue; + } let Ok(bytes) = read_lnd_admin_macaroon().await else { bad_minutes = 0; // no LND on this node (or not set up yet) continue; @@ -161,6 +185,10 @@ pub(crate) fn spawn_lnd_health_watchdog() { bad_minutes = 0; // down/locked — not the wedge signature continue; }; + if !resp.status().is_success() { + bad_minutes = 0; + continue; + } let Ok(info) = resp.json::().await else { bad_minutes = 0; continue; @@ -182,7 +210,12 @@ pub(crate) fn spawn_lnd_health_watchdog() { .get("num_pending_channels") .and_then(|v| v.as_u64()) .unwrap_or(0); - let wedged = !synced || (channels > 0 && peers == 0); + let height = info.get("block_height").and_then(|v| v.as_u64()); + let progressing = height + .zip(last_height) + .is_some_and(|(now, before)| now > before); + last_height = height; + let wedged = !progressing && (!synced || (channels > 0 && peers == 0)); if !wedged { bad_minutes = 0; continue; @@ -239,3 +272,31 @@ impl RpcHandler { Ok((client, macaroon_hex)) } } + +#[cfg(test)] +mod watchdog_dependency_tests { + use super::bitcoin_ready_for_lnd_watchdog; + use crate::bitcoin_status::BitcoinNodeStatus; + use serde_json::json; + #[test] + fn initial_sync_warmup_outage_stale_and_unknown_never_trigger_lnd_restart() { + let mut status = BitcoinNodeStatus::default(); + assert!(!bitcoin_ready_for_lnd_watchdog(&status)); + status.ok = true; + status.blockchain_info = Some(json!({"initialblockdownload":true})); + assert!(!bitcoin_ready_for_lnd_watchdog(&status)); + status.blockchain_info = Some(json!({"initialblockdownload":false})); + assert!(bitcoin_ready_for_lnd_watchdog(&status)); + status.stale = true; + assert!(!bitcoin_ready_for_lnd_watchdog(&status)); + status.stale = false; + status.ok = false; + assert!(!bitcoin_ready_for_lnd_watchdog(&status)); + status.ok = true; + status.age_ms = 30_000; + assert!(!bitcoin_ready_for_lnd_watchdog(&status)); + status.age_ms = 0; + status.blockchain_info = Some(json!({})); + assert!(!bitcoin_ready_for_lnd_watchdog(&status)); + } +} diff --git a/core/archipelago/src/api/rpc/package/install.rs b/core/archipelago/src/api/rpc/package/install.rs index 9ee88898..86259562 100644 --- a/core/archipelago/src/api/rpc/package/install.rs +++ b/core/archipelago/src/api/rpc/package/install.rs @@ -326,6 +326,10 @@ impl RpcHandler { // an older version pins it so install_fresh resolves that image and the // update badge stays suppressed. See docs/bitcoin-multi-version-design.md. if matches!(package_id, "bitcoin-core" | "bitcoin-knots") { + if let Some(value) = params.get("prune") { + let prune = value.as_bool().context("prune must be a boolean")?; + crate::settings::bitcoin_storage::save(&self.config.data_dir, prune).await?; + } if let Some(version) = params.get("version").and_then(|v| v.as_str()) { persist_install_version_selection(package_id, version).await; } diff --git a/core/archipelago/src/api/rpc/package/set_config.rs b/core/archipelago/src/api/rpc/package/set_config.rs index f0c50222..9bf7d3e4 100644 --- a/core/archipelago/src/api/rpc/package/set_config.rs +++ b/core/archipelago/src/api/rpc/package/set_config.rs @@ -153,8 +153,18 @@ impl RpcHandler { let default = app_catalog::catalog_default_version(app_id); let cfg = version_config::read(app_id); let installed = installed_version(app_id).await; + let bitcoin_prune = if matches!(app_id, "bitcoin-core" | "bitcoin-knots") { + Some( + crate::settings::bitcoin_storage::load(&self.config.data_dir) + .await? + .prune, + ) + } else { + None + }; Ok(serde_json::json!({ + "bitcoinPrune": bitcoin_prune, "id": app_id, "supportsVersions": supports_versions(app_id), "default": default, diff --git a/core/archipelago/src/bitcoin_status.rs b/core/archipelago/src/bitcoin_status.rs index f53ae4a4..66b4747c 100644 --- a/core/archipelago/src/bitcoin_status.rs +++ b/core/archipelago/src/bitcoin_status.rs @@ -100,7 +100,11 @@ fn friendly_transient_error(has_cached_state: bool, err_msg: &str) -> String { .trim() .trim_end_matches('.'); let lower = detail.to_lowercase(); - let state = if lower.contains("verifying blocks") { + let state = if lower.contains("loading block index") { + Some("loading its block index. This can take a while after installation or restart") + } else if lower.contains("replaying blocks") { + Some("checking saved blocks before startup completes") + } else if lower.contains("verifying blocks") { Some("verifying blocks after restart") } else if lower.contains("connection reset") { Some("starting up and not yet accepting RPC connections") @@ -340,3 +344,21 @@ mod tests { assert!(msg.len() < 260); } } + +#[cfg(test)] +mod startup_message_tests { + #[test] + fn loading_block_index_is_explained_without_rpc_error_dump() { + for cached in [false, true] { + let message = super::friendly_transient_error( + cached, + r#"getblockchaininfo: Bitcoin RPC returned 500 Internal Server Error: {"error":{"code":-28,"message":"Loading block index…"}}"#, + ); + assert!(message.contains("loading its block index")); + for raw in ["500", "-28", "Detail:", "getblockchaininfo", "{", "RPC"] { + assert!(!message.contains(raw)); + } + assert_eq!(message.contains("last known state"), cached); + } + } +} diff --git a/core/archipelago/src/container/companion.rs b/core/archipelago/src/container/companion.rs index 9752454f..42f1b8da 100644 --- a/core/archipelago/src/container/companion.rs +++ b/core/archipelago/src/container/companion.rs @@ -313,7 +313,7 @@ async fn image_id(image_ref: &str) -> Option { /// should reference (`localhost/:latest` for build, registry /// URL for pull). async fn ensure_image_present(spec: &CompanionSpec) -> Result { - let local_image = format!("localhost/{}:latest", spec.image_base); + let mut local_image = format!("localhost/{}:latest", spec.image_base); let local_image_compat = format!("localhost/{}:local", spec.image_base); let registry_image = format!("{}/{}:latest", COMPANION_REGISTRY, spec.image_base); @@ -322,11 +322,13 @@ async fn ensure_image_present(spec: &CompanionSpec) -> Result { for dir in spec.build_dir_candidates { let dockerfile = PathBuf::from(dir).join("Dockerfile"); if fs::try_exists(&dockerfile).await.unwrap_or(false) { - // `:local` is a deliberate manual override — never auto-rebuild it. + // Older installers and self-update create :local themselves. It + // must receive source updates too; treating it as a permanent + // manual override silently kept the old LND UI after an OTA. if image_exists(&local_image_compat).await { - return Ok(local_image_compat); + local_image = local_image_compat.clone(); } - // Reuse the auto-built `:latest` only when the build context has NOT + // Reuse either local tag only when the build context has NOT // changed since it was built. Without this staleness check an // already-present image is reused forever, so edits to the baked-in // context (Dockerfile, nginx.conf, …) never reach the node — this is @@ -849,20 +851,43 @@ async fn needs_repair(spec: &CompanionSpec) -> Result { if !matches_known_shape { return Ok(true); } - if on_disk.contains(&local_image) && !on_disk.contains(&local_image_compat) { + if let Some(image) = managed_local_image(spec, &on_disk) { for dir in spec.build_dir_candidates { let dockerfile = PathBuf::from(dir).join("Dockerfile"); if fs::try_exists(&dockerfile).await.unwrap_or(false) { // Conservative on any timeout/error inside: reuse the cache. - return Ok(context_is_newer_than_image(dir, &local_image).await); + return Ok(context_is_newer_than_image(dir, &image).await); } } } Ok(false) } +fn managed_local_image(spec: &CompanionSpec, unit: &str) -> Option { + ["latest", "local"] + .iter() + .map(|tag| format!("localhost/{}:{tag}", spec.image_base)) + .find(|image| build_unit(spec, image).render() == unit) +} + #[cfg(test)] mod tests { + #[test] + fn legacy_installer_local_tag_is_checked_for_source_updates_like_latest() { + for spec in ALL_COMPANIONS.iter().flat_map(|group| group.iter()) { + for tag in ["local", "latest"] { + let image = format!("localhost/{}:{tag}", spec.image_base); + let unit = build_unit(spec, &image).render(); + assert_eq!(managed_local_image(spec, &unit), Some(image)); + } + let registry = format!("{}/{}:latest", COMPANION_REGISTRY, spec.image_base); + assert_eq!( + managed_local_image(spec, &build_unit(spec, ®istry).render()), + None + ); + } + } + use super::*; fn names(specs: &[&'static CompanionSpec]) -> Vec<&'static str> { diff --git a/core/archipelago/src/container/lnd.rs b/core/archipelago/src/container/lnd.rs index 19f24ee2..4b745d57 100644 --- a/core/archipelago/src/container/lnd.rs +++ b/core/archipelago/src/container/lnd.rs @@ -89,18 +89,74 @@ bitcoind.estimatemode=ECONOMICAL\n" Ok(EnsureOutcome::Written) } +/// Bitcoin can accept TCP while returning RPC_IN_WARMUP for many minutes. +/// Unlocking LND then triggers its short chain-backend timeout and a restart loop. +/// Leave the wallet intact and locked; the next reconciliation retries readiness. +async fn bitcoin_rpc_ready() -> bool { + let (user, password) = crate::bitcoin_rpc::bitcoin_rpc_credentials().await; + let client = match reqwest::Client::builder() + .no_proxy() + .timeout(std::time::Duration::from_secs(5)) + .build() + { + Ok(client) => client, + Err(_) => return false, + }; + let response = client.post(crate::constants::BITCOIN_RPC_URL) + .basic_auth(user, Some(password)) + .json(&serde_json::json!({"jsonrpc":"1.0","id":"lnd-readiness","method":"getblockchaininfo","params":[]})) + .send().await; + match response { + Ok(response) if response.status().is_success() => response + .json::() + .await + .is_ok_and(|value| bitcoin_readiness_response(&value)), + _ => false, + } +} + +fn bitcoin_readiness_response(value: &serde_json::Value) -> bool { + value.get("error").is_none_or(|e| e.is_null()) + && value + .pointer("/result/blocks") + .and_then(|v| v.as_u64()) + .is_some() + && value + .pointer("/result/initialblockdownload") + .and_then(|v| v.as_bool()) + .is_some() +} + pub async fn ensure_wallet_initialized() -> Result<()> { let admin_macaroon = "/var/lib/archipelago/lnd/data/chain/bitcoin/mainnet/admin.macaroon"; let wallet_db = "/var/lib/archipelago/lnd/data/chain/bitcoin/mainnet/wallet.db"; if file_exists_as_root(wallet_db).await { + // GetInfo can wait for Bitcoin sync even though the wallet is already + // unlocked. State RPC stays available during that normal startup phase. + let client = reqwest::Client::builder() + .no_proxy() + .timeout(std::time::Duration::from_secs(5)) + .danger_accept_invalid_certs(true) + .build()?; + if wallet_is_unlocked(wallet_state(&client).await.as_deref()) { + return Ok(()); + } if file_exists_as_root(admin_macaroon).await && lnd_getinfo_ready(admin_macaroon).await { return Ok(()); } + if !bitcoin_rpc_ready().await { + tracing::debug!("[lnd] waiting for Bitcoin RPC readiness before wallet unlock"); + return Ok(()); + } unlock_existing_wallet_no_wipe().await?; wait_for_admin_macaroon(admin_macaroon).await?; return Ok(()); } + if !bitcoin_rpc_ready().await { + tracing::debug!("[lnd] waiting for Bitcoin RPC readiness before wallet initialization"); + return Ok(()); + } init_wallet_via_rest().await?; wait_for_admin_macaroon(admin_macaroon).await } @@ -258,6 +314,9 @@ async fn unlock_existing_wallet_via_rest() -> Result { // exactly the nodes least able to afford it. Waiting longer costs nothing — // a wrong password still exits on the first pass via `all_rejected`. for _ in 0..UNLOCK_NOT_READY_ATTEMPTS { + if wallet_is_unlocked(wallet_state(&client).await.as_deref()) { + return Ok(true); + } let mut all_rejected = true; for pw in &candidates { match try_unlock_once(&client, pw).await { @@ -294,6 +353,10 @@ pub(crate) async fn unlock_existing_wallet_no_wipe() -> Result<()> { } } +fn wallet_is_unlocked(state: Option<&str>) -> bool { + matches!(state, Some("UNLOCKED" | "RPC_ACTIVE" | "SERVER_ACTIVE")) +} + /// Current LND wallet state via the unauthenticated `/v1/state` endpoint /// (NON_EXISTING / LOCKED / UNLOCKED / RPC_ACTIVE / …). None if unreachable. async fn wallet_state(client: &reqwest::Client) -> Option { @@ -1089,3 +1152,44 @@ mod tests { .is_empty()); } } + +#[cfg(test)] +mod bitcoin_readiness_tests { + use super::bitcoin_readiness_response; + use serde_json::json; + #[test] + fn only_usable_bitcoin_rpc_allows_wallet_unlock() { + for response in [ + json!({}), + json!({"error":{"code":-28,"message":"Loading block index"},"result":null}), + json!({"result":{"blocks":null}}), + ] { + assert!(!bitcoin_readiness_response(&response)); + } + // Initial sync is supported by LND. Loading the database is not. + for ibd in [true, false] { + assert!(bitcoin_readiness_response( + &json!({"result":{"blocks":100,"initialblockdownload":ibd},"error":null}) + )); + } + } +} + +#[cfg(test)] +mod syncing_wallet_state_tests { + #[test] + fn an_unlocked_wallet_waiting_for_chain_sync_is_never_unlocked_again() { + for state in ["UNLOCKED", "RPC_ACTIVE", "SERVER_ACTIVE"] { + assert!(super::wallet_is_unlocked(Some(state))); + } + for state in [ + None, + Some("LOCKED"), + Some("NON_EXISTING"), + Some("WAITING_TO_START"), + Some("unknown"), + ] { + assert!(!super::wallet_is_unlocked(state)); + } + } +} diff --git a/core/archipelago/src/container/prod_orchestrator.rs b/core/archipelago/src/container/prod_orchestrator.rs index d7e6299a..801193ef 100644 --- a/core/archipelago/src/container/prod_orchestrator.rs +++ b/core/archipelago/src/container/prod_orchestrator.rs @@ -798,6 +798,10 @@ fn host_port_bindings_drifted( } async fn ensure_user_podman_socket() -> Result<()> { + // Unit tests inject a runtime; they must not restart the host Podman API. + if cfg!(test) { + return Ok(()); + } let socket_path = "/run/user/1000/podman/podman.sock"; if podman_socket_accepts_connections(socket_path).await { return Ok(()); @@ -1170,15 +1174,21 @@ impl ReconcileReport { fn cascade_pairs_for_report<'r>( report: &'r ReconcileReport, user_stopped: &std::collections::HashSet, + changed_backends: &HashSet, ) -> Vec<(&'r str, &'static str)> { let mut pairs = Vec::new(); for (backend, action) in &report.actions { if !matches!( action, - ReconcileAction::Installed | ReconcileAction::Started + ReconcileAction::NoOp | ReconcileAction::Started | ReconcileAction::Installed ) { continue; } + // A successful systemctl start can be a no-op after a transient + // Podman inspect failure. Require a witnessed lifecycle change. + if !changed_backends.contains(backend) { + continue; + } for dep in crate::app_ops::address_caching_dependents(backend) { let dep_untouched = report .actions @@ -1192,6 +1202,25 @@ fn cascade_pairs_for_report<'r>( pairs } +/// Only positive runtime evidence permits disrupting an address-caching wallet. +/// A known absent/stopped backend becoming running, a new container ID, or a +/// changed start timestamp qualifies. A failed observation never does. +fn backend_instance_changed(before: Option<&ContainerStatus>, after: &ContainerStatus) -> bool { + if after.state != ContainerState::Running || after.id.is_empty() { + return false; + } + let Some(before) = before else { + return true; + }; + if before.id.is_empty() { + return false; + } + if before.id != after.id || before.state != ContainerState::Running { + return true; + } + matches!((&before.started_at, &after.started_at), (Some(a), Some(b)) if !a.is_empty() && !b.is_empty() && a != b) +} + #[derive(Debug, Default)] pub struct AdoptionReport { pub adopted: Vec, @@ -1905,14 +1934,40 @@ impl ProdContainerOrchestrator { _ => 2, }); // Live container names (any state), for the same recovery check. - let present_containers: std::collections::HashSet = self - .runtime - .list_containers() - .await - .map(|cs| cs.into_iter().map(|c| c.name).collect()) + let listed_containers = self.runtime.list_containers().await.ok(); + let present_containers: HashSet = listed_containers + .as_ref() + .map(|cs| cs.iter().map(|c| c.name.clone()).collect()) .unwrap_or_default(); + // Keep unknown distinct from confirmed absence. Runtime queries can + // fail under load while systemd still has a healthy running backend. + let mut backend_before: HashMap> = HashMap::new(); + for lm in &manifests { + let id = &lm.manifest.app.id; + if crate::app_ops::address_caching_dependents(id).is_empty() { + continue; + } + let name = compute_container_name(&lm.manifest); + match self.runtime.get_container_status(&name).await { + Ok(status) => { + backend_before.insert(id.clone(), Some(status)); + } + Err(_) if listed_containers.is_some() && !present_containers.contains(&name) => { + backend_before.insert(id.clone(), None); + } + Err(err) => { + tracing::warn!(backend = %id, error = %err, + "cannot observe backend before reconcile; will not infer a dependency restart from an action report"); + } + } + } let mut report = ReconcileReport::default(); let disk_gb = self.disk_gb().await; + let bitcoin_pruned = disk_gb < ARCHIVAL_BITCOIN_DISK_GB + || crate::settings::bitcoin_storage::load(&self.data_dir) + .await + .map(|settings| settings.prune) + .unwrap_or(true); // Register every candidate before the (sequential, possibly slow) // pass so the scanner overlays queued-but-down apps as Restarting // instead of Stopped. Each app is deregistered as its turn finishes, @@ -1952,7 +2007,7 @@ impl ProdContainerOrchestrator { } if mode == ReconcileMode::ExistingOnly && requires_archival_bitcoin(&app_id) - && disk_gb < ARCHIVAL_BITCOIN_DISK_GB + && bitcoin_pruned { report.record( &app_id, @@ -2087,7 +2142,20 @@ impl ProdContainerOrchestrator { // state recovery, repair recreate, boot InstallMissing) moves the // address behind a running dependent's back — §C "restart lnd after // ANY bitcoin recreate". - for (backend, dep) in cascade_pairs_for_report(&report, &user_stopped) { + let mut changed_backends = HashSet::new(); + for (backend, before) in &backend_before { + let Some(name) = container_name_by_app_id.get(backend) else { + continue; + }; + if let Ok(after) = self.runtime.get_container_status(name).await { + if backend_instance_changed(before.as_ref(), &after) { + changed_backends.insert(backend.clone()); + } + } + } + // A user stop during a slow reconcile pass still takes precedence. + let user_stopped = crate::crash_recovery::load_user_stopped(&self.data_dir).await; + for (backend, dep) in cascade_pairs_for_report(&report, &user_stopped, &changed_backends) { // Same rule as the RPC cascade: hold the dependent's op lock // across the restart; skip when a worker is mid-sequence. let lock = crate::app_ops::op_lock(dep); @@ -3226,6 +3294,9 @@ impl ProdContainerOrchestrator { } async fn ensure_container_network(&self, manifest: &AppManifest) -> Result<()> { + if cfg!(test) { + return Ok(()); + } let Some(network) = manifest.app.container.network.as_deref() else { return Ok(()); }; @@ -3720,6 +3791,17 @@ impl ProdContainerOrchestrator { } let mut env = manifest.app.environment.clone(); env.extend(manifest.app.container.resolve_derived_env(&facts)); + if matches!(manifest.app.id.as_str(), "bitcoin-core" | "bitcoin-knots") { + let storage = crate::settings::bitcoin_storage::load(&self.data_dir).await?; + env.retain(|entry| !entry.starts_with("BITCOIN_PRUNE=")); + if storage.prune { + anyhow::ensure!( + manifest.app.container.custom_args.iter().any(|arg| arg.contains("BITCOIN_PRUNE")), + "This Bitcoin app definition cannot honor the pruning choice. Refresh the app catalog and try again." + ); + env.push("BITCOIN_PRUNE=1".to_string()); + } + } // FM_BITCOIND_URL now comes from the manifest's {{BITCOIN_HOST}} // derived_env (works on Knots/Core/any distro). The old hardcoded @@ -6073,6 +6155,48 @@ app: ); } + #[tokio::test] + async fn bitcoin_storage_choice_is_applied_and_old_catalog_cannot_silently_ignore_it() { + let rt = Arc::new(MockRuntime::default()); + let mut orch = orch_with(rt).await; + let dir = tempfile::tempdir().unwrap(); + orch.set_data_dir(dir.path().to_path_buf()); + for id in ["bitcoin-core", "bitcoin-knots"] { + let mut old = pull_manifest(id, "docker.io/bitcoin/bitcoin:28"); + // No preference: existing containers need no new environment flag. + crate::settings::bitcoin_storage::save(dir.path(), false) + .await + .unwrap(); + orch.resolve_dynamic_env(&mut old).await.unwrap(); + assert!(!old + .app + .environment + .iter() + .any(|s| s.starts_with("BITCOIN_PRUNE="))); + crate::settings::bitcoin_storage::save(dir.path(), true) + .await + .unwrap(); + assert!(orch + .resolve_dynamic_env(&mut old) + .await + .unwrap_err() + .to_string() + .contains("cannot honor")); + let mut current = pull_manifest(id, "docker.io/bitcoin/bitcoin:28"); + current + .app + .container + .custom_args + .push("if [ ${BITCOIN_PRUNE:-0} = 1 ]; then :; fi".into()); + orch.resolve_dynamic_env(&mut current).await.unwrap(); + assert!(current + .app + .environment + .iter() + .any(|s| s == "BITCOIN_PRUNE=1")); + } + } + #[tokio::test] async fn install_resolves_derived_and_secret_env_before_create() { let rt = Arc::new(MockRuntime::default()); @@ -6344,6 +6468,67 @@ app: ); } + #[test] + fn backend_cascade_requires_observed_instance_change() { + let running = ContainerStatus { + id: "container-1".into(), + name: "bitcoin-core".into(), + state: ContainerState::Running, + started_at: Some("start-1".into()), + health: None, + exit_code: None, + image: "bitcoin:1".into(), + created: "created-1".into(), + ports: vec![], + lan_address: None, + }; + assert!(!backend_instance_changed(Some(&running), &running)); + assert!(backend_instance_changed(None, &running)); + let mut before = running.clone(); + before.state = ContainerState::Exited; + assert!(backend_instance_changed(Some(&before), &running)); + before = running.clone(); + before.id = "old-container".into(); + assert!(backend_instance_changed(Some(&before), &running)); + before = running.clone(); + before.started_at = Some("earlier-start".into()); + assert!(backend_instance_changed(Some(&before), &running)); + before.started_at = None; + assert!(!backend_instance_changed(Some(&before), &running)); + before.id.clear(); + assert!(!backend_instance_changed(Some(&before), &running)); + let mut after = running.clone(); + after.state = ContainerState::Exited; + assert!(!backend_instance_changed(None, &after)); + after = running.clone(); + after.id.clear(); + assert!(!backend_instance_changed(None, &after)); + } + + #[test] + fn cascade_ignores_false_started_report_but_detects_real_exec_drift() { + let none = HashSet::new(); + let mut report = ReconcileReport { + actions: vec![ + ("bitcoin-core".into(), ReconcileAction::Started), + ("lnd".into(), ReconcileAction::NoOp), + ], + failures: vec![], + }; + // systemctl start of an already active unit does not move its address. + assert!(cascade_pairs_for_report(&report, &none, &none).is_empty()); + // A unit exec rewrite can restart Bitcoin while the outer reconcile + // action remains NoOp. Runtime evidence still requires LND to reconnect. + let changed = ["bitcoin-core".into()].into(); + report.actions[0].1 = ReconcileAction::NoOp; + assert_eq!( + cascade_pairs_for_report(&report, &none, &changed), + vec![("bitcoin-core", "lnd")] + ); + report.actions[0].1 = ReconcileAction::Left("lifecycle-op-in-flight".into()); + assert!(cascade_pairs_for_report(&report, &none, &changed).is_empty()); + } + #[test] fn cascade_pairs_cover_backend_recreate_with_running_dependent() { use std::collections::HashSet; @@ -6355,6 +6540,7 @@ app: failures: vec![], }; let none = HashSet::new(); + let changed: HashSet = ["bitcoin-core".into(), "bitcoin-knots".into()].into(); // Backend recreated while lnd sat running (NoOp) → cascade. let r = report(vec![ @@ -6362,7 +6548,7 @@ app: ("lnd", ReconcileAction::NoOp), ]); assert_eq!( - cascade_pairs_for_report(&r, &none), + cascade_pairs_for_report(&r, &none, &changed), vec![("bitcoin-knots", "lnd")] ); @@ -6372,7 +6558,7 @@ app: ("lnd", ReconcileAction::NoOp), ]); assert_eq!( - cascade_pairs_for_report(&r, &none), + cascade_pairs_for_report(&r, &none, &changed), vec![("bitcoin-core", "lnd")] ); @@ -6381,7 +6567,7 @@ app: ("bitcoin-knots", ReconcileAction::NoOp), ("lnd", ReconcileAction::NoOp), ]); - assert!(cascade_pairs_for_report(&r, &none).is_empty()); + assert!(cascade_pairs_for_report(&r, &none, &none).is_empty()); // Dependent itself (re)started this pass → it already resolved the // fresh address; no cascade. @@ -6389,7 +6575,7 @@ app: ("bitcoin-knots", ReconcileAction::Installed), ("lnd", ReconcileAction::Started), ]); - assert!(cascade_pairs_for_report(&r, &none).is_empty()); + assert!(cascade_pairs_for_report(&r, &none, &changed).is_empty()); // User-stopped dependent is never bounced. let r = report(vec![ @@ -6397,14 +6583,14 @@ app: ("lnd", ReconcileAction::NoOp), ]); let stopped: HashSet = ["lnd".to_string()].into(); - assert!(cascade_pairs_for_report(&r, &stopped).is_empty()); + assert!(cascade_pairs_for_report(&r, &stopped, &changed).is_empty()); // Non-backend recreates don't cascade anything. let r = report(vec![ ("grafana", ReconcileAction::Installed), ("lnd", ReconcileAction::NoOp), ]); - assert!(cascade_pairs_for_report(&r, &none).is_empty()); + assert!(cascade_pairs_for_report(&r, &none, &changed).is_empty()); } #[tokio::test] diff --git a/core/archipelago/src/container/quadlet.rs b/core/archipelago/src/container/quadlet.rs index 00a5b7fc..a6150398 100644 --- a/core/archipelago/src/container/quadlet.rs +++ b/core/archipelago/src/container/quadlet.rs @@ -184,6 +184,7 @@ pub struct QuadletUnit { pub no_new_privileges: bool, pub cpu_quota: Option, pub restart_policy: RestartPolicy, + pub stop_grace_secs: Option, } impl QuadletUnit { @@ -216,6 +217,10 @@ impl QuadletUnit { let _ = writeln!(s, "[Container]"); let _ = writeln!(s, "ContainerName={}", self.name); let _ = writeln!(s, "Image={}", self.image); + let grace = self + .stop_grace_secs + .unwrap_or_else(|| archipelago_container::runtime::stop_grace_secs_for(&self.name)); + let _ = writeln!(s, "StopTimeout={grace}"); // Pull=never: companions are pre-pulled or built. A missing image // must surface as a unit start failure, not a silent retry storm. let _ = writeln!(s, "Pull=never"); @@ -350,6 +355,15 @@ impl QuadletUnit { // the unit stuck in deactivating. Health/status remains app-level state, // not a systemd start gate. let _ = writeln!(s, "TimeoutStartSec=0"); + let _ = writeln!(s, "TimeoutStopSec={}", grace.saturating_add(15)); + // Stop explicitly before Quadlet's generated `podman rm -f`. The + // existing container may still carry Podman's old 10-second default; + // StopTimeout alone only protects containers created after migration. + let _ = writeln!(s, "ExecStop="); + let _ = writeln!( + s, + "ExecStop=/usr/bin/podman stop --ignore --time={grace} --cidfile=%t/%N.cid" + ); // Restart policy + 10s backoff. RestartSec keeps a crash-loop // from saturating the journal. Companions: Always. Backends: // OnFailure (clean stops stay stopped). @@ -525,6 +539,9 @@ impl QuadletUnit { // Always, not OnFailure: with quadlet's `--rm`, OnFailure left a // cleanly-exited app deleted and unrestarted. See RestartPolicy. restart_policy: RestartPolicy::Always, + stop_grace_secs: Some(super::prod_orchestrator::resolve_stop_grace_secs( + manifest, name, + )), } } } @@ -676,6 +693,13 @@ pub async fn unit_exists(name: &str) -> bool { /// Resolve the per-user quadlet dir under $HOME. Created if missing. pub async fn unit_dir() -> Result { + #[cfg(test)] + { + static TEST_UNITS: std::sync::OnceLock = std::sync::OnceLock::new(); + return Ok(TEST_UNITS + .get_or_init(|| tempfile::tempdir().unwrap().keep()) + .clone()); + } let home = std::env::var_os("HOME") .map(PathBuf::from) .ok_or_else(|| anyhow!("HOME not set; cannot locate quadlet unit dir"))?; @@ -785,7 +809,11 @@ pub async fn stop_service(service: &str) -> Result<()> { /// corruption — so the orchestrator passes the per-app grace here. Never waits /// less than `QUADLET_STOP_TIMEOUT`. pub async fn stop_service_with_timeout(service: &str, timeout: Duration) -> Result<()> { - let timeout = timeout.max(QUADLET_STOP_TIMEOUT); + let name = service.strip_suffix(".service").unwrap_or(service); + let body = fs::read_to_string(unit_dir().await?.join(format!("{name}.container"))) + .await + .unwrap_or_default(); + let timeout = timeout.max(stop_wait_timeout(name, &body)); match systemctl_user_status(&["stop", service], timeout).await { Ok(status) if status.success() => Ok(()), Ok(status) => Err(anyhow!("systemctl --user stop {service} exited {status}")), @@ -806,10 +834,29 @@ pub async fn stop_service_with_timeout(service: &str, timeout: Duration) -> Resu } } +/// The command waiter must outlive both the container grace and systemd's +/// stop deadline. Restart/repair callers must not kill Bitcoin at 45 seconds. +fn stop_wait_timeout(name: &str, unit_body: &str) -> Duration { + Duration::from_secs(stop_grace_from_unit(name, unit_body).saturating_add(30)) + .max(QUADLET_STOP_TIMEOUT) +} + +fn stop_grace_from_unit(name: &str, unit_body: &str) -> u64 { + directive_values(unit_body, "StopTimeout=") + .last() + .and_then(|value| value.parse::().ok()) + .unwrap_or_else(|| archipelago_container::runtime::stop_grace_secs_for(name)) +} + async fn systemctl_user_status( args: &[&str], timeout: Duration, ) -> Result { + #[cfg(test)] + { + use std::os::unix::process::ExitStatusExt; + return Ok(std::process::ExitStatus::from_raw(0)); + } let mut cmd = Command::new("systemctl"); cmd.arg("--user").args(args); cmd.kill_on_drop(true); @@ -856,6 +903,10 @@ async fn wait_not_deactivating(service: &str, timeout: Duration) -> bool { } async fn systemctl_user_output(args: &[&str], timeout: Duration) -> Result { + #[cfg(test)] + { + anyhow::bail!("Unit tests have no real user service manager"); + } let mut cmd = Command::new("systemctl"); cmd.arg("--user").args(args); cmd.kill_on_drop(true); @@ -923,6 +974,10 @@ fn directive_values(unit_body: &str, prefix: &str) -> Vec { /// that systemd no longer knows about. pub async fn disable_remove(unit_name: &str, dir: &Path) -> Result<()> { let svc = format!("{unit_name}.service"); + let path = dir.join(format!("{unit_name}.container")); + let body = fs::read_to_string(&path).await.unwrap_or_default(); + let timeout = stop_wait_timeout(unit_name, &body); + let grace = stop_grace_from_unit(unit_name, &body).to_string(); // Stop first; ignore failure (unit may already be down). BOUNDED — on // rootless podman a generated unit can wedge in "deactivating" while // `podman rm -f` hangs underneath it, and an unbounded `systemctl stop` @@ -930,13 +985,12 @@ pub async fn disable_remove(unit_name: &str, dir: &Path) -> Result<()> { // the package entry is stranded in `Removing` (a ghost in My Apps that also // blocks reinstall). If the graceful stop times out, escalate to // SIGKILL + reset-failed so teardown always proceeds. - if systemctl_user_status(&["stop", &svc], QUADLET_STOP_TIMEOUT) + if systemctl_user_status(&["stop", &svc], timeout) .await .is_err() { let _ = kill_and_reset_service(&svc).await; } - let path = dir.join(format!("{unit_name}.container")); if fs::try_exists(&path).await.unwrap_or(false) { match fs::remove_file(&path).await { Ok(()) => {} @@ -949,9 +1003,9 @@ pub async fn disable_remove(unit_name: &str, dir: &Path) -> Result<()> { // Bounded so a hung podman store can't re-introduce the stall this function // exists to avoid. let _ = tokio::time::timeout( - QUADLET_STOP_TIMEOUT, + timeout, Command::new("podman") - .args(["rm", "-f", unit_name]) + .args(["rm", "-f", "--ignore", "--time", &grace, unit_name]) .status(), ) .await; @@ -960,6 +1014,9 @@ pub async fn disable_remove(unit_name: &str, dir: &Path) -> Result<()> { /// Is the quadlet-generated service currently active? pub async fn is_active(service: &str) -> bool { + if cfg!(test) { + return false; + } Command::new("systemctl") .args(["--user", "is-active", "--quiet", service]) .status() @@ -973,6 +1030,118 @@ mod tests { use super::*; use tempfile::tempdir; + #[test] + fn shutdown_grace_covers_container_systemd_and_caller() { + for (name, grace) in [ + ("bitcoin-core", 600), + ("bitcoin-knots", 600), + ("lnd", 330), + ("electrumx", 300), + ("other", 30), + ] { + let unit = QuadletUnit { + name: name.into(), + ..Default::default() + }; + let body = unit.render(); + assert!(body.contains(&format!("StopTimeout={grace}\n"))); + assert!(body.contains(&format!("TimeoutStopSec={}\n", grace + 15))); + assert!(body.contains(&format!("podman stop --ignore --time={grace} --cidfile="))); + assert_eq!( + stop_wait_timeout(name, &body), + Duration::from_secs(grace + 30) + ); + // Legacy units have no StopTimeout directive yet. + assert_eq!(stop_wait_timeout(name, ""), Duration::from_secs(grace + 30)); + } + } + + #[test] + fn custom_stop_grace_survives_render_and_restart_budget() { + let manifest: AppManifest = serde_yaml::from_str( + r#" +app: + id: custom-db + name: Custom database + version: 1.0.0 + stop_grace_secs: 900 + container: + image: example/db:1 +"#, + ) + .unwrap(); + let unit = QuadletUnit::from_manifest(&manifest, "custom-db"); + assert_eq!(unit.stop_grace_secs, Some(900)); + assert_eq!( + stop_wait_timeout("custom-db", &unit.render()), + Duration::from_secs(930) + ); + assert_eq!( + stop_wait_timeout("lnd", "StopTimeout=invalid"), + Duration::from_secs(360) + ); + } + + #[test] + fn stop_grace_migration_does_not_request_an_execution_restart() { + let unit = sample_unit(); + let new = unit.render(); + let old = new + .lines() + .filter(|line| { + !line.starts_with("StopTimeout=") + && !line.starts_with("TimeoutStopSec=") + && !line.starts_with("ExecStop=") + }) + .collect::>() + .join("\n"); + assert!(!exec_changed(&old, &new)); + assert!(!publish_ports_changed(&old, &new)); + assert!(!network_aliases_changed(&old, &new)); + assert!(!health_cmd_changed(&old, &new)); + } + + #[test] + fn actual_quadlet_generator_stops_before_forced_removal() { + let generator = Path::new("/usr/lib/systemd/system-generators/podman-system-generator"); + if !generator.exists() { + eprintln!( + "Quadlet generator unavailable; run this regression on the Linux release host" + ); + return; + } + let dir = tempdir().unwrap(); + let unit = QuadletUnit { + name: "grace-test".into(), + image: "localhost/test:latest".into(), + stop_grace_secs: Some(600), + ..Default::default() + }; + std::fs::write(dir.path().join("grace-test.container"), unit.render()).unwrap(); + let output = std::process::Command::new(generator) + .args(["--user", "--dryrun"]) + .env("QUADLET_UNIT_DIRS", dir.path()) + .output() + .unwrap(); + assert!( + output.status.success(), + "{}", + String::from_utf8_lossy(&output.stderr) + ); + let generated = String::from_utf8_lossy(&output.stdout).to_string() + + &String::from_utf8_lossy(&output.stderr); + let stop = generated + .find("ExecStop=/usr/bin/podman stop --ignore --time=600") + .unwrap(); + let remove = generated.find("ExecStop=/usr/bin/podman rm ").unwrap(); + assert!( + stop < remove, + "Legacy container must stop gracefully before removal" + ); + assert!(generated.contains("--stop-timeout 600")); + assert!(generated.contains("TimeoutStopSec=615")); + } + #[test] fn render_emits_secret_env_by_reference_never_value() { let u = QuadletUnit { diff --git a/core/archipelago/src/content_server.rs b/core/archipelago/src/content_server.rs index 72a981f5..0858fda0 100644 --- a/core/archipelago/src/content_server.rs +++ b/core/archipelago/src/content_server.rs @@ -7,7 +7,7 @@ use anyhow::{Context, Result}; use serde::{Deserialize, Serialize}; use std::path::{Path, PathBuf}; use tokio::fs; -use tracing::{debug, info, warn}; +use tracing::{debug, warn}; const CATALOG_FILE: &str = "content/catalog.json"; const CONTENT_DIR: &str = "content/files"; @@ -241,6 +241,8 @@ pub enum ServeResult { /// The catalog entry and file exist but this node can't read the file. /// Returned before any payment is taken. Unavailable, + /// Requested byte range cannot be served; no payment was taken. + RangeNotSatisfiable(u64), } /// Serve a content item by ID with access control and optional range request. @@ -255,6 +257,39 @@ pub async fn serve_content( range: Option, owner_session: bool, ) -> Result { + serve_content_with( + data_dir, + id, + payment_token, + invoice_hash, + peer_did, + range, + owner_session, + |path, range, mime| prepare_content(data_dir, path, range, mime), + |token, amount| async move { verify_payment_token(data_dir, &token, amount).await }, + ) + .await +} + +// Inject only the read and payment boundaries, so tests can prove ordering +// without mint access, file-permission assumptions or privileged commands. +async fn serve_content_with( + data_dir: &Path, + id: &str, + payment_token: Option<&str>, + invoice_hash: Option<&str>, + peer_did: Option<&str>, + range: Option, + owner_session: bool, + read: R, + verify: V, +) -> Result +where + R: FnOnce(PathBuf, Option, String) -> RF, + RF: std::future::Future>, + V: FnOnce(String, u64) -> VF, + VF: std::future::Future, +{ let catalog = load_catalog(data_dir).await?; let item = match catalog.items.iter().find(|i| i.id == id) { Some(i) => i, @@ -317,18 +352,28 @@ pub async fn serve_content( return Ok(ServeResult::NotFound); } - // Confirm the file is readable BEFORE the paid gate below redeems the - // buyer's token. Reading it only afterwards meant a permission error - // surfaced after the sale: the buyer was charged and got an error - // instead of the file (2026-09-29, a FileBrowser upload left 0640). - if let Err(e) = ensure_readable(&file_path).await { - warn!( - content_id = %id, - path = %file_path.display(), - "shared content file is not readable by this node: {e:#}" - ); - return Ok(ServeResult::Unavailable); + // Refuse unauthorized viewers before opening or reading any bytes. + if !owner_session && matches!(item.access, AccessControl::PeersOnly) && !is_known_peer { + return Ok(ServeResult::Forbidden); } + if !owner_session { + if let AccessControl::Paid { price_sats, .. } = &item.access { + if payment_token.is_none() && invoice_hash.is_none() { + return Ok(ServeResult::PaymentRequired(*price_sats)); + } + } + } + + // Finish all file I/O before consuming bearer payment. Merely opening then + // reopening after charging still lost payments on read errors or deletion. + let prepared = match read(file_path, range, item.mime_type.clone()).await { + Ok(result @ (ServeResult::Ok(..) | ServeResult::Partial { .. })) => result, + Ok(other) => return Ok(other), + Err(error) => { + warn!(content_id = %id, "Cannot prepare shared content: {error:#}"); + return Ok(ServeResult::Unavailable); + } + }; // Check access control if !owner_session { @@ -341,9 +386,13 @@ pub async fn serve_content( // Each path only counts when the sharer accepts that method. let mut authorized = false; if let Some(token) = payment_token { - if (method_accepted(&item.access, "ecash") - || method_accepted(&item.access, "fedimint")) - && verify_payment_token(data_dir, token, *price_sats).await + let method = if token.trim().starts_with("cashu") { + "ecash" + } else { + "fedimint" + }; + if method_accepted(&item.access, method) + && verify(token.to_owned(), *price_sats).await { authorized = true; } @@ -370,56 +419,127 @@ pub async fn serve_content( } } + Ok(prepared) +} - let metadata = fs::metadata(&file_path) +async fn prepare_content( + data_dir: &Path, + path: PathBuf, + range: Option, + mime: String, +) -> Result { + use tokio::io::{AsyncReadExt, AsyncSeekExt}; + let mut file = match fs::OpenOptions::new() + .read(true) + .custom_flags(libc::O_NONBLOCK) + .open(&path) .await - .context("Failed to read file metadata")?; - let total_size = metadata.len(); - - // Handle range request for streaming - if let Some(range) = range { - let start = range.start.min(total_size.saturating_sub(1)); - let end = range - .end - .map(|e| e.min(total_size - 1)) - .unwrap_or(total_size - 1); - - if start > end || start >= total_size { - return Ok(ServeResult::NotFound); + { + Ok(file) => file, + Err(error) if error.kind() == std::io::ErrorKind::PermissionDenied => { + let bytes = read_filebrowser_via_userns(data_dir, &path).await?; + return slice_prepared_content(bytes, range, mime); } - - let len = (end - start + 1) as usize; - use tokio::io::{AsyncReadExt, AsyncSeekExt}; - let mut file = tokio::fs::File::open(&file_path) + Err(error) => return Err(error).context("Opening shared content"), + }; + let metadata = file.metadata().await?; + anyhow::ensure!(metadata.is_file(), "Shared content is not a regular file"); + let total = metadata.len(); + if let Some(range) = range { + let Some((start, end)) = checked_range(&range, total) else { + return Ok(ServeResult::RangeNotSatisfiable(total)); + }; + file.seek(std::io::SeekFrom::Start(start)).await?; + let len = usize::try_from(end - start + 1).context("Content range is too large")?; + let mut bytes = vec![0; len]; + file.read_exact(&mut bytes) .await - .context("Failed to open content file")?; - file.seek(std::io::SeekFrom::Start(start)) - .await - .context("Failed to seek")?; - let mut buf = vec![0u8; len]; - file.read_exact(&mut buf) - .await - .context("Failed to read range")?; - - debug!( - "Serving content '{}' range {}-{}/{} ({} bytes)", - id, start, end, total_size, len - ); + .context("Reading shared content range")?; return Ok(ServeResult::Partial { - bytes: buf, - mime_type: item.mime_type.clone(), + bytes, + mime_type: mime, start, end, - total: total_size, + total, }); } - - let bytes = fs::read(&file_path) + let mut bytes = Vec::new(); + file.read_to_end(&mut bytes) .await - .context("Failed to read content file")?; + .context("Reading shared content")?; + Ok(ServeResult::Ok(bytes, mime)) +} - debug!("Serving content '{}' ({} bytes)", id, bytes.len()); - Ok(ServeResult::Ok(bytes, item.mime_type.clone())) +fn checked_range(range: &ByteRange, total: u64) -> Option<(u64, u64)> { + let last = total.checked_sub(1)?; + let end = range.end.unwrap_or(last).min(last); + (range.start <= end && range.start < total).then_some((range.start, end)) +} + +fn slice_prepared_content( + bytes: Vec, + range: Option, + mime: String, +) -> Result { + let total = bytes.len() as u64; + match range { + None => Ok(ServeResult::Ok(bytes, mime)), + Some(range) => match checked_range(&range, total) { + Some((start, end)) => Ok(ServeResult::Partial { + bytes: bytes[start as usize..=end as usize].to_vec(), + mime_type: mime, + start, + end, + total, + }), + None => Ok(ServeResult::RangeNotSatisfiable(total)), + }, + } +} + +/// Read only an explicitly shared, regular file within FileBrowser storage. +/// Do not change its mode or grant world-readable access to paid/private data. +async fn filebrowser_read_path(data_dir: &Path, path: &Path) -> Result { + let root = fs::canonicalize(data_dir.join("filebrowser")).await?; + let target = fs::canonicalize(path).await?; + anyhow::ensure!( + target.starts_with(&root) && target != root, + "Shared file is outside Files storage" + ); + anyhow::ensure!( + fs::metadata(&target).await?.is_file(), + "Shared content is not a regular file" + ); + Ok(target) +} + +async fn read_filebrowser_via_userns(data_dir: &Path, path: &Path) -> Result> { + let path = filebrowser_read_path(data_dir, path).await?; + // Tests exercise the boundary explicitly; they never launch the host Podman. + #[cfg(test)] + { + let _ = path; + anyhow::bail!("Files namespace read disabled in unit tests") + } + #[cfg(not(test))] + { + let output = tokio::time::timeout( + std::time::Duration::from_secs(900), + tokio::process::Command::new("podman") + .args(["unshare", "cat", "--"]) + .arg(path) + .kill_on_drop(true) + .output(), + ) + .await + .context("Files namespace read timed out")??; + anyhow::ensure!( + output.status.success(), + "Files namespace read failed: {}", + output.status + ); + Ok(output.stdout) + } } /// Result of attempting to serve a preview. @@ -589,64 +709,8 @@ pub async fn serve_content_preview(data_dir: &Path, id: &str) -> Result Result<()> { - ensure_readable_with(path, grant_read_access).await -} - -async fn ensure_readable_with(path: &Path, grant: F) -> Result<()> -where - F: FnOnce(PathBuf) -> Fut, - Fut: std::future::Future>, -{ - match fs::File::open(path).await { - Ok(_) => return Ok(()), - Err(e) if e.kind() == std::io::ErrorKind::PermissionDenied => {} - Err(e) => return Err(e).context("Failed to open content file"), - } - grant(path.to_path_buf()).await?; - info!("Granted read access to shared content file {}", path.display()); - fs::File::open(path) - .await - .context("Content file still unreadable after chmod")?; - Ok(()) -} - -// Tests must not shell out to podman: whether it exists (and can chmod a -// file the test user owns) would decide the outcome. -#[cfg(not(test))] -use grant_read_via_podman as grant_read_access; - -#[cfg(test)] -async fn grant_read_access(_path: PathBuf) -> Result<()> { - anyhow::bail!("granting read access is disabled in tests") -} - -#[cfg_attr(test, allow(dead_code))] -async fn grant_read_via_podman(path: PathBuf) -> Result<()> { - let out = tokio::process::Command::new("podman") - .args(["unshare", "chmod", "a+r"]) - .arg(&path) - .output() - .await - .context("Failed to run podman unshare chmod")?; - if !out.status.success() { - anyhow::bail!( - "podman unshare chmod a+r failed: {}", - String::from_utf8_lossy(&out.stderr).trim() - ); - } - Ok(()) -} - /// Verify a payment token covers the required amount. -/// Accepts both cashuA tokens (real Cashu) and legacy cashuSend_ format. +/// Accepts real Cashu tokens and Fedimint notes. /// Swaps proofs at the mint to verify they're unspent before accepting. async fn verify_payment_token(data_dir: &Path, token: &str, required_sats: u64) -> bool { match crate::wallet::ecash::verify_and_receive_payment(data_dir, token, required_sats).await { @@ -800,135 +864,299 @@ mod prune_missing_content_tests { } #[cfg(test)] -mod unreadable_content_tests { +mod paid_read_order_tests { use super::*; - use std::os::unix::fs::PermissionsExt; + use std::sync::atomic::{AtomicUsize, Ordering}; - /// Writes `bytes` to the FileBrowser area and makes it unreadable, the - /// way a 0640 upload owned by a container subuid looks to this service. - /// `None` when the test runs as root, where mode bits don't stop reads. - fn unreadable_file(data_dir: &Path, name: &str) -> Option { - let dir = data_dir.join("filebrowser").join("Music"); - std::fs::create_dir_all(&dir).unwrap(); - let path = dir.join(name); - std::fs::write(&path, b"audio").unwrap(); - std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o000)).unwrap(); - std::fs::File::open(&path).is_err().then_some(path) - } - - fn paid_item(filename: &str) -> ContentItem { - ContentItem { - id: "paid-item".to_string(), - filename: filename.to_string(), - mime_type: "audio/mpeg".to_string(), - size_bytes: 5, - description: String::new(), - access: AccessControl::Paid { - price_sats: 10, - accepted: vec!["ecash".to_string()], - }, - availability: Availability::AllPeers, - added_at: "2026-01-01T00:00:00Z".to_string(), - } - } - - /// Regression (2026-09-29): the seller redeemed the buyer's token and - /// only then failed to read the file, so the buyer paid for nothing. - /// An unreadable file must be refused before the payment gate runs, - /// which is why a token that would never verify still gets Unavailable - /// rather than PaymentRequired. - #[tokio::test] - async fn an_unreadable_paid_file_is_refused_before_any_payment_is_taken() { + async fn fixture(bytes: &[u8]) -> tempfile::TempDir { let dir = tempfile::tempdir().unwrap(); - let data_dir = dir.path(); - let Some(_path) = unreadable_file(data_dir, "song.mp3") else { - return; // running as root - }; + fs::create_dir_all(dir.path().join("content/files")) + .await + .unwrap(); + fs::write(dir.path().join("content/files/test.bin"), bytes) + .await + .unwrap(); save_catalog( - data_dir, + dir.path(), &ContentCatalog { - items: vec![paid_item("Music/song.mp3")], + items: vec![ContentItem { + id: "paid".into(), + filename: "test.bin".into(), + mime_type: "application/octet-stream".into(), + size_bytes: bytes.len() as u64, + description: String::new(), + access: AccessControl::Paid { + price_sats: 10, + accepted: vec!["ecash".into()], + }, + availability: Availability::AllPeers, + added_at: "2026-09-30".into(), + }], }, ) .await .unwrap(); + dir + } - let result = serve_content( - data_dir, - "paid-item", - Some("cashuBnot-a-real-token"), + #[tokio::test] + async fn all_read_failures_precede_redemption_even_as_root() { + for kind in [ + std::io::ErrorKind::PermissionDenied, + std::io::ErrorKind::UnexpectedEof, + std::io::ErrorKind::NotFound, + std::io::ErrorKind::Other, + ] { + let dir = fixture(b"abc").await; + let charged = AtomicUsize::new(0); + let result = serve_content_with( + dir.path(), + "paid", + Some("cashuBtest"), + None, + None, + None, + false, + |_, _, _| async move { Err(std::io::Error::from(kind).into()) }, + |_, _| async { + charged.fetch_add(1, Ordering::SeqCst); + true + }, + ) + .await + .unwrap(); + assert!(matches!(result, ServeResult::Unavailable)); + assert_eq!(charged.load(Ordering::SeqCst), 0); + assert_eq!(load_catalog(dir.path()).await.unwrap().items.len(), 1); + } + } + + #[tokio::test] + async fn deletion_during_payment_cannot_lose_prepared_bytes() { + let dir = fixture(b"original").await; + let result = serve_content_with( + dir.path(), + "paid", + Some("cashuBtest"), None, None, None, false, - ) - .await - .unwrap(); - assert!(matches!(result, ServeResult::Unavailable)); - // An unreadable file is not a missing one: keep the catalog entry. - assert_eq!(load_catalog(data_dir).await.unwrap().items.len(), 1); - } - - #[tokio::test] - async fn a_readable_paid_file_still_demands_payment() { - let dir = tempfile::tempdir().unwrap(); - let data_dir = dir.path(); - let music = data_dir.join("filebrowser").join("Music"); - std::fs::create_dir_all(&music).unwrap(); - std::fs::write(music.join("song.mp3"), b"audio").unwrap(); - save_catalog( - data_dir, - &ContentCatalog { - items: vec![paid_item("Music/song.mp3")], + |path, range, mime| prepare_content(dir.path(), path, range, mime), + |_, amount| { + assert_eq!(amount, 10); + async { + fs::remove_file(dir.path().join("content/files/test.bin")) + .await + .unwrap(); + true + } }, ) .await .unwrap(); - - let result = serve_content(data_dir, "paid-item", None, None, None, None, false) - .await - .unwrap(); - assert!(matches!(result, ServeResult::PaymentRequired(10))); + assert!(matches!(result, ServeResult::Ok(bytes, _) if bytes == b"original")); } #[tokio::test] - async fn ensure_readable_grants_access_once_then_reopens() { - let dir = tempfile::tempdir().unwrap(); - let Some(path) = unreadable_file(dir.path(), "a.mp3") else { - return; - }; - let calls = std::sync::atomic::AtomicUsize::new(0); - ensure_readable_with(&path, |p| { - calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst); - async move { - std::fs::set_permissions(&p, std::fs::Permissions::from_mode(0o644))?; - Ok(()) - } - }) + async fn empty_out_of_bounds_and_reversed_ranges_never_charge() { + for (bytes, start, end) in [ + (b"".as_slice(), 0, None), + (b"abc".as_slice(), 3, None), + (b"abc".as_slice(), 2, Some(1)), + ] { + let dir = fixture(bytes).await; + let result = serve_content_with( + dir.path(), + "paid", + Some("cashuBtest"), + None, + None, + Some(ByteRange { start, end }), + false, + |path, range, mime| prepare_content(dir.path(), path, range, mime), + |_, _| async { panic!("invalid range reached payment") }, + ) + .await + .unwrap(); + assert!( + matches!(result, ServeResult::RangeNotSatisfiable(n) if n == bytes.len() as u64) + ); + } + } + + #[tokio::test] + async fn prepared_range_survives_file_change_while_payment_is_verified() { + let dir = fixture(b"abcdef").await; + let result = serve_content_with( + dir.path(), + "paid", + Some("cashuBtest"), + None, + None, + Some(ByteRange { + start: 2, + end: Some(999), + }), + false, + |path, range, mime| prepare_content(dir.path(), path, range, mime), + |_, _| async { + fs::write(dir.path().join("content/files/test.bin"), b"x") + .await + .unwrap(); + true + }, + ) .await .unwrap(); - assert_eq!(calls.load(std::sync::atomic::Ordering::SeqCst), 1); + assert!( + matches!(result, ServeResult::Partial { bytes, start: 2, end: 5, total: 6, .. } if bytes == b"cdef") + ); } #[tokio::test] - async fn ensure_readable_leaves_a_readable_file_alone() { - let dir = tempfile::tempdir().unwrap(); - let path = dir.path().join("ok.mp3"); - std::fs::write(&path, b"x").unwrap(); - ensure_readable_with(&path, |_| async { anyhow::bail!("must not grant") }) + async fn payment_denial_never_returns_prepared_content() { + let dir = fixture(b"secret").await; + let charged = AtomicUsize::new(0); + let result = serve_content_with( + dir.path(), + "paid", + Some("cashuBtest"), + None, + None, + None, + false, + |path, range, mime| prepare_content(dir.path(), path, range, mime), + |_, _| async { + charged.fetch_add(1, Ordering::SeqCst); + false + }, + ) + .await + .unwrap(); + assert!(matches!(result, ServeResult::PaymentRequired(10))); + assert_eq!(charged.load(Ordering::SeqCst), 1); + } + + #[tokio::test] + async fn missing_payment_and_peer_restrictions_precede_file_reads() { + let dir = fixture(b"secret").await; + let result = serve_content_with( + dir.path(), + "paid", + None, + None, + None, + None, + false, + |_, _, _| async { panic!("unauthorized file read") }, + |_, _| async { panic!("unexpected payment") }, + ) + .await + .unwrap(); + assert!(matches!(result, ServeResult::PaymentRequired(10))); + let mut catalog = load_catalog(dir.path()).await.unwrap(); + catalog.items[0].access = AccessControl::PeersOnly; + save_catalog(dir.path(), &catalog).await.unwrap(); + let result = serve_content_with( + dir.path(), + "paid", + None, + None, + None, + None, + false, + |_, _, _| async { panic!("unauthorized file read") }, + |_, _| async { panic!("unexpected payment") }, + ) + .await + .unwrap(); + assert!(matches!(result, ServeResult::Forbidden)); + } + + #[tokio::test] + async fn owner_reads_paid_content_without_redemption() { + let dir = fixture(b"own file").await; + let result = serve_content_with( + dir.path(), + "paid", + None, + None, + None, + None, + true, + |path, range, mime| prepare_content(dir.path(), path, range, mime), + |_, _| async { panic!("owner charged") }, + ) + .await + .unwrap(); + assert!(matches!(result, ServeResult::Ok(bytes, _) if bytes == b"own file")); + } + + #[tokio::test] + async fn directory_in_place_of_file_does_not_charge() { + let dir = fixture(b"abc").await; + let path = dir.path().join("content/files/test.bin"); + fs::remove_file(&path).await.unwrap(); + fs::create_dir(&path).await.unwrap(); + let result = serve_content_with( + dir.path(), + "paid", + Some("cashuBtest"), + None, + None, + None, + false, + |path, range, mime| prepare_content(dir.path(), path, range, mime), + |_, _| async { panic!("directory charged") }, + ) + .await + .unwrap(); + assert!(matches!(result, ServeResult::Unavailable)); + } + + #[tokio::test] + async fn files_namespace_read_is_scoped_to_regular_files_and_keeps_mode() { + use std::os::unix::fs::{symlink, PermissionsExt}; + let dir = fixture(b"outside").await; + let root = dir.path().join("filebrowser"); + fs::create_dir(&root).await.unwrap(); + let inside = root.join("song"); + fs::write(&inside, b"song").await.unwrap(); + fs::set_permissions(&inside, std::fs::Permissions::from_mode(0o640)) .await .unwrap(); + assert_eq!( + filebrowser_read_path(dir.path(), &inside).await.unwrap(), + inside + ); + assert_eq!( + fs::metadata(&inside).await.unwrap().permissions().mode() & 0o777, + 0o640 + ); + let outside = dir.path().join("content/files/test.bin"); + symlink(&outside, root.join("escape")).unwrap(); + for path in [outside, root.join("escape"), root.clone()] { + assert!(filebrowser_read_path(dir.path(), &path).await.is_err()); + } } - #[tokio::test] - async fn ensure_readable_reports_a_failed_grant() { - let dir = tempfile::tempdir().unwrap(); - let Some(path) = unreadable_file(dir.path(), "b.mp3") else { - return; - }; - let err = ensure_readable_with(&path, |_| async { anyhow::bail!("no podman") }) - .await - .unwrap_err(); - assert!(err.to_string().contains("no podman")); + #[test] + fn user_namespace_bytes_use_the_same_range_rules() { + assert!(matches!( + slice_prepared_content( + vec![], + Some(ByteRange { + start: 0, + end: None + }), + "x".into() + ) + .unwrap(), + ServeResult::RangeNotSatisfiable(0) + )); + assert!( + matches!(slice_prepared_content(b"abc".to_vec(), Some(ByteRange { start: 1, end: None }), "x".into()).unwrap(), ServeResult::Partial { bytes, start: 1, end: 2, total: 3, .. } if bytes == b"bc") + ); } } diff --git a/core/archipelago/src/fips/dial.rs b/core/archipelago/src/fips/dial.rs index df4a35ba..02838efc 100644 --- a/core/archipelago/src/fips/dial.rs +++ b/core/archipelago/src/fips/dial.rs @@ -132,7 +132,21 @@ pub fn client() -> reqwest::Client { /// before the Tor fallback ever gets a chance. The generous `connect_timeout` /// is preserved so a cold hole-punched path still gets time to establish. pub fn client_with_timeout(timeout: Duration) -> reqwest::Client { + client_with_delivery_policy(timeout, false) +} + +fn delivery_redirect_policy(single: bool) -> reqwest::redirect::Policy { + if single { + reqwest::redirect::Policy::none() + } else { + reqwest::redirect::Policy::default() + } +} + +fn client_with_delivery_policy(timeout: Duration, single: bool) -> reqwest::Client { reqwest::Client::builder() + .no_proxy() + .redirect(delivery_redirect_policy(single)) .timeout(timeout) .connect_timeout(Duration::from_secs(8)) .user_agent("archipelago-fips/1") @@ -488,7 +502,7 @@ impl<'a> PeerRequest<'a> { // Use the FIPS reply unless it's one a Tor retry could // fix (404 path-not-served / 5xx) and we're allowed to // fall back. FIPS-only never falls back. - if pref == TransportPref::Fips || !fips_should_fall_back(resp.status()) { + if fips_answer_is_final(pref, self.single_delivery, resp.status()) { telemetry::record_fips_ok(); self.spawn_record(crate::transport::TransportKind::Fips); return Ok((resp, crate::transport::TransportKind::Fips)); @@ -597,13 +611,21 @@ impl<'a> PeerRequest<'a> { } else { budget }; - let c = client_with_timeout(per_attempt); + let c = client_with_delivery_policy(per_attempt, self.single_delivery); let mut rb = c.post(&url).json(body); for (k, v) in &self.headers { rb = rb.header(*k, v); } - match tokio::time::timeout(budget, send_with_retry(rb)).await { + let single = self.single_delivery; + let attempt = send_with_retry_if(rb, |e| fips_retryable(single, e)); + match tokio::time::timeout(budget, attempt).await { Ok(Ok(r)) => Ok(Some(r)), + Ok(Err(e)) if single && !e.is_connect() => Err(anyhow::anyhow!( + "FIPS POST failed after possible delivery; not replaying: {e}" + )), + Err(_) if single => Err(anyhow::anyhow!( + "FIPS POST exceeded its budget after possible delivery; not replaying" + )), Ok(Err(e)) => { telemetry::record_fallback(FallbackReason::ConnectFail); tracing::info!( @@ -658,7 +680,7 @@ impl<'a> PeerRequest<'a> { } else { budget }; - let c = client_with_timeout(per_attempt); + let c = client_with_delivery_policy(per_attempt, self.single_delivery); let mut rb = c.get(&url); for (k, v) in &self.headers { rb = rb.header(*k, v); @@ -737,6 +759,7 @@ impl<'a> PeerRequest<'a> { .context("Invalid Tor SOCKS proxy URL")?; reqwest::Client::builder() .proxy(proxy) + .redirect(delivery_redirect_policy(self.single_delivery)) .timeout(self.timeout) .build() .context("Build Tor HTTP client") @@ -909,3 +932,92 @@ mod tests { assert!(fips_retryable(true, &err)); } } + +#[cfg(test)] +mod delivery_redirect_tests { + use super::*; + use hyper::{ + service::{make_service_fn, service_fn}, + Body, Response, Server, + }; + use std::{ + convert::Infallible, + sync::{ + atomic::{AtomicUsize, Ordering}, + Arc, + }, + }; + + #[tokio::test] + async fn paid_bearer_request_does_not_follow_redirects_but_normal_get_does() { + let seen = Arc::new(AtomicUsize::new(0)); + let counter = seen.clone(); + let server = Server::bind(&([127, 0, 0, 1], 0).into()); + let address = server.local_addr(); + let service = make_service_fn(move |_| { + let counter = counter.clone(); + async move { + Ok::<_, Infallible>(service_fn(move |request: hyper::Request| { + let counter = counter.clone(); + async move { + counter.fetch_add(1, Ordering::SeqCst); + let response = if request.uri().path() == "/first" { + Response::builder() + .status(302) + .header("Location", "/replay") + .body(Body::empty()) + .unwrap() + } else { + Response::new(Body::from("replayed")) + }; + Ok::<_, Infallible>(response) + } + })) + } + }); + let task = tokio::spawn(server.serve(service)); + let url = format!("http://{address}/first"); + let response = client_with_delivery_policy(Duration::from_secs(2), true) + .get(&url) + .header("X-Payment-Token", "dummy-test-token") + .send() + .await + .unwrap(); + assert_eq!(response.status(), reqwest::StatusCode::FOUND); + assert_eq!(seen.load(Ordering::SeqCst), 1); + let response = client_with_delivery_policy(Duration::from_secs(2), false) + .get(url) + .send() + .await + .unwrap(); + assert_eq!(response.status(), reqwest::StatusCode::OK); + assert_eq!(seen.load(Ordering::SeqCst), 3); + task.abort(); + } + + #[tokio::test] + async fn paid_request_is_not_resent_when_peer_disconnects_after_reading_it() { + use tokio::io::AsyncReadExt; + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let seen = Arc::new(AtomicUsize::new(0)); + let counter = seen.clone(); + let task = tokio::spawn(async move { + while let Ok((mut stream, _)) = listener.accept().await { + let mut buf = [0; 4096]; + let _ = stream.read(&mut buf).await; + counter.fetch_add(1, Ordering::SeqCst); + drop(stream); + } + }); + let c = client_with_delivery_policy(Duration::from_secs(2), true); + let error = send_with_retry_if(c.get(format!("http://{address}/")), |e| { + fips_retryable(true, e) + }) + .await + .unwrap_err(); + assert!(!error.is_connect()); + assert_eq!(seen.load(Ordering::SeqCst), 1); + task.abort(); + } +} diff --git a/core/archipelago/src/settings/bitcoin_storage.rs b/core/archipelago/src/settings/bitcoin_storage.rs new file mode 100644 index 00000000..17b1c4e5 --- /dev/null +++ b/core/archipelago/src/settings/bitcoin_storage.rs @@ -0,0 +1,51 @@ +//! Install-time pruning preference, shared by Bitcoin Core and Knots. +//! Missing preference preserves the existing disk-based automatic selection. +use anyhow::{Context, Result}; +use serde::{Deserialize, Serialize}; +use std::path::Path; + +#[derive(Default, Serialize, Deserialize)] +pub struct BitcoinStorage { + pub prune: bool, +} + +pub async fn load(data_dir: &Path) -> Result { + match tokio::fs::read(data_dir.join("settings/bitcoin-storage.json")).await { + Ok(bytes) => serde_json::from_slice(&bytes).context("Invalid Bitcoin storage settings"), + Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(BitcoinStorage::default()), + Err(e) => Err(e.into()), + } +} + +pub async fn save(data_dir: &Path, prune: bool) -> Result<()> { + let dir = data_dir.join("settings"); + tokio::fs::create_dir_all(&dir).await?; + let path = dir.join("bitcoin-storage.json"); + let temporary = dir.join("bitcoin-storage.json.tmp"); + tokio::fs::write(&temporary, serde_json::to_vec(&BitcoinStorage { prune })?).await?; + tokio::fs::rename(temporary, path).await?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + #[tokio::test] + async fn missing_setting_keeps_auto_and_explicit_pruning_survives_reload() { + let dir = tempfile::tempdir().unwrap(); + assert!(!load(dir.path()).await.unwrap().prune); + save(dir.path(), true).await.unwrap(); + assert!(load(dir.path()).await.unwrap().prune); + save(dir.path(), false).await.unwrap(); + assert!(!load(dir.path()).await.unwrap().prune); + } + #[tokio::test] + async fn corrupt_setting_is_not_silently_changed_to_archival() { + let dir = tempfile::tempdir().unwrap(); + save(dir.path(), true).await.unwrap(); + tokio::fs::write(dir.path().join("settings/bitcoin-storage.json"), "broken") + .await + .unwrap(); + assert!(load(dir.path()).await.is_err()); + } +} diff --git a/core/archipelago/src/settings/mod.rs b/core/archipelago/src/settings/mod.rs index db3bd5a5..8ffaf648 100644 --- a/core/archipelago/src/settings/mod.rs +++ b/core/archipelago/src/settings/mod.rs @@ -7,3 +7,5 @@ pub mod ai_permissions; pub mod session_policy; pub mod transport; + +pub mod bitcoin_storage; diff --git a/core/archipelago/src/update.rs b/core/archipelago/src/update.rs index 51455303..8c27b1f9 100644 --- a/core/archipelago/src/update.rs +++ b/core/archipelago/src/update.rs @@ -1481,6 +1481,21 @@ pub async fn cancel_download(data_dir: &Path) -> Result<()> { /// service unit that inherits systemd's default protections (i.e. none /// of ours), escaping the namespace. pub(crate) async fn host_sudo(args: &[&str]) -> Result { + #[cfg(test)] + { + anyhow::ensure!( + std::env::var("ARCHY_TEST_ISOLATED").as_deref() == Ok("1"), + "Host-operation tests require scripts/test-backend-isolated.sh" + ); + let (program, args) = args.split_first().context("Missing test command")?; + // Run inside the test namespace, never escape through sudo/systemd-run. + return tokio::process::Command::new(program) + .args(args) + .status() + .await + .context("isolated test command failed"); + } + let mut full: Vec<&str> = vec![ "systemd-run", "--wait", @@ -1505,6 +1520,21 @@ pub(crate) async fn host_sudo(args: &[&str]) -> Result /// Same mechanism as `host_sudo` but captures stdout — for read-only probes /// (e.g. `stat`) where the answer is in the output, not the exit status. pub(crate) async fn host_sudo_output(args: &[&str]) -> Result { + #[cfg(test)] + { + anyhow::ensure!( + std::env::var("ARCHY_TEST_ISOLATED").as_deref() == Ok("1"), + "Host-operation tests require scripts/test-backend-isolated.sh" + ); + let (program, args) = args.split_first().context("Missing test command")?; + // Run inside the test namespace, never escape through sudo/systemd-run. + return tokio::process::Command::new(program) + .args(args) + .output() + .await + .context("isolated test command failed"); + } + let mut full: Vec<&str> = vec![ "systemd-run", "--wait", diff --git a/core/archipelago/src/wallet/ecash.rs b/core/archipelago/src/wallet/ecash.rs index e77647a8..ecc5c932 100644 --- a/core/archipelago/src/wallet/ecash.rs +++ b/core/archipelago/src/wallet/ecash.rs @@ -775,7 +775,9 @@ pub async fn send_token_at(data_dir: &Path, mint_url: &str, amount_sats: u64) -> let mut all_target: Vec = send_denoms.clone(); all_target.extend(&change_denoms); - let swap_result = client.swap(&selected_proofs, &all_target).await?; + let swap_result = client + .swap_at_least(&selected_proofs, &all_target, amount_sats) + .await?; // Mark original proofs as spent wallet.mark_spent(&indices); @@ -1192,7 +1194,11 @@ pub async fn receive_token(data_dir: &Path, token_str: &str) -> Result { // Verify all mints in the token are accepted let accepted = load_accepted_mints(data_dir).await?; for mint_url in token.mint_urls() { - if !accepted.mints.iter().any(|m| m == mint_url) { + if !accepted + .mints + .iter() + .any(|m| m.trim_end_matches('/') == mint_url.trim_end_matches('/')) + { anyhow::bail!("Mint '{}' is not in accepted mints list", mint_url); } } @@ -1217,7 +1223,7 @@ pub async fn receive_token(data_dir: &Path, token_str: &str) -> Result { received_total += amount; } Err(e) => { - warn!("Failed to swap proofs from mint {}: {:#}", entry.mint, e); + warn!("Failed to swap proofs from mint {}: {}", entry.mint, e); all_already_redeemed &= e.is::(); last_reason = Some(e.to_string()); // Continue with other mints if any @@ -1298,22 +1304,10 @@ pub async fn verify_and_receive_payment( token_str: &str, required_sats: u64, ) -> Result { - // Handle legacy tokens + let token_str = token_str.trim(); + // Synthetic legacy balances are not cryptographic proof of payment. if token_str.starts_with("cashuSend_") { - let amount = token_str - .split('_') - .nth(1) - .and_then(|s| s.parse::().ok()) - .unwrap_or(0); - if amount < required_sats { - anyhow::bail!( - "Insufficient payment: {} sats, need {} sats", - amount, - required_sats - ); - } - let received = receive_legacy_token(data_dir, token_str).await?; - return Ok(received); + anyhow::bail!("Legacy ecash cannot authorize a paid download"); } // Fedimint notes (#3): a buyer whose balance is in Fedimint pays with notes @@ -1336,52 +1330,45 @@ pub async fn verify_and_receive_payment( // Parse and validate the token (cashuA or cashuB) let token = CashuToken::deserialize(token_str)?; - let total = token.total_amount(); - + if token.unit.as_deref().unwrap_or("sat") != "sat" { + anyhow::bail!("Payment must be denominated in sats"); + } + // A sale must redeem atomically at one mint. Otherwise a later mint + // failure can consume earlier inputs without delivering the purchase. + let entry = match token.token.as_slice() { + [entry] => entry, + _ => anyhow::bail!("Use a single-mint token for this payment"), + }; + let total = entry + .proofs + .iter() + .try_fold(0u64, |sum, p| sum.checked_add(p.amount)) + .ok_or_else(|| anyhow::anyhow!("Payment amount overflow"))?; if total < required_sats { - anyhow::bail!( - "Insufficient payment: {} sats, need {} sats", - total, - required_sats - ); + anyhow::bail!("Insufficient payment: {total} sats, need {required_sats} sats"); } - - // Verify mints are accepted let accepted = load_accepted_mints(data_dir).await?; - for mint_url in token.mint_urls() { - if !accepted.mints.iter().any(|m| m == mint_url) { - anyhow::bail!("Mint '{}' not accepted", mint_url); - } + if !accepted + .mints + .iter() + .any(|m| m.trim_end_matches('/') == entry.mint.trim_end_matches('/')) + { + anyhow::bail!("Mint is not in the seller's accepted mints list"); } - // Swap proofs at mint (this verifies they're unspent and gives us fresh proofs) + let client = mint_client(data_dir, &entry.mint).await?; + let result = client + .swap_at_least( + &entry.proofs, + &amount_to_denominations(total), + required_sats, + ) + .await?; + let received_total = result.new_proofs.iter().map(|p| p.amount).sum(); + // Load after the network call, so an unrelated wallet update during the + // swap is not overwritten with a pre-swap snapshot. let mut wallet = load_wallet(data_dir).await?; - let mut received_total = 0u64; - - for entry in &token.token { - let client = mint_client(data_dir, &entry.mint).await?; - let entry_total: u64 = entry.proofs.iter().map(|p| p.amount).sum(); - let target_amounts = amount_to_denominations(entry_total); - - match client.swap(&entry.proofs, &target_amounts).await { - Ok(result) => { - let amount: u64 = result.new_proofs.iter().map(|p| p.amount).sum(); - wallet.add_proofs(&entry.mint, result.new_proofs); - received_total += amount; - } - Err(e) => { - warn!("Payment verification failed at mint {}: {}", entry.mint, e); - } - } - } - - if received_total < required_sats { - anyhow::bail!( - "Payment verification failed: only {} of {} sats verified", - received_total, - required_sats - ); - } + wallet.add_proofs(entry.mint.trim_end_matches('/'), result.new_proofs); wallet.record_tx( TransactionType::Receive, @@ -2465,3 +2452,7 @@ mod tests { assert_eq!(w.mint_url, "https://mint.minibits.cash/Bitcoin"); } } + +#[cfg(test)] +#[path = "payment_tests.rs"] +mod payment_tests; diff --git a/core/archipelago/src/wallet/mint_client.rs b/core/archipelago/src/wallet/mint_client.rs index 91b65055..85ed629b 100644 --- a/core/archipelago/src/wallet/mint_client.rs +++ b/core/archipelago/src/wallet/mint_client.rs @@ -153,6 +153,20 @@ fn mint_error(op: &str, status: reqwest::StatusCode, body: &str) -> anyhow::Erro cause.context(describe_mint_error_body(status, body)) } +fn fee_adjusted_targets(requested: &[u64], mut available: u64) -> Vec { + let mut outputs = Vec::new(); + for &amount in requested { + if available >= amount { + outputs.push(amount); + available -= amount; + } else { + outputs.extend(amount_to_denominations(available)); + break; + } + } + outputs +} + /// HTTP client for a single Cashu mint. pub struct MintClient { url: String, @@ -512,30 +526,57 @@ impl MintClient { /// Swap proofs for new proofs of different denominations. /// This is how we "receive" a token — swap it for fresh proofs that only we know. pub async fn swap(&self, inputs: &[Proof], target_amounts: &[u64]) -> Result { - let keyset = self.get_active_sat_keyset().await?; + self.swap_at_least(inputs, target_amounts, 0).await + } - // cashuB tokens carry NUT-02 v2 keyset ids in their 8-byte short form, - // which the mint rejects (bare 422). Repair here, not at each caller: - // the paid-download seller path swapped directly and every Minibits - // payment failed once the mint rotated to a v2 keyset. - let repaired = self.resolve_truncated_keyset_ids(inputs).await; - let inputs: &[Proof] = &repaired; + /// Refuse a payment whose mint fees would leave the seller underpaid, + /// before consuming any input proofs. + pub async fn swap_at_least( + &self, + inputs: &[Proof], + target_amounts: &[u64], + minimum: u64, + ) -> Result { + // V4 tokens carry short keyset IDs. Every swap path (including paid + // files and streams) must expand these, not only wallet imports. + let resolved = self.resolve_truncated_keyset_ids(inputs).await?; + let inputs = resolved.as_slice(); + let keyset = self.get_active_sat_keyset().await?; // NUT-02: a mint may charge a per-input fee, and it rejects the swap // outright unless outputs == inputs - fee (`11005 Transaction inputs // should equal outputs less fee`). Applied here rather than at each // call site so send, receive and cross-mint swaps are all covered. // Fee-free mints (Minibits) compute 0 and are unaffected. - let inputs_total: u64 = inputs.iter().map(|p| p.amount).sum(); - let fee = match self.get_keysets().await { - Ok(ks) => super::cashu::swap_fee_for(inputs, &ks), - Err(e) => { - debug!("Could not read keyset fees ({e:#}) — assuming fee-free mint"); - 0 - } - }; + anyhow::ensure!(!inputs.is_empty(), "No input proofs to swap"); + let inputs_total = inputs + .iter() + .try_fold(0u64, |sum, p| sum.checked_add(p.amount)) + .context("Input amount overflow")?; + let keysets = self.get_keysets().await?; + let mut fee_ppk = 0u64; + for proof in inputs { + let input_keyset = keysets + .iter() + .find(|k| k.id == proof.id) + .context("The mint does not recognize an input keyset")?; + anyhow::ensure!( + input_keyset.unit == "sat", + "Input keyset is not denominated in sats" + ); + fee_ppk = fee_ppk + .checked_add(input_keyset.input_fee_ppk) + .context("Mint fee overflow")?; + } + let fee = fee_ppk.div_ceil(1000); let spendable = inputs_total.saturating_sub(fee); - let requested: u64 = target_amounts.iter().sum(); + if spendable < minimum { + anyhow::bail!("Payment would leave {spendable} sats after mint fees; need {minimum} sats. No proofs were redeemed."); + } + let requested = target_amounts + .iter() + .try_fold(0u64, |sum, amount| sum.checked_add(*amount)) + .context("Output amount overflow")?; let owned_targets: Vec; let target_amounts: &[u64] = if requested > spendable { if spendable == 0 { @@ -546,7 +587,10 @@ impl MintClient { debug!( "Reducing swap outputs {requested} -> {spendable} to cover a {fee} sat mint fee" ); - owned_targets = amount_to_denominations(spendable); + // Callers put payment outputs before change. Keep that prefix + // intact while fees reduce change; re-splitting the entire sum + // can omit a payment denomination after consuming the inputs. + owned_targets = fee_adjusted_targets(target_amounts, spendable); &owned_targets } else { target_amounts @@ -591,6 +635,9 @@ impl MintClient { let mut new_proofs = Vec::new(); for (sig, (secret, r, amount)) in signatures.iter().zip(blinding_data.iter()) { + if sig.amount != *amount || sig.id != keyset.id { + anyhow::bail!("Mint returned a swap signature for an unexpected amount or keyset"); + } let c_prime = sig.c_prime_as_pubkey()?; let mint_key = keyset.key_for_amount(*amount)?; let c = bdhke::unblind_signature(&c_prime, r, &mint_key)?; @@ -737,43 +784,35 @@ impl MintClient { /// Repair proofs whose keyset id is a truncated NUT-02 **v2** id. /// /// A v2 keyset id is 33 bytes (version byte `0x01` + 32-byte hash), but - /// wallets written against the original 8-byte format truncate it when - /// they build a token. The mint then reads the `0x01` version, expects 33 + /// compact V4 tokens carry an 8-byte short ID. The swap endpoint needs + /// the full ID restored from the mint's keyset list. The mint then reads the `0x01` version, expects 33 /// bytes, and rejects the swap — reported as /// `inputs[0].id: NUT02: ID length invalid` behind a bare 422 (seen with /// a Minibits-issued token, 2026-08-17). /// /// The id only names which keyset signed the proof, so restoring the full /// id the mint advertises is exactly what the sender meant. It is also - /// safe to attempt: an id that names the wrong keyset fails signature - /// verification at the mint and no coins move. Anything already valid, or - /// with no unambiguous match, is passed through untouched so the mint's - /// own error is what the operator sees. - async fn resolve_truncated_keyset_ids(&self, proofs: &[Proof]) -> Vec { + /// safe to attempt: the mint still verifies the proof signature. Unknown + /// or ambiguous short IDs are rejected before redemption. + async fn resolve_truncated_keyset_ids(&self, proofs: &[Proof]) -> Result> { let needs_repair = proofs.iter().any(|p| is_truncated_v2_keyset_id(&p.id)); if !needs_repair { - return proofs.to_vec(); + return Ok(proofs.to_vec()); } // The mint's own keyset list, in the reference implementation's shape // so its NUT-02 resolver can consume it directly. - let known = match self.get_cdk_keysets().await { - Ok(k) => k, - Err(e) => { - debug!("Could not list keysets to repair truncated keyset ids: {e:#}"); - return proofs.to_vec(); - } - }; + let known = self.get_cdk_keysets().await?; proofs .iter() .cloned() .map(|mut p| { - if let Some(full) = super::cashu::resolve_keyset_id(&p.id, &known) { - debug!("Expanded short keyset id {} to {} for swap", p.id, full); - p.id = full; + if is_truncated_v2_keyset_id(&p.id) { + p.id = super::cashu::resolve_keyset_id(&p.id, &known) + .context("The mint cannot resolve this short keyset ID unambiguously")?; } - p + Ok(p) }) .collect() } @@ -809,7 +848,7 @@ impl MintClient { let mut all_new_proofs = Vec::new(); for entry in &token.token { - if entry.mint != self.url { + if entry.mint.trim_end_matches('/') != self.url { debug!( "Skipping proofs from different mint {} (ours: {})", entry.mint, self.url @@ -869,136 +908,4 @@ mod tests { let client = MintClient::new("http://mint.example.com").unwrap(); assert_eq!(client.url(), "http://mint.example.com"); } - - /// A minimal mint on 127.0.0.1 answering `/v1/keys` and `/v1/keysets` - /// with one v2 keyset, and rejecting every `/v1/swap` as already spent. - /// Each swap request body is sent back on the returned channel. - async fn stub_mint( - full_id: &'static str, - ) -> (String, tokio::sync::mpsc::UnboundedReceiver) { - use tokio::io::{AsyncReadExt, AsyncWriteExt}; - let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - let addr = listener.local_addr().unwrap(); - let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); - tokio::spawn(async move { - loop { - let Ok((mut stream, _)) = listener.accept().await else { - return; - }; - let mut buf = Vec::new(); - let mut chunk = [0u8; 4096]; - // Read headers, then as much body as Content-Length says. - let (head, body) = loop { - let n = stream.read(&mut chunk).await.unwrap_or(0); - if n == 0 { - break (String::new(), Vec::new()); - } - buf.extend_from_slice(&chunk[..n]); - let Some(pos) = buf.windows(4).position(|w| w == b"\r\n\r\n") else { - continue; - }; - let head = String::from_utf8_lossy(&buf[..pos]).to_string(); - let len = head - .lines() - .find_map(|l| { - let (k, v) = l.split_once(':')?; - k.eq_ignore_ascii_case("content-length") - .then(|| v.trim().parse::().ok())? - }) - .unwrap_or(0); - while buf.len() < pos + 4 + len { - let n = stream.read(&mut chunk).await.unwrap_or(0); - if n == 0 { - break; - } - buf.extend_from_slice(&chunk[..n]); - } - break (head, buf[pos + 4..].to_vec()); - }; - let request_line = head.lines().next().unwrap_or_default().to_string(); - let (status, reply) = if request_line.starts_with("GET /v1/keys ") { - let key = "0279be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798"; - let keys: serde_json::Map = (0..16) - .map(|i| ((1u64 << i).to_string(), serde_json::json!(key))) - .collect(); - ( - "200 OK", - serde_json::json!({"keysets": [ - {"id": full_id, "unit": "sat", "active": true, "keys": keys} - ]}), - ) - } else if request_line.starts_with("GET /v1/keysets ") { - ( - "200 OK", - serde_json::json!({"keysets": [ - {"id": full_id, "unit": "sat", "active": true, "input_fee_ppk": 0} - ]}), - ) - } else if request_line.starts_with("POST /v1/swap ") { - let _ = tx.send(serde_json::from_slice(&body).unwrap_or_default()); - ( - "400 Bad Request", - serde_json::json!({"code": 11001, "detail": "Token Already Spent"}), - ) - } else { - ("404 Not Found", serde_json::json!({})) - }; - let reply = reply.to_string(); - let _ = stream - .write_all( - format!( - "HTTP/1.1 {status}\r\nContent-Type: application/json\r\n\ - Content-Length: {}\r\nConnection: close\r\n\r\n{reply}", - reply.len() - ) - .as_bytes(), - ) - .await; - } - }); - (format!("http://{addr}"), rx) - } - - fn proof_with_id(id: &str) -> Proof { - Proof { - amount: 8, - id: id.to_string(), - secret: "test-secret".to_string(), - c: "02".to_string() + &"11".repeat(32), - } - } - - /// Regression (2026-09-29): the paid-download seller called `swap` - /// directly with a cashuB token's short v2 keyset id and the mint - /// answered 422. `swap` itself must send the full id. - #[tokio::test] - async fn swap_expands_a_short_v2_keyset_id_before_calling_the_mint() { - const FULL: &str = "01fc0ec0e59cd6fa01b7a88f8cd77fce81fd1e64bca67d752e984992b7a3c3a821"; - let (url, mut swaps) = stub_mint(FULL).await; - let client = MintClient::new(&url).unwrap(); - - let Err(err) = client - .swap(&[proof_with_id("01fc0ec0e59cd6fa")], &[8]) - .await - else { - panic!("stub mint rejects every swap"); - }; - assert!(err.is::(), "unexpected error: {err:#}"); - - let body = swaps.recv().await.expect("swap reached the mint"); - assert_eq!(body["inputs"][0]["id"], FULL); - } - - #[tokio::test] - async fn swap_passes_complete_keyset_ids_through_unchanged() { - const FULL: &str = "01fc0ec0e59cd6fa01b7a88f8cd77fce81fd1e64bca67d752e984992b7a3c3a821"; - let (url, mut swaps) = stub_mint(FULL).await; - let client = MintClient::new(&url).unwrap(); - - for id in [FULL, "009a1f293253e41e"] { - let _ = client.swap(&[proof_with_id(id)], &[8]).await; - let body = swaps.recv().await.expect("swap reached the mint"); - assert_eq!(body["inputs"][0]["id"], id); - } - } } diff --git a/core/archipelago/src/wallet/payment_tests.rs b/core/archipelago/src/wallet/payment_tests.rs new file mode 100644 index 00000000..9fb40266 --- /dev/null +++ b/core/archipelago/src/wallet/payment_tests.rs @@ -0,0 +1,428 @@ +//! Real HTTP/curve-signature regressions for paid Cashu redemption. +use super::*; +use crate::wallet::{bdhke, cashu::Proof}; +use bitcoin::secp256k1::{PublicKey, Scalar, Secp256k1, SecretKey}; +use hyper::{ + service::{make_service_fn, service_fn}, + Body, Request, Response, Server, +}; +use serde_json::{json, Value}; +use std::{ + convert::Infallible, + sync::{Arc, Mutex}, +}; + +const ACTIVE: &str = "0011223344556677"; +const V2: &str = "011111111111111111111111111111111111111111111111111111111111111111"; + +struct Mint { + url: String, + requests: Arc>>, + task: tokio::task::JoinHandle<()>, + failure: Arc, +} +impl Drop for Mint { + fn drop(&mut self) { + self.task.abort(); + } +} + +fn signing_key() -> SecretKey { + SecretKey::from_slice(&[7; 32]).unwrap() +} +fn signed_point(point: PublicKey) -> String { + point + .mul_tweak(&Secp256k1::new(), &Scalar::from(signing_key())) + .unwrap() + .to_string() +} +fn proof(id: &str, amount: u64) -> Proof { + let secret = format!("test-{id}-{amount}"); + Proof { + amount, + id: id.into(), + c: signed_point(bdhke::hash_to_curve(secret.as_bytes()).unwrap()), + secret, + } +} +impl Mint { + async fn start(fee: u64, failure: Option) -> Self { + 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 requests = Arc::new(Mutex::new(Vec::new())); + let seen = requests.clone(); + let failure = Arc::new(std::sync::atomic::AtomicU16::new(failure.unwrap_or(0))); + let rejection = failure.clone(); + let spent = Arc::new(Mutex::new(std::collections::HashSet::::new())); + let service = make_service_fn(move |_| { + let seen = seen.clone(); + let rejection = rejection.clone(); + let spent = spent.clone(); + async move { + Ok::<_, Infallible>(service_fn(move |req: Request| { + let seen = seen.clone(); + let rejection = rejection.clone(); + let spent = spent.clone(); + async move { + let mut status = 200; + let body = match req.uri().path() { + "/v1/keysets" => json!({"keysets":[ + {"id": ACTIVE,"unit":"sat","active":true,"input_fee_ppk":fee}, + {"id": V2,"unit":"sat","active":false,"input_fee_ppk":fee} + ]}), + "/v1/keys" => { + let public = + PublicKey::from_secret_key(&Secp256k1::new(), &signing_key()) + .to_string(); + let keys: serde_json::Map = (0..16) + .map(|i| ((1u64 << i).to_string(), json!(public))) + .collect(); + json!({"keysets":[{"id": ACTIVE,"unit":"sat","keys":keys}]}) + } + "/v1/swap" => { + let body: Value = serde_json::from_slice( + &hyper::body::to_bytes(req.into_body()).await.unwrap(), + ) + .unwrap(); + seen.lock().unwrap().push(body.clone()); + let inputs = body["inputs"].as_array().unwrap(); + let outputs = body["outputs"].as_array().unwrap(); + let code = rejection.load(std::sync::atomic::Ordering::SeqCst); + if code != 0 { + status = code; + json!({"detail":"mock mint rejection"}) + } else if inputs.iter().any(|p| p["id"] != V2 && p["id"] != ACTIVE) + { + status = 422; + json!({"detail":[{"msg":"NUT02: ID length invalid"}]}) + } else if inputs.iter().any(|p| { + spent + .lock() + .unwrap() + .contains(p["secret"].as_str().unwrap()) + }) { + status = 400; + json!({"code":11001,"detail":"Token Already Spent"}) + } else { + let total: u64 = + inputs.iter().map(|p| p["amount"].as_u64().unwrap()).sum(); + let out: u64 = + outputs.iter().map(|p| p["amount"].as_u64().unwrap()).sum(); + assert_eq!( + out, + total - (inputs.len() as u64 * fee).div_ceil(1000) + ); + for p in inputs { + spent + .lock() + .unwrap() + .insert(p["secret"].as_str().unwrap().into()); + } + json!({"signatures":outputs.iter().map(|o| json!({ + "amount":o["amount"],"id":ACTIVE, + "C_":signed_point(o["B_"].as_str().unwrap().parse().unwrap()) + })).collect::>()}) + } + } + _ => { + status = 404; + json!({}) + } + }; + Ok::<_, Infallible>( + Response::builder() + .status(status) + .header("Content-Type", "application/json") + .body(Body::from(body.to_string())) + .unwrap(), + ) + } + })) + } + }); + let server = Server::from_tcp(listener).unwrap().serve(service); + let task = tokio::spawn(async move { + server.await.unwrap(); + }); + Self { + url, + requests, + task, + failure, + } + } + async fn wallet(&self) -> tempfile::TempDir { + let dir = tempfile::tempdir().unwrap(); + save_accepted_mints( + dir.path(), + &AcceptedMints { + mints: vec![format!("{}/", self.url)], + }, + ) + .await + .unwrap(); + dir + } +} + +#[tokio::test] +async fn paid_v4_inactive_v2_keyset_is_expanded_and_cryptographic_proofs_saved() { + let mint = Mint::start(0, None).await; + let dir = mint.wallet().await; + let token = CashuToken::new(&mint.url, vec![proof(V2, 64), proof(V2, 32), proof(V2, 4)]) + .serialize_v4() + .unwrap(); + let decoded = CashuToken::deserialize(&token).unwrap(); + assert_eq!( + decoded.token[0].proofs[0].id.len(), + 16, + "reproduce the short V4 ID" + ); + assert_eq!( + verify_and_receive_payment(dir.path(), &token, 100) + .await + .unwrap(), + 100 + ); + let wallet = load_wallet(dir.path()).await.unwrap(); + assert_eq!(wallet.balance(), 100); + for p in wallet.proofs { + assert_eq!( + p.proof.c, + signed_point(bdhke::hash_to_curve(p.proof.secret.as_bytes()).unwrap()) + ); + } + assert!(mint.requests.lock().unwrap()[0]["inputs"] + .as_array() + .unwrap() + .iter() + .all(|p| p["id"] == V2)); + assert!(verify_and_receive_payment(dir.path(), &token, 100) + .await + .is_err()); + assert_eq!(load_wallet(dir.path()).await.unwrap().balance(), 100); +} + +#[tokio::test] +async fn paid_v3_full_v2_and_v1_ids_work() { + for id in [V2, ACTIVE] { + let mint = Mint::start(0, None).await; + let dir = mint.wallet().await; + let token = CashuToken::new(&mint.url, vec![proof(id, 128)]) + .serialize() + .unwrap(); + assert_eq!( + verify_and_receive_payment(dir.path(), &token, 100) + .await + .unwrap(), + 128 + ); + } +} + +#[tokio::test] +async fn fees_cannot_consume_underpayment_and_allowed_fees_credit_actual_value() { + let mint = Mint::start(1000, None).await; + let dir = mint.wallet().await; + let token = CashuToken::new(&mint.url, vec![proof(V2, 128)]) + .serialize_v4() + .unwrap(); + assert!(verify_and_receive_payment(dir.path(), &token, 128) + .await + .unwrap_err() + .to_string() + .contains("after mint fees")); + assert!(mint.requests.lock().unwrap().is_empty()); + assert_eq!( + verify_and_receive_payment(dir.path(), &token, 127) + .await + .unwrap(), + 127 + ); + assert_eq!(load_wallet(dir.path()).await.unwrap().balance(), 127); +} + +#[tokio::test] +async fn rejected_mint_response_does_not_credit_wallet() { + for status in [200, 400, 422, 500, 503] { + let mint = Mint::start(0, Some(status)).await; + let dir = mint.wallet().await; + let token = CashuToken::new(&mint.url, vec![proof(V2, 128)]) + .serialize_v4() + .unwrap(); + assert!(verify_and_receive_payment(dir.path(), &token, 100) + .await + .is_err()); + assert_eq!(load_wallet(dir.path()).await.unwrap().balance(), 0); + } +} + +#[tokio::test] +async fn invalid_untrusted_multimint_and_underpaid_tokens_never_reach_swap() { + let mint = Mint::start(0, None).await; + let dir = mint.wallet().await; + let token = CashuToken::new(&mint.url, vec![proof(V2, 128)]); + let mut invalid = vec![ + "cashuSend_500_abc_1700000000".into(), + "cashuBinvalid".into(), + ]; + let mut wrong_unit = token.clone(); + wrong_unit.unit = Some("usd".into()); + invalid.push(wrong_unit.serialize().unwrap()); + let mut multi = token.clone(); + multi.token.push(token.token[0].clone()); + invalid.push(multi.serialize().unwrap()); + let mut untrusted = token.clone(); + untrusted.token[0].mint = "http://127.0.0.1:1".into(); + invalid.push(untrusted.serialize().unwrap()); + for id in ["00ffffffffffffff", "01ffffffffffffff"] { + invalid.push( + CashuToken::new(&mint.url, vec![proof(id, 128)]) + .serialize() + .unwrap(), + ); + } + for value in invalid { + assert!(verify_and_receive_payment(dir.path(), &value, 100) + .await + .is_err()); + } + assert!( + verify_and_receive_payment(dir.path(), &token.serialize().unwrap(), 129) + .await + .is_err() + ); + assert!(mint.requests.lock().unwrap().is_empty()); + assert_eq!(load_wallet(dir.path()).await.unwrap().balance(), 0); +} + +#[tokio::test] +async fn buyer_token_rejected_by_seller_can_be_refunded_without_balance_loss() { + let mint = Mint::start(0, Some(422)).await; + let buyer = mint.wallet().await; + let seller = mint.wallet().await; + let mut wallet = load_wallet(buyer.path()).await.unwrap(); + wallet.mint_url = mint.url.clone(); + wallet.add_proofs(&mint.url, vec![proof(V2, 64), proof(V2, 32), proof(V2, 4)]); + save_wallet(buyer.path(), &wallet).await.unwrap(); + let token = send_token(buyer.path(), 100).await.unwrap(); + assert_eq!(load_wallet(buyer.path()).await.unwrap().balance(), 0); + assert!(verify_and_receive_payment(seller.path(), &token, 100) + .await + .is_err()); + mint.failure.store(0, std::sync::atomic::Ordering::SeqCst); + assert_eq!(receive_token(buyer.path(), &token).await.unwrap(), 100); + assert_eq!(load_wallet(buyer.path()).await.unwrap().balance(), 100); + assert_eq!(load_wallet(seller.path()).await.unwrap().balance(), 0); + assert!(receive_token(buyer.path(), &token).await.is_err()); + assert_eq!(load_wallet(buyer.path()).await.unwrap().balance(), 100); +} + +#[tokio::test] +async fn unreachable_mint_does_not_credit_seller() { + let mint = Mint::start(0, None).await; + let dir = mint.wallet().await; + let token = CashuToken::new(&mint.url, vec![proof(V2, 128)]) + .serialize_v4() + .unwrap(); + mint.task.abort(); + tokio::task::yield_now().await; + assert!(verify_and_receive_payment(dir.path(), &token, 100) + .await + .is_err()); + assert_eq!(load_wallet(dir.path()).await.unwrap().balance(), 0); +} + +#[tokio::test] +async fn send_with_fees_preserves_payment_denominations_and_saves_change() { + // 128 inputs - 2 fee = 126. Splitting 126 as one sum omits 1, + // which is needed for a 65-sat payment, after consuming the inputs. + let mint = Mint::start(1000, None).await; + let buyer = mint.wallet().await; + let mut wallet = load_wallet(buyer.path()).await.unwrap(); + wallet.mint_url = mint.url.clone(); + let first = proof(V2, 64); + let mut second = first.clone(); + second.secret.push_str("-second"); + second.c = signed_point(bdhke::hash_to_curve(second.secret.as_bytes()).unwrap()); + wallet.add_proofs(&mint.url, vec![first, second]); + save_wallet(buyer.path(), &wallet).await.unwrap(); + let encoded = send_token(buyer.path(), 65).await.unwrap(); + assert_eq!( + CashuToken::deserialize(&encoded).unwrap().total_amount(), + 65 + ); + assert_eq!(load_wallet(buyer.path()).await.unwrap().balance(), 61); +} + +#[tokio::test] +async fn paid_file_gate_delivers_bytes_only_after_payment_and_does_not_charge_missing_files() { + use crate::content_server::{ + self, AccessControl, Availability, ContentCatalog, ContentItem, ServeResult, + }; + for (exists, accepts_cashu, price) in [ + (true, true, 100), + (true, false, 100), + (false, true, 100), + (true, true, 129), + ] { + let mint = Mint::start(0, None).await; + let seller = mint.wallet().await; + let item = ContentItem { + id: "paid-test".into(), + filename: "test.txt".into(), + mime_type: "text/plain".into(), + size_bytes: 5, + description: String::new(), + added_at: String::new(), + availability: Availability::AllPeers, + access: AccessControl::Paid { + price_sats: price, + accepted: vec![if accepts_cashu { "ecash" } else { "fedimint" }.into()], + }, + }; + content_server::save_catalog(seller.path(), &ContentCatalog { items: vec![item] }) + .await + .unwrap(); + if exists { + tokio::fs::create_dir_all(seller.path().join("content/files")) + .await + .unwrap(); + tokio::fs::write(seller.path().join("content/files/test.txt"), b"hello") + .await + .unwrap(); + } + let token = CashuToken::new(&mint.url, vec![proof(V2, 128)]) + .serialize_v4() + .unwrap(); + let result = content_server::serve_content( + seller.path(), + "paid-test", + Some(&token), + None, + None, + None, + false, + ) + .await + .unwrap(); + if exists && accepts_cashu && price <= 128 { + match result { + ServeResult::Ok(bytes, mime) => { + assert_eq!(bytes, b"hello"); + assert_eq!(mime, "text/plain"); + } + _ => panic!("paid content was not delivered"), + } + assert_eq!(load_wallet(seller.path()).await.unwrap().balance(), 128); + } else { + assert!(matches!( + result, + ServeResult::NotFound | ServeResult::PaymentRequired(_) + )); + assert_eq!(load_wallet(seller.path()).await.unwrap().balance(), 0); + assert!(mint.requests.lock().unwrap().is_empty()); + } + } +} diff --git a/docker/lnd-ui/index.html b/docker/lnd-ui/index.html index c5cf8da5..64dd2ad9 100644 --- a/docker/lnd-ui/index.html +++ b/docker/lnd-ui/index.html @@ -989,7 +989,7 @@ // ── State ─────────────────────────────────────────────────────── let unit = 'sats'; - let state = { info: null, channels: [], pending: null, peers: [], payments: [], invoices: [], txns: [], fees: null, graph: null }; + let state = { readiness: null, info: null, channels: [], pending: null, peers: [], payments: [], invoices: [], txns: [], fees: null, graph: null }; let peerSort = { col: 'peer', dir: 1 }; let activityFilter = 'all'; let logsLoaded = false; @@ -1142,9 +1142,19 @@ } async function refreshAll() { + if (state.refreshing) return; + state.refreshing = true; const icon = document.getElementById('refreshIcon'); if (icon) icon.classList.add('animate-spin-slow'); try { + state.readiness = await lndSafe('/archy-status', null); + if (state.readiness && state.readiness.state.startsWith('waiting_')) { + state.info = null; + state.onchainStale = true; + state.chanbalStale = true; + renderAll(); + return; + } const [info, channels, pending, peers, fees, graph, payments, invoices, txns] = await Promise.all([ lndSafe('/v1/getinfo', null), lndSafe('/v1/channels', { channels: [] }), @@ -1166,10 +1176,17 @@ state.invoices = (invoices && invoices.invoices) || []; state.txns = (txns && txns.transactions) || []; - // Balances are separate so one failing endpoint can't blank the rest. - state.onchain = await lndSafe('/v1/balance/blockchain', null); - state.chanbal = await lndSafe('/v1/balance/channels', null); + // Preserve known balances on outage; never decode an error as zero. + const [onchain, chanbal] = await Promise.all([ + lndSafe('/v1/balance/blockchain', null), + lndSafe('/v1/balance/channels', null), + ]); + state.onchainStale = !validBalance(onchain && (onchain.confirmed_balance ?? onchain.total_balance)); + state.chanbalStale = !validBalance(chanbal && (chanbal.local_balance?.sat ?? chanbal.balance)); + if (!state.onchainStale) state.onchain = onchain; + if (!state.chanbalStale) state.chanbal = chanbal; } finally { + state.refreshing = false; if (icon) icon.classList.remove('animate-spin-slow'); } renderAll(); @@ -1192,11 +1209,17 @@ const pill = document.getElementById('headerStatusPill'); const dot = document.getElementById('headerStatusDot'); - if (!g) { - setText('headerStatusText', 'Unreachable'); - pill.className = 'pill bad'; - dot.className = 'status-dot-sm bg-red'; - document.getElementById('syncCard').style.display = 'none'; + const waiting = state.readiness && state.readiness.state.startsWith('waiting_'); + if (!g || waiting) { + setText('headerStatusText', waiting ? state.readiness.message : 'Connecting to LND'); + pill.className = 'pill warn'; + dot.className = 'status-dot-sm bg-yellow'; + document.getElementById('syncCard').style.display = ''; + setText('syncSubtitle', waiting ? state.readiness.message + '. Lightning will become available automatically.' : 'Checking Lightning availability. Retrying automatically.'); + setText('syncBlockLabel', ''); + setText('syncPercent', ''); + document.getElementById('syncProgressBar').style.width = '0%'; + for (const id of ['syncChain', 'syncGraph', 'syncHeight', 'syncPeers']) setText(id, '—'); return; } @@ -1237,6 +1260,11 @@ } // ── Balances ──────────────────────────────────────────────────── + function validBalance(value) { + return (typeof value === 'number' || (typeof value === 'string' && /^\d+$/.test(value))) + && Number.isSafeInteger(Number(value)) && Number(value) >= 0; + } + function renderBalances() { const onchainConfirmed = num(state.onchain && (state.onchain.confirmed_balance ?? state.onchain.total_balance)); const onchainUnconfirmed = num(state.onchain && state.onchain.unconfirmed_balance); @@ -1253,22 +1281,23 @@ const haveOnchain = !!state.onchain; const haveChan = !!cb; - setBalance('balTotal', haveOnchain || haveChan ? onchainConfirmed + lnLocal : null); - setText('balTotalSub', haveOnchain || haveChan ? 'on-chain + lightning' : 'waiting for LND'); + setBalance('balTotal', haveOnchain && haveChan ? onchainConfirmed + lnLocal : null); + setText('balTotalSub', state.onchainStale || state.chanbalStale ? 'balance unavailable · last known values' : haveOnchain && haveChan ? 'on-chain + lightning' : 'waiting for LND'); setBalance('balLightning', haveChan ? lnLocal : null); - setText('balLightningSub', !haveChan ? 'waiting for LND' + setText('balLightningSub', !haveChan ? 'waiting for LND' : state.chanbalStale ? 'last known balance' : lnPending > 0 ? fmtAmount(lnPending) + ' pending open' : 'spendable over channels'); setBalance('balOnchain', haveOnchain ? onchainConfirmed : null); - setText('balOnchainSub', !haveOnchain ? 'waiting for LND' + setText('balOnchainSub', !haveOnchain ? 'waiting for LND' : state.onchainStale ? 'last known balance' : onchainUnconfirmed > 0 ? fmtAmount(onchainUnconfirmed) + ' unconfirmed' : 'confirmed'); - setText('liqLocal', fmtAmount(lnLocal)); - setText('liqRemote', fmtAmount(lnRemote)); + const liquidityReady = haveChan && !state.chanbalStale && !!state.info; + setText('liqLocal', liquidityReady ? fmtAmount(lnLocal) : '—'); + setText('liqRemote', liquidityReady ? fmtAmount(lnRemote) : '—'); const total = lnLocal + lnRemote; const localPct = total > 0 ? (lnLocal / total) * 100 : 50; - document.getElementById('liqBarLocal').style.width = localPct + '%'; - document.getElementById('liqBarRemote').style.width = (100 - localPct) + '%'; - setText('liqHint', total > 0 + document.getElementById('liqBarLocal').style.width = (liquidityReady ? localPct : 0) + '%'; + document.getElementById('liqBarRemote').style.width = (liquidityReady ? 100 - localPct : 0) + '%'; + setText('liqHint', !liquidityReady ? 'Channel capacity is unavailable while waiting for LND.' : total > 0 ? Math.round(localPct) + '% of your channel capacity is outbound (sendable).' : 'Open a channel to start sending and receiving over Lightning.'); } @@ -1284,6 +1313,15 @@ function renderSummary() { const g = state.info; + if (!g) { + for (const id of ['statPeers', 'statActiveChannels', 'statCapacity', 'statRoutingMonth', 'healthHeight', 'healthPending', 'chActive', 'chInactive', 'chPending', 'chCapacity']) setText(id, '—'); + for (const id of ['statChannelsSub', 'channelsLinkSub']) setText(id, 'Waiting for LND'); + for (const id of ['healthChain', 'healthGraph']) { + const pill = document.getElementById(id); + pill.textContent = '—'; pill.className = 'pill warn'; + } + return; + } const chans = state.channels; const active = chans.filter(c => c.active).length; const inactive = chans.length - active; @@ -1321,6 +1359,10 @@ function renderChannels() { const el = document.getElementById('channelList'); if (!el) return; + if (!state.info) { + el.innerHTML = '
Waiting for LND. Existing channels will appear when it is ready.
'; + return; + } const q = (document.getElementById('channelFilter').value || '').toLowerCase(); let list = state.channels.slice(); if (q) list = list.filter(c => String(c.remote_pubkey || '').toLowerCase().includes(q) || String(c.chan_id || '').includes(q)); diff --git a/docs/TODO.md b/docs/TODO.md index 1db93d91..6e5d53fb 100644 --- a/docs/TODO.md +++ b/docs/TODO.md @@ -3,14 +3,47 @@ Working backlog of forward-looking items not yet scoped into a dedicated plan doc. See [`ROADMAP.md`](ROADMAP.md) for the curated, public-facing direction. -## Blocking incident — before unrelated work +## Framework incident — closed with operator acceptance -- **OPEN: Framework LND startup / missing Receive address / false zero balance.** - User requires investigation and a verified fix on the actual node before later - unrelated work. Access is pending; a manual LND restart is only a workaround. +- **CLOSED WITH OPERATOR ACCEPTANCE (2026-09-30): Framework LND startup / + missing Receive address / false zero balance.** Startup, native balances, + Cashu address and source integration were verified; the operator accepted the + remaining display check and authorized release. See the incident record for evidence. See [incident evidence and closure criteria](incident-framework-lnd-startup.md) and the repository `AGENTS.md` session-start instructions. +## Next release after 1.8.21 — reported 2026-09-30 + +- [ ] **ThinkPad X250 kiosk: Bitcoin installation version selector is unreadable + and appears underneath the pruning information.** Operator reports white + styling with invisible text on the actual kiosk; the same flow works in remote + Brave. Reproduce on the X250's kiosk engine and record its version, display + scale and resolution. Inspect the native `