From 2e761666d5701a5d929508f6bc5f5c31e729c48b Mon Sep 17 00:00:00 2001 From: archipelago Date: Tue, 6 Oct 2026 06:43:19 -0400 Subject: [PATCH] Stream confirmed on-chain file deliveries into the owned cache --- core/archipelago/src/api/rpc/content.rs | 64 ++++++++++++++++--- core/archipelago/src/api/rpc/content_tests.rs | 31 +++++++++ neode-ui/src/views/PeerFiles.vue | 10 +-- .../__tests__/PeerFilesLightning.test.ts | 15 +++++ 4 files changed, 105 insertions(+), 15 deletions(-) diff --git a/core/archipelago/src/api/rpc/content.rs b/core/archipelago/src/api/rpc/content.rs index 89ff3303..f046791b 100644 --- a/core/archipelago/src/api/rpc/content.rs +++ b/core/archipelago/src/api/rpc/content.rs @@ -1203,6 +1203,28 @@ impl RpcHandler { return Err(anyhow::anyhow!("Invalid address")); } + 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 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 cached_purchase_response( + &self.config.data_dir, + onion, + content_id, + cache_only, + 0, + ) + .await; + } let (data, _) = self.state_manager.get_snapshot().await; let local_did = crate::identity::did_key_from_pubkey_hex(&data.server_info.pubkey)?; let fips_npub = crate::federation::fips_npub_for_onion(&self.config.data_dir, onion).await; @@ -1251,16 +1273,38 @@ impl RpcHandler { })); } - let bytes = response - .bytes() - .await - .context("Failed to read response body")?; - use base64::Engine; - let encoded = base64::engine::general_purpose::STANDARD.encode(&bytes); - Ok(serde_json::json!({ - "data": encoded, - "size": bytes.len(), - })) + let mime = response + .headers() + .get(reqwest::header::CONTENT_TYPE) + .and_then(|v| v.to_str().ok()) + .unwrap_or("application/octet-stream") + .split(';') + .next() + .unwrap_or("application/octet-stream") + .to_string(); + let filename = params + .get("filename") + .and_then(|v| v.as_str()) + .unwrap_or(content_id); + let item = cache_peer_response( + &self.config.data_dir, + onion, + content_id, + filename, + &mime, + params + .get("price_sats") + .and_then(|v| v.as_u64()) + .unwrap_or(0), + "onchain", + response, + ) + .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!("On-chain purchase cached; optional Files copy failed: {error:#}"); + } + cached_purchase_response(&self.config.data_dir, onion, content_id, cache_only, 0).await } /// Fetch a preview of paid content from a peer (no payment required). diff --git a/core/archipelago/src/api/rpc/content_tests.rs b/core/archipelago/src/api/rpc/content_tests.rs index 0697246a..76f01bc0 100644 --- a/core/archipelago/src/api/rpc/content_tests.rs +++ b/core/archipelago/src/api/rpc/content_tests.rs @@ -191,3 +191,34 @@ async fn inline_download_limits_known_and_chunked_bodies_without_draining_them() assert!(bounded_content_bytes(response, 4).await.is_err()); assert_eq!(consumed.load(Ordering::SeqCst), 2); } + +#[tokio::test] +async fn onchain_delivery_streams_to_owned_cache_and_preserves_incomplete_recovery() { + use tokio::io::AsyncReadExt; + let dir = tempfile::tempdir().unwrap(); + let response: reqwest::Response = hyper::Response::builder() + .header("content-length", "2097152") + .body(hyper::Body::from(vec![71u8; 2 * 1024 * 1024])) + .unwrap() + .into(); + let item = cache_peer_response(dir.path(), "seller.onion", "film", "film.mp4", "video/mp4", 1, "onchain", response).await.unwrap(); + assert!(item.download_complete); + assert_eq!(item.ecash_backend, "onchain"); + let (_, mut file) = crate::content_owned::open_owned(dir.path(), "seller.onion", "film").await.unwrap().unwrap(); + let mut bytes = Vec::new(); + file.read_to_end(&mut bytes).await.unwrap(); + assert_eq!(bytes, vec![71u8; 2 * 1024 * 1024]); + let reply = cached_purchase_response(dir.path(), "seller.onion", "film", true, 0).await.unwrap(); + assert_eq!(reply["owned"], true); + assert!(reply.get("data").is_none()); + + let truncated: reqwest::Response = hyper::Response::builder() + .header("content-length", "100") + .body(hyper::Body::from("short")) + .unwrap() + .into(); + assert!(cache_peer_response(dir.path(), "seller.onion", "interrupted", "film.mp4", "video/mp4", 1, "onchain", truncated).await.is_err()); + let entries = crate::content_owned::list_owned_checked(dir.path()).await.unwrap(); + assert!(!entries.iter().find(|item| item.content_id == "interrupted").unwrap().download_complete); + assert!(existing_paid_content(dir.path(), "seller.onion", "interrupted", None, true).await.is_err()); +} diff --git a/neode-ui/src/views/PeerFiles.vue b/neode-ui/src/views/PeerFiles.vue index 1f6b4e10..d131137a 100644 --- a/neode-ui/src/views/PeerFiles.vue +++ b/neode-ui/src/views/PeerFiles.vue @@ -1288,14 +1288,14 @@ async function pollOnchain(address: string) { timeout: 30000, }) if (res?.paid) { - const dl = await rpcClient.call<{ data?: string; 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-onchain', - params: { onion, content_id: item.id, address }, - timeout: 120000, + params: { onion, content_id: item.id, address, filename: item.filename, price_sats: getItemPrice(item.access), cache_only: true }, + timeout: 960000, }) onchainPaying.value = false - if (dl?.data) { - openPurchased(item, dl.data, dl.mime_type) + if (dl?.data !== undefined || dl?.owned === true) { + openPurchased(item, dl.data, dl.mime_type, onion, dl.owned_content_id) } else { lnError.value = dl?.error || 'Paid, but the download failed. Try again shortly.' } diff --git a/neode-ui/src/views/__tests__/PeerFilesLightning.test.ts b/neode-ui/src/views/__tests__/PeerFilesLightning.test.ts index 611ce64b..d46dde5c 100644 --- a/neode-ui/src/views/__tests__/PeerFilesLightning.test.ts +++ b/neode-ui/src/views/__tests__/PeerFilesLightning.test.ts @@ -30,6 +30,21 @@ beforeEach(() => { vi.mocked(rpcClient.payLightningInvoice).mockResolvedValue({ status: 'succeeded' } as never) }) describe('Lightning file delivery recovery', () => { + it('opens a confirmed on-chain delivery from HTTP cache without another payment', async () => { + vi.mocked(rpcClient.call).mockImplementation(async ({ method }) => { + if (method === 'content.onchain-status') return { paid: true } + if (method === 'content.download-peer-onchain') return { owned: true, mime_type: 'video/mp4', size_bytes: 200000000 } + return { items: [] } + }) + const { wrapper, vm } = await open() + await vm.pollOnchain('bc1test') + expect(vm.viewerUrl).toBe('/api/peer-content/peer.onion/paid-file') + expect(vm.viewerMime).toBe('video/mp4') + const calls = vi.mocked(rpcClient.call).mock.calls.map(([call]) => call) + expect(calls.find(call => call.method === 'content.download-peer-onchain')?.params).toMatchObject({ cache_only: true, address: 'bc1test', filename: 'bought.txt' }) + expect(calls.some(call => ['lnd.sendcoins', 'content.request-onchain'].includes(call.method))).toBe(false) + wrapper.unmount() + }) 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 }