Merge branch 'work/post-190-onchain-stream' into work/post-190-session-key
This commit is contained in:
@@ -1203,6 +1203,28 @@ impl RpcHandler {
|
|||||||
return Err(anyhow::anyhow!("Invalid address"));
|
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 (data, _) = self.state_manager.get_snapshot().await;
|
||||||
let local_did = crate::identity::did_key_from_pubkey_hex(&data.server_info.pubkey)?;
|
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;
|
let fips_npub = crate::federation::fips_npub_for_onion(&self.config.data_dir, onion).await;
|
||||||
@@ -1251,16 +1273,38 @@ impl RpcHandler {
|
|||||||
}));
|
}));
|
||||||
}
|
}
|
||||||
|
|
||||||
let bytes = response
|
let mime = response
|
||||||
.bytes()
|
.headers()
|
||||||
.await
|
.get(reqwest::header::CONTENT_TYPE)
|
||||||
.context("Failed to read response body")?;
|
.and_then(|v| v.to_str().ok())
|
||||||
use base64::Engine;
|
.unwrap_or("application/octet-stream")
|
||||||
let encoded = base64::engine::general_purpose::STANDARD.encode(&bytes);
|
.split(';')
|
||||||
Ok(serde_json::json!({
|
.next()
|
||||||
"data": encoded,
|
.unwrap_or("application/octet-stream")
|
||||||
"size": bytes.len(),
|
.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).
|
/// Fetch a preview of paid content from a peer (no payment required).
|
||||||
|
|||||||
@@ -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!(bounded_content_bytes(response, 4).await.is_err());
|
||||||
assert_eq!(consumed.load(Ordering::SeqCst), 2);
|
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());
|
||||||
|
}
|
||||||
|
|||||||
@@ -1288,14 +1288,14 @@ async function pollOnchain(address: string) {
|
|||||||
timeout: 30000,
|
timeout: 30000,
|
||||||
})
|
})
|
||||||
if (res?.paid) {
|
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',
|
method: 'content.download-peer-onchain',
|
||||||
params: { onion, content_id: item.id, address },
|
params: { onion, content_id: item.id, address, filename: item.filename, price_sats: getItemPrice(item.access), cache_only: true },
|
||||||
timeout: 120000,
|
timeout: 960000,
|
||||||
})
|
})
|
||||||
onchainPaying.value = false
|
onchainPaying.value = false
|
||||||
if (dl?.data) {
|
if (dl?.data !== undefined || dl?.owned === true) {
|
||||||
openPurchased(item, dl.data, dl.mime_type)
|
openPurchased(item, dl.data, dl.mime_type, onion, dl.owned_content_id)
|
||||||
} else {
|
} else {
|
||||||
lnError.value = dl?.error || 'Paid, but the download failed. Try again shortly.'
|
lnError.value = dl?.error || 'Paid, but the download failed. Try again shortly.'
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -30,6 +30,21 @@ beforeEach(() => {
|
|||||||
vi.mocked(rpcClient.payLightningInvoice).mockResolvedValue({ status: 'succeeded' } as never)
|
vi.mocked(rpcClient.payLightningInvoice).mockResolvedValue({ status: 'succeeded' } as never)
|
||||||
})
|
})
|
||||||
describe('Lightning file delivery recovery', () => {
|
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 () => {
|
it('opens a cached ecash purchase without transferring base64 into the UI', async () => {
|
||||||
vi.mocked(rpcClient.call).mockImplementation(async ({ method }) => {
|
vi.mocked(rpcClient.call).mockImplementation(async ({ method }) => {
|
||||||
if (method === 'content.download-peer-paid') return { owned: true, mime_type: 'video/mp4', size_bytes: 200000000 }
|
if (method === 'content.download-peer-paid') return { owned: true, mime_type: 'video/mp4', size_bytes: 200000000 }
|
||||||
|
|||||||
Reference in New Issue
Block a user