From 9ce04627ddc047d53cfe85e2dc16218ca5fa3631 Mon Sep 17 00:00:00 2001 From: archipelago Date: Tue, 6 Oct 2026 05:48:05 -0400 Subject: [PATCH] Stream purchased files into durable cache and avoid duplicate concurrent payments --- core/archipelago/src/api/rpc/content.rs | 303 +++++++++++------ core/archipelago/src/container/filebrowser.rs | 99 +++++- core/archipelago/src/content_owned.rs | 312 ++++++++++++++++-- docs/post-1.9.0-progress-20261006.md | 36 ++ .../__tests__/usePaidItemViewer.test.ts | 20 ++ neode-ui/src/composables/usePaidItemViewer.ts | 23 +- neode-ui/src/views/PeerFiles.vue | 40 +-- .../__tests__/PeerFilesLightning.test.ts | 29 ++ neode-ui/src/views/web5/Web5SharedContent.vue | 25 +- 9 files changed, 709 insertions(+), 178 deletions(-) diff --git a/core/archipelago/src/api/rpc/content.rs b/core/archipelago/src/api/rpc/content.rs index 61c54492..934aa8b1 100644 --- a/core/archipelago/src/api/rpc/content.rs +++ b/core/archipelago/src/api/rpc/content.rs @@ -43,6 +43,22 @@ async fn reclaim_spent_ecash(data_dir: &std::path::Path, token: &str, backend: & } } +async fn bounded_seller_error(mut response: reqwest::Response) -> String { + let mut bytes = Vec::new(); + let _ = tokio::time::timeout(std::time::Duration::from_secs(5), async { + while bytes.len() < 4096 { + match response.chunk().await { + Ok(Some(chunk)) => { + bytes.extend_from_slice(&chunk[..chunk.len().min(4096 - bytes.len())]) + } + _ => break, + } + } + }) + .await; + String::from_utf8_lossy(&bytes).into_owned() +} + /// 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 { @@ -80,6 +96,7 @@ async fn existing_paid_content( onion: &str, content_id: &str, filename: Option<&str>, + cache_only: bool, ) -> Result> { let owned = crate::content_owned::list_owned_checked(data_dir) .await @@ -93,35 +110,80 @@ async fn existing_paid_content( }) else { return Ok(None); }; - let (mime, bytes) = crate::content_owned::read_owned(data_dir, &item.onion, &item.content_id) - .await.context("This purchase is recorded, but its cached file is unavailable. No new payment was sent. Restore the cached file or contact the seller.")?; - let mut response = paid_content_response(&bytes, &mime, 0); + let mut response = cached_purchase_response(data_dir, &item.onion, &item.content_id, cache_only, 0).await + .context("This purchase is recorded, but its cached file is unavailable. No new payment was sent. Recover delivery without paying again.")?; response["already_owned"] = serde_json::json!(true); response["filename"] = serde_json::json!(item.filename); Ok(Some(response)) } -// Updated clients open the persisted file through the Range-capable HTTP -// endpoint. Avoid putting two base64 copies of a large video in a JSON reply. -// Keep older clients compatible until both sides have upgraded. -fn invoice_download_response(bytes: &[u8], mime: &str, cache_only: bool) -> serde_json::Value { - if cache_only { - serde_json::json!({ "owned": true, "mime_type": mime, "size_bytes": bytes.len() }) - } else { - paid_content_response(bytes, mime, 0) +async fn cached_purchase_response( + data_dir: &std::path::Path, + onion: &str, + content_id: &str, + cache_only: bool, + paid_sats: u64, +) -> Result { + use tokio::io::AsyncReadExt; + let (mime, file) = crate::content_owned::open_owned(data_dir, onion, content_id) + .await? + .context("Purchased content is not cached")?; + let size = file.metadata().await?.len(); + if cache_only || size > 16 * 1024 * 1024 { + return Ok( + serde_json::json!({"owned":true,"mime_type":mime,"size":size,"size_bytes":size,"paid_sats":paid_sats,"owned_content_id":content_id}), + ); } + let mut bytes = Vec::with_capacity(size as usize); + file.take(16 * 1024 * 1024 + 1) + .read_to_end(&mut bytes) + .await?; + anyhow::ensure!( + bytes.len() as u64 == size, + "Purchased file changed during reading" + ); + let mut result = paid_content_response(&bytes, &mime, paid_sats); + result["owned_content_id"] = serde_json::json!(content_id); + Ok(result) } -/// File purchases through an atomic no-clobber write in Files' own namespace. -async fn file_purchase_in_files( +async fn cache_peer_response( data_dir: &std::path::Path, + onion: &str, + content_id: &str, filename: &str, mime: &str, - bytes: &[u8], -) -> Result { - let folder = if mime.starts_with("image/") || mime.starts_with("video/") { + paid_sats: u64, + backend: &str, + response: reqwest::Response, +) -> Result { + let expected = response.content_length(); + crate::content_owned::record_purchase_stream( + data_dir, + crate::content_owned::OwnedItem { + onion: onion.into(), + content_id: content_id.into(), + filename: filename.into(), + mime_type: mime.into(), + size_bytes: expected.unwrap_or(0), + paid_sats, + ecash_backend: backend.into(), + purchased_at: chrono::Utc::now().to_rfc3339(), + download_complete: false, + }, + response.bytes_stream(), + expected, + ) + .await +} + +async fn file_cached_purchase_in_files( + data_dir: &std::path::Path, + item: &crate::content_owned::OwnedItem, +) -> Result<()> { + let folder = if item.mime_type.starts_with("image/") || item.mime_type.starts_with("video/") { "Photos" - } else if mime.starts_with("audio/") { + } else if item.mime_type.starts_with("audio/") { "Music" } else { "Documents" @@ -131,19 +193,16 @@ async fn file_purchase_in_files( tokio::fs::metadata(&root).await?.is_dir(), "Files storage is unavailable" ); - let name = std::path::Path::new(filename) + let name = std::path::Path::new(&item.filename) .file_name() .and_then(|n| n.to_str()) .filter(|n| !n.is_empty()) .unwrap_or("download"); - let path = - crate::container::filebrowser::save_new_file(&root.join(folder), name, bytes).await?; - Ok(format!( - "{folder}/{}", - path.file_name() - .and_then(|n| n.to_str()) - .context("Invalid Files name")? - )) + let (_, file) = crate::content_owned::open_owned(data_dir, &item.onion, &item.content_id) + .await? + .context("Purchase unavailable")?; + crate::container::filebrowser::save_new_file_from(&root.join(folder), name, file).await?; + Ok(()) } impl RpcHandler { @@ -540,6 +599,9 @@ impl RpcHandler { return Err(anyhow::anyhow!("Invalid v3 onion address")); } + crate::content_owned::validate_identity(onion, content_id)?; + let _purchase_lock = crate::content_owned::lock_seller_purchases(onion).await; + // NEVER pay twice for content we already own (2026-07-22: a file // shared twice produced two catalog ids for the same bytes and the // buyer paid both). Guard BEFORE any ecash is minted, matching both @@ -551,6 +613,10 @@ impl RpcHandler { onion, content_id, params.get("filename").and_then(|v| v.as_str()), + params + .get("cache_only") + .and_then(|v| v.as_bool()) + .unwrap_or(false), ) .await? { @@ -668,12 +734,11 @@ impl RpcHandler { if response.status() == reqwest::StatusCode::PAYMENT_REQUIRED { // 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(); + drop(response); tracing::warn!( - "paid download: seller {onion} rejected {used_backend} payment of {price_sats} sats: {body}" + "paid download: seller rejected {used_backend} payment of {price_sats} sats" ); - // Seller couldn't redeem the token — reclaim it so the buyer keeps - // their funds (the spent-but-unredeemed-notes case the user hit). + // Reclaim only proofs the mint still considers unspent. let refund = reclaim_spent_ecash(&self.config.data_dir, &token_str, used_backend).await; return Ok(serde_json::json!({ "error": format!("The seller could not verify the payment. {refund}") @@ -682,8 +747,8 @@ impl RpcHandler { if !response.status().is_success() { let status = response.status(); - let body = response.text().await.unwrap_or_default(); - tracing::warn!("paid download: seller {onion} returned {status}: {body}"); + let body = bounded_seller_error(response).await; + tracing::warn!("paid download: seller {onion} returned {status}"); let refund = reclaim_spent_ecash(&self.config.data_dir, &token_str, used_backend).await; return Ok(serde_json::json!({ "error": format!("{} {refund}", seller_error_message(status, &body)) @@ -700,59 +765,26 @@ impl RpcHandler { .filter(|s| !s.is_empty()) .unwrap_or_else(|| "application/octet-stream".to_string()); - 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 - // it unblurred and views it in-app from this cache — no re-payment and no - // reliance on a browser download (which silently fails on the mobile - // companion, the original "paid but never unlocked" report). Best-effort: - // a cache-write failure must not fail an already-paid download. let filename = params .get("filename") .and_then(|v| v.as_str()) - .unwrap_or(content_id) - .to_string(); - let purchased_at = chrono::Utc::now().to_rfc3339(); - if let Err(e) = crate::content_owned::record_purchase( + .unwrap_or(content_id); + let item = cache_peer_response(&self.config.data_dir,onion,content_id,filename,&mime_type,price_sats,used_backend,response) + .await.context("Paid file delivery could not be saved. Do not send another payment; recover this purchase first")?; + if let Err(error) = file_cached_purchase_in_files(&self.config.data_dir, &item).await { + tracing::warn!("Purchase cached; optional Files copy failed: {error:#}"); + } + let mut result = cached_purchase_response( &self.config.data_dir, onion, content_id, - &filename, - &mime_type, - &bytes, + params + .get("cache_only") + .and_then(|v| v.as_bool()) + .unwrap_or(false), price_sats, - used_backend, - &purchased_at, ) - .await - { - tracing::warn!("paid download: failed to cache purchased content (non-fatal): {e:#}"); - } - - // The durable purchased-content cache above is primary. A Files copy - // remains optional: a stopped FileBrowser must not undo a paid download. - let filed = - file_purchase_in_files(&self.config.data_dir, &filename, &mime_type, &bytes).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}" - ), - } - - tracing::info!("paid download: received {} bytes from {onion} (paid {price_sats} sats via {used_backend})", bytes.len()); - let mut result = paid_content_response(&bytes, &mime_type, price_sats); + .await?; result["ecash_backend"] = serde_json::json!(used_backend); Ok(result) } @@ -900,14 +932,27 @@ impl RpcHandler { return Err(anyhow::anyhow!("Invalid payment_hash")); } + crate::content_owned::validate_identity(onion, content_id)?; + let _purchase_lock = crate::content_owned::lock_seller_purchases(onion).await; let cache_only = params .get("cache_only") .and_then(|v| v.as_bool()) .unwrap_or(false); - if let Some((mime, bytes)) = - crate::content_owned::read_owned(&self.config.data_dir, onion, content_id).await + if crate::content_owned::list_owned_checked(&self.config.data_dir) + .await? + .iter() + .any(|item| { + item.onion == onion && item.content_id == content_id && item.download_complete + }) { - return Ok(invoice_download_response(&bytes, &mime, cache_only)); + return cached_purchase_response( + &self.config.data_dir, + onion, + content_id, + cache_only, + 0, + ) + .await; } // Older sellers only mark settlement during status polling. Always // perform that handshake before requesting bytes; retries never pay. @@ -976,36 +1021,29 @@ impl RpcHandler { .next() .unwrap_or("application/octet-stream") .to_string(); - let bytes = response - .bytes() - .await - .context("Paid file transfer interrupted; retry the download without paying again")?; let filename = params .get("filename") .and_then(|v| v.as_str()) .unwrap_or(content_id); - crate::content_owned::record_purchase( + let item = cache_peer_response( &self.config.data_dir, onion, content_id, filename, &mime, - &bytes, params .get("price_sats") .and_then(|v| v.as_u64()) .unwrap_or(0), "lightning", - &chrono::Utc::now().to_rfc3339(), + response, ) .await - .context("Paid file could not be saved; retry the download without paying again")?; - if let Err(error) = - file_purchase_in_files(&self.config.data_dir, filename, &mime, &bytes).await - { + .context("Paid file could not be saved; retry delivery without paying again")?; + if let Err(error) = file_cached_purchase_in_files(&self.config.data_dir, &item).await { tracing::warn!("Lightning purchase cached; optional Files copy failed: {error:#}"); } - Ok(invoice_download_response(&bytes, &mime, cache_only)) + cached_purchase_response(&self.config.data_dir, onion, content_id, cache_only, 0).await } /// Buyer side (#46): ask the seller for a fresh on-chain address to pay. @@ -1462,20 +1500,18 @@ impl RpcHandler { .and_then(|v| v.as_str()) .ok_or_else(|| anyhow::anyhow!("Missing content_id"))?; - match crate::content_owned::read_owned(&self.config.data_dir, onion, content_id).await { - Some((mime_type, bytes)) => { - use base64::Engine; - let encoded = base64::engine::general_purpose::STANDARD.encode(&bytes); - Ok(serde_json::json!({ - "data": encoded, - "size": bytes.len(), - "mime_type": mime_type, - })) - } - None => Ok(serde_json::json!({ - "error": "You don't own this item yet, or its cached copy is missing." - })), - } + existing_paid_content( + &self.config.data_dir, + onion, + content_id, + params.get("filename").and_then(|v| v.as_str()), + params + .get("cache_only") + .and_then(|v| v.as_bool()) + .unwrap_or(false), + ) + .await? + .context("Purchased content is not cached") } } @@ -1486,15 +1522,62 @@ mod tests; #[cfg(test)] mod invoice_delivery_response_tests { use super::*; - #[test] - fn cached_delivery_avoids_base64_but_keeps_old_clients_compatible() { - let cached = invoice_download_response(b"paid bytes", "video/mp4", true); + #[tokio::test] + async fn cached_delivery_avoids_base64_and_legacy_small_reads_remain_compatible() { + let dir = tempfile::tempdir().unwrap(); + crate::content_owned::record_purchase( + dir.path(), + "seller.onion", + "film", + "film", + "video/mp4", + b"paid bytes", + 1, + "cashu", + "now", + ) + .await + .unwrap(); + let cached = cached_purchase_response(dir.path(), "seller.onion", "film", true, 0) + .await + .unwrap(); assert_eq!(cached["owned"], true); assert_eq!(cached["size_bytes"], 10); assert!(cached.get("data").is_none()); assert!(cached.get("data_base64").is_none()); - let legacy = invoice_download_response(b"paid bytes", "video/mp4", false); + let alias = existing_paid_content( + dir.path(), + "seller.onion", + "new-catalog-id", + Some("film"), + true, + ) + .await + .unwrap() + .unwrap(); + assert_eq!(alias["owned_content_id"], "film"); + let legacy = cached_purchase_response(dir.path(), "seller.onion", "film", false, 0) + .await + .unwrap(); assert_eq!(legacy["data"], "cGFpZCBieXRlcw=="); assert_eq!(legacy["data"], legacy["data_base64"]); + let mut entry = crate::content_owned::list_owned_checked(dir.path()) + .await + .unwrap() + .remove(0); + entry.content_id = "incomplete".into(); + let stream = futures_util::stream::iter([Ok::<_, std::io::Error>( + bytes::Bytes::from_static(b"part"), + )]); + assert!( + crate::content_owned::record_purchase_stream(dir.path(), entry, stream, Some(10)) + .await + .is_err() + ); + assert!( + existing_paid_content(dir.path(), "seller.onion", "incomplete", None, true) + .await + .is_err() + ); } } diff --git a/core/archipelago/src/container/filebrowser.rs b/core/archipelago/src/container/filebrowser.rs index 95f1d7f2..6ca3d9b1 100644 --- a/core/archipelago/src/container/filebrowser.rs +++ b/core/archipelago/src/container/filebrowser.rs @@ -212,6 +212,27 @@ pub async fn save_new_file(dir: &Path, name: &str, bytes: &[u8]) -> Result Result { + use tokio::io::AsyncSeekExt; + validate_filename(name)?; + match fs::symlink_metadata(dir).await { + Ok(meta) => anyhow::ensure!(meta.is_dir(), "Files destination is not a directory"), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => {} + Err(error) => return Err(error.into()), + } + source.seek(std::io::SeekFrom::Start(0)).await?; + match write_direct_stream(dir, name, &mut source).await { + Ok(path) => Ok(path), + Err(error) if error.kind() == std::io::ErrorKind::PermissionDenied => { + source.seek(std::io::SeekFrom::Start(0)).await?; + write_via_userns_stream(dir.to_owned(), name.to_owned(), source).await + } + Err(error) => Err(error).context("Saving purchased file"), + } +} + fn validate_filename(name: &str) -> Result<()> { anyhow::ensure!( !name.is_empty() @@ -291,8 +312,15 @@ impl Drop for PendingFile { } async fn write_direct(dir: &Path, name: &str, bytes: &[u8]) -> std::io::Result { + write_direct_stream(dir, name, &mut &bytes[..]).await +} + +async fn write_direct_stream( + dir: &Path, + name: &str, + source: &mut R, +) -> std::io::Result { use std::os::unix::fs::PermissionsExt; - use tokio::io::AsyncWriteExt; fs::create_dir_all(dir).await?; let temp_path = dir.join(format!(".archy-saving-{}", uuid::Uuid::new_v4())); let mut file = fs::OpenOptions::new() @@ -302,7 +330,7 @@ async fn write_direct(dir: &Path, name: &str, bytes: &[u8]) -> std::io::Result

) -> Result< .context("Files namespace writer timed out")? } +async fn write_via_userns_stream( + dir: PathBuf, + name: String, + mut source: fs::File, +) -> Result { + let expected = source.metadata().await?.len(); + let mut child = tokio::process::Command::new("podman") + .args(["unshare", "sh", "-c", WRITE_VIA_USERNS, "sh"]) + .arg(&dir) + .arg(&name) + .arg(expected.to_string()) + .kill_on_drop(true) + .stdin(std::process::Stdio::piped()) + .stdout(std::process::Stdio::piped()) + .stderr(std::process::Stdio::null()) + .spawn() + .context("Starting Files namespace writer")?; + let mut stdin = child.stdin.take().context("Files writer stdin missing")?; + tokio::time::timeout(std::time::Duration::from_secs(900), async { + let count = tokio::io::copy(&mut source, &mut stdin).await?; + drop(stdin); + let output = child.wait_with_output().await?; + anyhow::ensure!( + count == expected && output.status.success(), + "Files copy failed; purchased cache is retained" + ); + let chosen = String::from_utf8(output.stdout).context("Invalid Files response")?; + validate_filename(&chosen)?; + anyhow::ensure!( + (1..=100).any(|n| numbered_name(&name, n) == chosen), + "Unexpected Files destination" + ); + Ok(dir.join(chosen)) + }) + .await + .context("Files namespace writer timed out")? +} + #[cfg(test)] mod tests { use super::*; + #[tokio::test] + async fn streamed_purchase_preserves_existing_file_and_exact_large_copy() { + let dir = tempfile::tempdir().unwrap(); + fs::write(dir.path().join("film.mp4"), b"keep") + .await + .unwrap(); + let source = dir.path().join("source"); + let bytes = vec![17; 2 * 1024 * 1024]; + fs::write(&source, &bytes).await.unwrap(); + let path = save_new_file_from( + dir.path(), + "film.mp4", + fs::File::open(&source).await.unwrap(), + ) + .await + .unwrap(); + assert_eq!(path.file_name().unwrap(), "film (2).mp4"); + assert_eq!(fs::read(path).await.unwrap(), bytes); + assert_eq!( + fs::read(dir.path().join("film.mp4")).await.unwrap(), + b"keep" + ); + assert!(!std::fs::read_dir(dir.path()).unwrap().any(|entry| entry + .unwrap() + .file_name() + .to_string_lossy() + .starts_with(".archy-saving"))); + } + #[tokio::test] async fn cloud_credentials_use_unique_record_and_never_default_password() { let dir = tempfile::tempdir().unwrap(); diff --git a/core/archipelago/src/content_owned.rs b/core/archipelago/src/content_owned.rs index be7c4846..d5113166 100644 --- a/core/archipelago/src/content_owned.rs +++ b/core/archipelago/src/content_owned.rs @@ -16,6 +16,29 @@ use tokio::{fs, io::AsyncWriteExt, sync::Mutex}; static PURCHASE_WRITES: Mutex<()> = Mutex::const_new(()); +// Serialize purchases per seller so duplicate UI requests cannot both pass the +// ownership check and spend. Weak entries are pruned between acquisitions. +pub async fn lock_seller_purchases(onion: &str) -> tokio::sync::OwnedMutexGuard<()> { + use std::sync::{Arc, OnceLock, Weak}; + static LOCKS: OnceLock>>>> = + OnceLock::new(); + let lock = { + let mut locks = LOCKS + .get_or_init(Default::default) + .lock() + .unwrap_or_else(|e| e.into_inner()); + locks.retain(|_, value| value.strong_count() > 0); + if let Some(lock) = locks.get(onion).and_then(Weak::upgrade) { + lock + } else { + let lock = Arc::new(Mutex::new(())); + locks.insert(onion.into(), Arc::downgrade(&lock)); + lock + } + }; + lock.lock_owned().await +} + const OWNED_DIR: &str = "purchased-content"; const OWNED_INDEX: &str = "owned.json"; @@ -32,6 +55,13 @@ pub struct OwnedItem { pub ecash_backend: String, /// RFC3339 timestamp; best-effort, empty if the clock was unavailable. pub purchased_at: String, + /// False means payment succeeded but delivery still needs recovery. + #[serde(default = "completed")] + pub download_complete: bool, +} + +fn completed() -> bool { + true } #[derive(Debug, Default, Serialize, Deserialize)] @@ -81,27 +111,49 @@ async fn load_index(data_dir: &Path) -> OwnedIndex { load_index_checked(data_dir).await.unwrap_or_default() } +struct PendingFile(PathBuf); +impl Drop for PendingFile { + fn drop(&mut self) { + let _ = std::fs::remove_file(&self.0); + } +} + +pub fn validate_identity(onion: &str, content_id: &str) -> Result<()> { + anyhow::ensure!( + !onion.is_empty() + && !content_id.is_empty() + && onion != "." + && content_id != "." + && !onion.contains("..") + && !content_id.contains("..") + && sanitize(onion) == onion + && sanitize(content_id) == content_id, + "Invalid purchase path" + ); + Ok(()) +} + async fn atomic_write(path: &Path, bytes: &[u8]) -> Result<()> { let parent = path.parent().context("Purchase path has no parent")?; fs::create_dir_all(parent).await?; - let temp = parent.join(format!(".purchase-{}.tmp", uuid::Uuid::new_v4())); + let temp = PendingFile(parent.join(format!(".purchase-{}.tmp", uuid::Uuid::new_v4()))); let result = async { let mut file = fs::OpenOptions::new() .write(true) .create_new(true) .mode(0o600) - .open(&temp) + .open(&temp.0) .await?; file.write_all(bytes).await?; file.sync_all().await?; drop(file); - fs::rename(&temp, path).await?; + fs::rename(&temp.0, path).await?; fs::File::open(parent).await?.sync_all().await?; Ok::<_, anyhow::Error>(()) } .await; if result.is_err() { - let _ = fs::remove_file(&temp).await; + let _ = fs::remove_file(&temp.0).await; } result } @@ -131,6 +183,7 @@ pub async fn record_purchase( ecash_backend: &str, purchased_at: &str, ) -> Result<()> { + validate_identity(onion, content_id)?; // Read-modify-write must be one serialized transaction. Never replace a // damaged index with an empty one, and never expose partially written bytes. let _lock = PURCHASE_WRITES.lock().await; @@ -149,6 +202,7 @@ pub async fn record_purchase( paid_sats, ecash_backend: ecash_backend.to_string(), purchased_at: purchased_at.to_string(), + download_complete: true, }; if let Some(existing) = index .items @@ -162,6 +216,110 @@ pub async fn record_purchase( save_index(data_dir, &index).await } +fn require_cache_space(file: &fs::File, additional: u64) -> Result<()> { + use std::os::fd::AsRawFd; + let mut value = std::mem::MaybeUninit::::uninit(); + if unsafe { libc::fstatvfs(file.as_raw_fd(), value.as_mut_ptr()) } != 0 { + return Err(std::io::Error::last_os_error()).context("Checking purchase storage"); + } + let value = unsafe { value.assume_init() }; + let available = (value.f_bavail as u64).saturating_mul(value.f_frsize as u64); + anyhow::ensure!( + additional + .checked_add(256 * 1024 * 1024) + .is_some_and(|n| n <= available), + "Insufficient purchase storage; payment remains recorded for recovery" + ); + Ok(()) +} + +/// Save a peer response without holding its body in memory. A durable incomplete +/// ownership record precedes body consumption, so interrupted delivery cannot be +/// mistaken for permission to send another payment. Invoice retries may complete it. +pub async fn record_purchase_stream( + data_dir: &Path, + mut entry: OwnedItem, + mut stream: S, + expected: Option, +) -> Result +where + S: futures_util::Stream> + Unpin, + E: std::error::Error + Send + Sync + 'static, +{ + use futures_util::StreamExt; + validate_identity(&entry.onion, &entry.content_id)?; + let path = bytes_path(data_dir, &entry.onion, &entry.content_id); + let parent = path.parent().context("Purchase has no parent")?; + fs::create_dir_all(parent).await?; + let temporary = PendingFile(parent.join(format!(".purchase-{}.tmp", uuid::Uuid::new_v4()))); + let mut file = fs::OpenOptions::new() + .write(true) + .create_new(true) + .mode(0o600) + .open(&temporary.0) + .await?; + entry.download_complete = false; + entry.size_bytes = expected.unwrap_or(0); + { + let _lock = PURCHASE_WRITES.lock().await; + let mut index = load_index_checked(data_dir).await?; + if let Some(old) = index + .items + .iter_mut() + .find(|i| i.onion == entry.onion && i.content_id == entry.content_id) + { + anyhow::ensure!( + !old.download_complete, + "Purchase is already cached; use the owned copy" + ); + *old = entry.clone(); + } else { + index.items.push(entry.clone()); + } + save_index(data_dir, &index).await?; + } + require_cache_space(&file, expected.unwrap_or(0))?; + let mut received = 0u64; + while let Some(chunk) = stream.next().await { + let chunk = + chunk.context("Purchased file transfer interrupted; no new payment should be sent")?; + received = received + .checked_add(chunk.len() as u64) + .context("Content size overflow")?; + anyhow::ensure!( + expected.is_none_or(|n| received <= n), + "Seller exceeded the declared content length" + ); + require_cache_space(&file, chunk.len() as u64)?; + for part in chunk.chunks(64 * 1024) { + file.write_all(part).await?; + } + } + anyhow::ensure!( + expected.is_none_or(|n| received == n), + "Purchased file transfer ended early" + ); + file.sync_all().await?; + drop(file); + let _lock = PURCHASE_WRITES.lock().await; + let mut index = load_index_checked(data_dir).await?; + fs::rename(&temporary.0, &path).await?; + fs::File::open(parent).await?.sync_all().await?; + entry.size_bytes = received; + entry.download_complete = true; + if let Some(old) = index + .items + .iter_mut() + .find(|i| i.onion == entry.onion && i.content_id == entry.content_id) + { + *old = entry.clone(); + } else { + anyhow::bail!("Purchase ownership record disappeared; no new payment should be sent"); + } + save_index(data_dir, &index).await?; + Ok(entry) +} + /// Payment decisions must not interpret an unreadable index as no purchases. pub async fn list_owned_checked(data_dir: &Path) -> Result> { Ok(load_index_checked(data_dir).await?.items) @@ -197,17 +355,10 @@ pub async fn open_owned( else { return Ok(None); }; - // Reject path components even when called outside the HTTP route. + validate_identity(onion, content_id)?; anyhow::ensure!( - !onion.is_empty() - && !content_id.is_empty() - && onion != "." - && content_id != "." - && !onion.contains("..") - && !content_id.contains("..") - && sanitize(onion) == onion - && sanitize(content_id) == content_id, - "Invalid purchase path" + item.download_complete, + "Payment is recorded, but delivery is incomplete. Retry delivery without paying again." ); let file = fs::OpenOptions::new() .read(true) @@ -228,22 +379,135 @@ pub async fn read_owned( onion: &str, content_id: &str, ) -> Option<(String, Vec)> { - let bytes = fs::read(bytes_path(data_dir, onion, content_id)) - .await - .ok()?; - let mime = load_index(data_dir) - .await - .items - .into_iter() - .find(|i| i.onion == onion && i.content_id == content_id) - .map(|i| i.mime_type) - .unwrap_or_else(|| "application/octet-stream".to_string()); + use tokio::io::AsyncReadExt; + let (mime, mut file) = open_owned(data_dir, onion, content_id).await.ok()??; + let mut bytes = Vec::new(); + file.read_to_end(&mut bytes).await.ok()?; Some((mime, bytes)) } #[cfg(test)] mod tests { use super::*; + fn entry() -> OwnedItem { + OwnedItem { + onion: "seller.onion".into(), + content_id: "film".into(), + filename: "film.mp4".into(), + mime_type: "video/mp4".into(), + size_bytes: 0, + paid_sats: 1, + ecash_backend: "lightning".into(), + purchased_at: "now".into(), + download_complete: false, + } + } + #[tokio::test] + async fn stream_records_incomplete_delivery_then_recovers_without_losing_ownership() { + let dir = tempfile::tempdir().unwrap(); + let short = futures_util::stream::iter([Ok::<_, std::io::Error>( + bytes::Bytes::from_static(b"part"), + )]); + assert!(record_purchase_stream(dir.path(), entry(), short, Some(10)) + .await + .is_err()); + let pending = list_owned_checked(dir.path()).await.unwrap(); + assert_eq!(pending.len(), 1); + assert!(!pending[0].download_complete); + assert!(open_owned(dir.path(), "seller.onion", "film") + .await + .is_err()); + assert_eq!( + std::fs::read_dir(owned_root(dir.path()).join("seller.onion")) + .unwrap() + .count(), + 0 + ); + let chunks = futures_util::stream::iter( + (0..64).map(|_| Ok::<_, std::io::Error>(bytes::Bytes::from(vec![42; 65536]))), + ); + let item = record_purchase_stream(dir.path(), entry(), chunks, Some(4 * 1024 * 1024)) + .await + .unwrap(); + assert!(item.download_complete); + let (_, bytes) = read_owned(dir.path(), "seller.onion", "film") + .await + .unwrap(); + assert_eq!(bytes.len(), 4 * 1024 * 1024); + assert!(bytes.iter().all(|b| *b == 42)); + assert_eq!(list_owned_checked(dir.path()).await.unwrap().len(), 1); + } + #[tokio::test] + async fn cancelled_stream_cleans_partial_bytes_but_keeps_payment_record() { + let dir = tempfile::tempdir().unwrap(); + let started = std::sync::Arc::new(tokio::sync::Notify::new()); + let notify = started.clone(); + let root = dir.path().to_path_buf(); + let task = tokio::spawn(async move { + let stream = Box::pin(futures_util::stream::once(async move { + notify.notify_one(); + std::future::pending::>().await + })); + record_purchase_stream(&root, entry(), stream, Some(10)).await + }); + tokio::time::timeout(std::time::Duration::from_secs(10), started.notified()) + .await + .unwrap(); + task.abort(); + assert!(task.await.unwrap_err().is_cancelled()); + assert!(!list_owned_checked(dir.path()).await.unwrap()[0].download_complete); + assert_eq!( + std::fs::read_dir(owned_root(dir.path()).join("seller.onion")) + .unwrap() + .count(), + 0 + ); + } + #[tokio::test] + async fn unsafe_purchase_paths_are_rejected_before_writing() { + let dir = tempfile::tempdir().unwrap(); + for id in ["..", "../private", "/absolute", "a/b"] { + assert!(record_purchase( + dir.path(), + "seller.onion", + id, + "x", + "x", + b"x", + 1, + "cashu", + "now" + ) + .await + .is_err()); + } + assert!(!owned_root(dir.path()).exists()); + } + #[tokio::test] + async fn seller_purchase_lock_blocks_duplicates_but_not_another_seller() { + let first = lock_seller_purchases("one.onion").await; + assert!(tokio::time::timeout( + std::time::Duration::from_millis(10), + lock_seller_purchases("one.onion") + ) + .await + .is_err()); + let other = tokio::time::timeout( + std::time::Duration::from_secs(1), + lock_seller_purchases("two.onion"), + ) + .await + .unwrap(); + drop(first); + let _next = tokio::time::timeout( + std::time::Duration::from_secs(1), + lock_seller_purchases("one.onion"), + ) + .await + .unwrap(); + drop(other); + } + #[tokio::test] async fn owned_stream_preserves_ownership_on_missing_corrupt_or_symlinked_bytes() { let dir = tempfile::tempdir().unwrap(); diff --git a/docs/post-1.9.0-progress-20261006.md b/docs/post-1.9.0-progress-20261006.md index 673eb40b..20d4021e 100644 --- a/docs/post-1.9.0-progress-20261006.md +++ b/docs/post-1.9.0-progress-20261006.md @@ -268,3 +268,39 @@ Reticulum daemons, live Minibits and the subprocess permission helper (which its parent test executes separately). A production build of the earlier integration is running from detached `11f016a9`; it does not include this new preview fix and must not be described as the final release candidate. + + +### Bounded peer delivery and preview qualification + +Paid-preview source at `051dc7e3` passed all four isolated regression cases +(`/tmp/archy-preview-boundary-tests.log`). No live deployment is claimed. +The earlier production backend at `11f016a9` also built successfully and was +archived with its SHA256/source receipt under the private qualification directory; +it excludes these later fixes and is not the final candidate. + +`2fee0339` replaces whole-file seller buffering with bounded file-backed responses. +Bearer payments prepare a complete private anonymous snapshot before redemption; +free/owner/durable-invoice transfers can stream an open file directly. Rootless +Files reads consume bounded subprocess stdout and require successful completion +before payment. Malformed ranges are rejected, suffix ranges are supported, and +free public previews also stream. New tests cover source deletion during payment, +invalid payment, multi-gigabyte sparse files, truncated preparation and ranges. +Full isolated backend qualification is running from the separate frozen media +worktree (`/tmp/archy-bounded-media-backend-tests.log`); results remain pending. +Buyer-side caching still needs bounded transfer and recovery work. + +Buyer follow-up now streams successful ecash/Lightning deliveries into the owned +cache, records incomplete delivery before consuming the response, removes partial +temporary files on cancellation, and blocks duplicate concurrent ecash purchases +per seller. Owned-file opens and saves use local HTTP streaming rather than base64 +for current clients; small legacy reads remain supported. Optional Files copies +stream through the existing no-clobber namespace writer. Interrupted receipt/body +handling is not equivalent to a durable end-to-end ecash retry protocol: loss +before response headers or during mint settlement remains an explicit review gate. +No real funds were spent. Nineteen focused UI tests passed before the final two +stream-viewer cases were added; backend/production qualification is pending. + +Seller streaming's first full compile found a lifetime error in one new test; +`df7677d2` corrects it. The repeated isolated full suite is compiling from that +frozen source (`/tmp/archy-bounded-media-backend-tests-2.log`). The failed run is +retained and is not counted as a pass. diff --git a/neode-ui/src/composables/__tests__/usePaidItemViewer.test.ts b/neode-ui/src/composables/__tests__/usePaidItemViewer.test.ts index 1868ad46..f4f759fe 100644 --- a/neode-ui/src/composables/__tests__/usePaidItemViewer.test.ts +++ b/neode-ui/src/composables/__tests__/usePaidItemViewer.test.ts @@ -75,6 +75,26 @@ describe('usePaidItemViewer — UIFIX-04 (lightbox routing) + UIFIX-06 (loading/ vi.stubGlobal('atob', vi.fn(() => 'binarydata')) }) + it('opens a large cached video through authenticated HTTP without decoding or a blob allocation', async () => { + mockedRpc.call.mockResolvedValue({ owned: true, mime_type: 'video/mp4', size_bytes: 200000000 }) + const viewer = usePaidItemViewer() + await viewer.open(VIDEO_ITEM) + expect(await viewer.resolveBlobUrl(viewer.lightboxItems.value[0]!.path)).toBe('/api/peer-content/abc123.onion/content-2') + expect(createObjectURLSpy).not.toHaveBeenCalled() + expect(atob).not.toHaveBeenCalled() + expect(mockedRpc.call).toHaveBeenCalledWith(expect.objectContaining({ params: { onion: VIDEO_ITEM.onion, content_id: VIDEO_ITEM.content_id, cache_only: true } })) + }) + + it('does not open incomplete delivery or start another purchase', async () => { + mockedRpc.call.mockRejectedValue(new Error('Delivery incomplete')) + const viewer = usePaidItemViewer() + await viewer.open(VIDEO_ITEM) + expect(viewer.lightboxIndex.value).toBeNull() + expect(viewer.error.value).toBeTruthy() + expect(mockedRpc.call).toHaveBeenCalledTimes(1) + expect(createObjectURLSpy).not.toHaveBeenCalled() + }) + it('routes an image mime to the lightbox, not window.open', async () => { mockedRpc.call.mockResolvedValue({ data_base64: 'ZmFrZQ==', mime_type: 'image/jpeg' }) const viewer = usePaidItemViewer() diff --git a/neode-ui/src/composables/usePaidItemViewer.ts b/neode-ui/src/composables/usePaidItemViewer.ts index 27913bd2..3ce30f9e 100644 --- a/neode-ui/src/composables/usePaidItemViewer.ts +++ b/neode-ui/src/composables/usePaidItemViewer.ts @@ -22,7 +22,7 @@ export function paidItemKey(it: OwnedItemLike): string { } /** - * Fetch, decode, route-to-viewer and loading/error state for a purchased + * Validate ownership, stream to viewer and manage loading/error state for a purchased * item (UIFIX-04 / UIFIX-06). * * - image/video → routed into the caller's MediaLightbox via `lightboxItems` @@ -42,7 +42,7 @@ export function usePaidItemViewer() { const lightboxItems = ref([]) const lightboxIndex = ref(null) - // Synthetic-path -> already-fetched blob URL. Populated only for items + // Synthetic-path -> authenticated cache URL (or legacy blob URL). Populated only for items // routed to the lightbox; resolveBlobUrl reads from here and never fetches. const urlByPath = new Map() // One in-flight fetch per item key — a second open() for the same key @@ -61,21 +61,20 @@ export function usePaidItemViewer() { opening.value = key error.value = null try { - const res = await rpcClient.call<{ data_base64?: string; data?: string; mime_type?: string }>({ + const res = await rpcClient.call<{ data_base64?: string; data?: string; owned?: boolean; mime_type?: string; error?: string }>({ method: 'content.owned-get', - params: { onion: it.onion, content_id: it.content_id }, + params: { onion: it.onion, content_id: it.content_id, cache_only: true }, timeout: 60000, }) - const b64 = res.data_base64 || res.data - if (!b64) { - error.value = "Couldn't open this file — the peer returned no data." + const b64 = res.data_base64 ?? res.data + if (b64 === undefined && res.owned !== true) { + error.value = res.error || "Couldn't open this file — its cached copy is unavailable." return } - const bin = atob(b64) - const arr = new Uint8Array(bin.length) - for (let i = 0; i < bin.length; i++) arr[i] = bin.charCodeAt(i) const mime = res.mime_type || it.mime_type - const url = URL.createObjectURL(new Blob([arr], { type: mime })) + const url = b64 === undefined + ? `/api/peer-content/${encodeURIComponent(it.onion)}/${encodeURIComponent(it.content_id)}` + : URL.createObjectURL(new Blob([Uint8Array.from(atob(b64), c => c.charCodeAt(0))], { type: mime })) // Music ALWAYS plays in the global bottom-bar player — never a // lightbox (blob URL stays alive for the bar; it owns playback now). @@ -104,7 +103,7 @@ export function usePaidItemViewer() { // No in-app viewer (documents, etc.) — keep today's behaviour. window.open(url, '_blank', 'noopener') - setTimeout(() => URL.revokeObjectURL(url), 60000) + if (url.startsWith("blob:")) setTimeout(() => URL.revokeObjectURL(url), 60000) } catch { error.value = "Couldn't open this file — it may be unavailable right now." } finally { diff --git a/neode-ui/src/views/PeerFiles.vue b/neode-ui/src/views/PeerFiles.vue index 82c7e54f..1f6b4e10 100644 --- a/neode-ui/src/views/PeerFiles.vue +++ b/neode-ui/src/views/PeerFiles.vue @@ -701,7 +701,7 @@ function base64ToBlob(base64: string, mime: string): Blob { } // In-app viewer (lightbox) for owned content AND free images. Owned content is -// loaded as a blob URL from the purchase cache; free images point straight at +// loaded through local cached streaming; free images point straight at // the Range-capable streaming proxy (no base64 round-trip). const viewerItem = ref(null) const viewerUrl = ref(null) @@ -717,22 +717,23 @@ async function viewOwned(item: CatalogItem) { playing.value = item.id purchaseError.value = null try { - const res = await rpcClient.call<{ data?: string; mime_type?: string; error?: string }>({ + const res = await rpcClient.call<{ data?: string; owned?: boolean; owned_content_id?: string; mime_type?: string; error?: string }>({ method: 'content.owned-get', - params: { onion, content_id: item.id }, + params: { onion, content_id: item.id, filename: item.filename, cache_only: true }, timeout: 60000, }) - if (!res?.data) { purchaseError.value = res?.error || 'Could not open your purchased file'; return } + if (res?.data === undefined && res?.owned !== true) { purchaseError.value = res?.error || 'Could not open your purchased file'; return } const mime = res.mime_type || item.mime_type + const url = res.data !== undefined ? URL.createObjectURL(base64ToBlob(res.data, mime)) + : `/api/peer-content/${encodeURIComponent(onion)}/${encodeURIComponent(res.owned_content_id || item.id)}` // Audio always plays in the global bottom-bar player — never the lightbox // (the blob URL is intentionally not revoked while the bar plays it). if (mime.startsWith('audio/')) { - const url = URL.createObjectURL(base64ToBlob(res.data, mime)) audioPlayer.play(url, item.filename.split('/').pop() || item.filename) return } releaseViewerUrl() - viewerUrl.value = URL.createObjectURL(base64ToBlob(res.data, mime)) + viewerUrl.value = url viewerMime.value = mime viewerItem.value = item } catch (e: unknown) { @@ -797,12 +798,13 @@ async function saveOwned(item: CatalogItem) { const onion = props.peerId || currentPeer.value?.onion if (!onion) return try { - const res = await rpcClient.call<{ data?: string; mime_type?: string; error?: string }>({ + const res = await rpcClient.call<{ data?: string; owned?: boolean; owned_content_id?: string; mime_type?: string; error?: string }>({ method: 'content.owned-get', - params: { onion, content_id: item.id }, + params: { onion, content_id: item.id, filename: item.filename, cache_only: true }, timeout: 60000, }) - if (res?.data) triggerDownload(res.data, item) + if (res?.owned === true) streamDownload(`/api/peer-content/${encodeURIComponent(onion)}/${encodeURIComponent(res.owned_content_id || item.id)}`, item) + else if (res?.data) triggerDownload(res.data, item) else purchaseError.value = res?.error || 'Could not save your purchased file' } catch (e: unknown) { purchaseError.value = e instanceof Error ? e.message : 'Could not save your purchased file' @@ -1362,11 +1364,11 @@ async function prepareEcashPay() { * mobile companion ("paid but never unlocked"); the viewer's Save button * still offers an explicit download. */ -function openPurchased(item: CatalogItem, base64Data: string | undefined, mimeType?: string, seller?: string) { +function openPurchased(item: CatalogItem, base64Data: string | undefined, mimeType?: string, seller?: string, cachedId?: string) { const onion = seller || props.peerId || currentPeer.value?.onion const url = base64Data !== undefined ? URL.createObjectURL(base64ToBlob(base64Data, mimeType || item.mime_type)) - : `/api/peer-content/${encodeURIComponent(onion || "")}/${encodeURIComponent(item.id)}` + : `/api/peer-content/${encodeURIComponent(onion || "")}/${encodeURIComponent(cachedId || item.id)}` if (onion) { try { localStorage.removeItem(receiptKey(onion, item.id)) } catch { /* owned cache remains authoritative */ } lnReceipt.value = null @@ -1395,20 +1397,20 @@ async function confirmEcashPay() { const item = payItem.value const onion = props.peerId || currentPeer.value?.onion const method = ecashPlan.value?.chosen - if (!item || !onion || !method) return + if (!item || !onion || !method || downloading.value) return const price = getItemPrice(item.access) downloading.value = item.id purchaseError.value = null try { - const result = await rpcClient.call<{ data?: string; error?: string; ecash_backend?: string; mime_type?: string }>({ + const result = await rpcClient.call<{ data?: string; owned?: boolean; owned_content_id?: string; error?: string; ecash_backend?: string; mime_type?: string }>({ method: 'content.download-peer-paid', - params: { onion, content_id: item.id, price_sats: price, method, filename: item.filename }, - timeout: 120000, + params: { onion, content_id: item.id, price_sats: price, method, filename: item.filename, cache_only: true }, + timeout: 960000, }) - if (result?.data) { + if (result?.data !== undefined || result?.owned === true) { // The purchase is now cached + owned by this node (backend persisted it). - openPurchased(item, result.data, result.mime_type) + openPurchased(item, result.data, result.mime_type, onion, result.owned_content_id) } else if (result?.error) { // Keep the confirm screen open so the user can switch backend and retry. purchaseError.value = result.error @@ -1486,7 +1488,7 @@ async function payWithLightning() { } } lnReceipt.value = inv - const dl = await rpcClient.call<{ data?: string; owned?: boolean; mime_type?: string; error?: string }>({ + const dl = await rpcClient.call<{ data?: string; owned?: boolean; owned_content_id?: string; mime_type?: string; error?: string }>({ method: 'content.download-peer-invoice', params: { onion, content_id: item.id, payment_hash: inv.payment_hash, filename: item.filename, price_sats: inv.price_sats, cache_only: true }, timeout: 960000, @@ -1517,7 +1519,7 @@ async function pollInvoice() { if (res?.paid) { // Settled — pull the file using the payment hash as the gate token. invoiceWaiting.value = false - const dl = await rpcClient.call<{ data?: string; owned?: boolean; mime_type?: string; error?: string }>({ + const dl = await rpcClient.call<{ data?: string; owned?: boolean; owned_content_id?: string; mime_type?: string; error?: string }>({ method: 'content.download-peer-invoice', params: { onion, content_id: item.id, payment_hash: inv.payment_hash, filename: item.filename, price_sats: inv.price_sats, cache_only: true }, timeout: 960000, diff --git a/neode-ui/src/views/__tests__/PeerFilesLightning.test.ts b/neode-ui/src/views/__tests__/PeerFilesLightning.test.ts index 8d71b05b..611ce64b 100644 --- a/neode-ui/src/views/__tests__/PeerFilesLightning.test.ts +++ b/neode-ui/src/views/__tests__/PeerFilesLightning.test.ts @@ -30,6 +30,35 @@ beforeEach(() => { vi.mocked(rpcClient.payLightningInvoice).mockResolvedValue({ status: 'succeeded' } as never) }) describe('Lightning file delivery recovery', () => { + it('opens a cached ecash purchase without transferring base64 into the UI', async () => { + vi.mocked(rpcClient.call).mockImplementation(async ({ method }) => { + if (method === 'content.download-peer-paid') return { owned: true, mime_type: 'video/mp4', size_bytes: 200000000 } + return { items: [] } + }) + const { wrapper, vm } = await open() + vm.ecashPlan = { cashu: 10, fedimint: 0, ark: 0, total: 10, chosen: 'cashu' } + await vm.confirmEcashPay() + expect(vm.viewerUrl).toBe('/api/peer-content/peer.onion/paid-file') + expect(vm.viewerMime).toBe('video/mp4') + expect(vi.mocked(rpcClient.call).mock.calls.find(([v]) => v.method === 'content.download-peer-paid')![0].params).toMatchObject({ cache_only: true, method: 'cashu' }) + wrapper.unmount() + }) + it('does not issue a duplicate ecash purchase while delivery is pending', async () => { + let finish!: (value: unknown) => void + vi.mocked(rpcClient.call).mockImplementation(async ({ method }) => { + if (method === 'content.download-peer-paid') return await new Promise(resolve => { finish = resolve }) + return { items: [] } + }) + const { wrapper, vm } = await open() + vm.ecashPlan = { cashu: 10, fedimint: 0, ark: 0, total: 10, chosen: 'cashu' } + const first = vm.confirmEcashPay() + await vm.confirmEcashPay() + expect(vi.mocked(rpcClient.call).mock.calls.filter(([v]) => v.method === 'content.download-peer-paid')).toHaveLength(1) + finish({ error: 'Delivery needs recovery; do not pay again' }) + await first + expect(vm.purchaseError).toContain('do not pay again') + wrapper.unmount() + }) it('retries delivery after a seller rejection without paying or requesting another invoice', async () => { download.mockResolvedValue({ error: 'Seller has not registered this payment yet' }) const { wrapper, vm } = await open() diff --git a/neode-ui/src/views/web5/Web5SharedContent.vue b/neode-ui/src/views/web5/Web5SharedContent.vue index 3712b478..7cd07232 100644 --- a/neode-ui/src/views/web5/Web5SharedContent.vue +++ b/neode-ui/src/views/web5/Web5SharedContent.vue @@ -612,26 +612,29 @@ async function purchaseAndDownload(item: PeerContentItem) { purchasingId.value = item.id try { - const result = await rpcClient.call<{ data?: string; error?: string }>({ + const result = await rpcClient.call<{ data?: string; owned?: boolean; owned_content_id?: string; error?: string }>({ method: 'content.download-peer-paid', - params: { onion: browsePeerOnion.value, content_id: item.id, price_sats: price }, - timeout: 120000, + params: { onion: browsePeerOnion.value, content_id: item.id, price_sats: price, filename: item.filename, cache_only: true }, + timeout: 960000, }) - if (result?.data) { - const blob = new Blob( - [Uint8Array.from(atob(result.data), c => c.charCodeAt(0))], - { type: item.mime_type }, - ) + if (result?.owned === true) { + const a = document.createElement('a') + a.href = `/api/peer-content/${encodeURIComponent(browsePeerOnion.value)}/${encodeURIComponent(result.owned_content_id || item.id)}` + a.download = item.filename.split('/').pop() || item.filename + a.click() + emit('toast', 'Purchase saved — download started') + } else if (result?.data) { + const blob = new Blob([Uint8Array.from(atob(result.data), c => c.charCodeAt(0))], { type: item.mime_type }) const url = URL.createObjectURL(blob) const a = document.createElement('a') a.href = url a.download = item.filename.split('/').pop() || item.filename a.click() - URL.revokeObjectURL(url) - emit('toast', `Downloaded for ${price} sats`) + setTimeout(() => URL.revokeObjectURL(url), 60000) + emit('toast', 'Purchase saved — download started') } else { - emit('toast', 'Purchase failed — no data received') + emit('toast', result?.error || 'Purchase could not be completed') } } catch (e: unknown) { emit('toast', e instanceof Error ? e.message : 'Purchase failed')