diff --git a/core/archipelago/src/api/handler/mod.rs b/core/archipelago/src/api/handler/mod.rs
index f615912d..77d7b8c1 100644
--- a/core/archipelago/src/api/handler/mod.rs
+++ b/core/archipelago/src/api/handler/mod.rs
@@ -1,16 +1,16 @@
-mod purchase;
-mod cloud_purchase;
mod blob;
mod cdp;
+mod cloud_purchase;
mod content;
-mod registered_media;
-mod rental_playback;
mod dwn;
mod model_proxy;
mod node_message;
mod proxy;
+mod purchase;
+mod registered_media;
mod remote_input;
mod remote_relay;
+mod rental_playback;
mod routstr_proxy;
mod websocket;
@@ -389,7 +389,9 @@ impl ApiHandler {
let method = req.method().clone();
if path.starts_with("/api/rental-playback/") {
- return self.handle_local_rental_request(&method, &path, req.headers()).await;
+ return self
+ .handle_local_rental_request(&method, &path, req.headers())
+ .await;
}
// Handle CORS preflight for all routes
@@ -453,13 +455,28 @@ impl ApiHandler {
}
// Purchase routes bound the original body before the generic buffer.
- if method == Method::POST && matches!(path.as_str(),
- crate::content_purchase_protocol::OFFER_ROUTE | crate::content_purchase_protocol::ACCEPT_ROUTE
- | crate::content_purchase_protocol::SETTLE_ROUTE | crate::content_purchase_protocol::STATUS_ROUTE
- | crate::content_purchase_protocol::CANCEL_ROUTE) {
+ if method == Method::POST
+ && matches!(
+ path.as_str(),
+ crate::content_purchase_protocol::PREPARE_OFFER_ROUTE
+ | crate::content_purchase_protocol::OFFER_ROUTE
+ | crate::content_purchase_protocol::ACCEPT_ROUTE
+ | crate::content_purchase_protocol::SETTLE_ROUTE
+ | crate::content_purchase_protocol::STATUS_ROUTE
+ | crate::content_purchase_protocol::CANCEL_ROUTE
+ )
+ {
return self.handle_purchase_request(req).await;
}
+ if method == Method::POST
+ && path.starts_with("/content/registered_")
+ && path.contains("/rental/")
+ && (path.ends_with("/prepare") || path.ends_with("/start"))
+ {
+ return self.handle_rental_control(req).await;
+ }
+
// Convert body to bytes for non-WS routes
let headers = req.headers().clone();
let query_string = req.uri().query().map(|s| s.to_string()).unwrap_or_default();
@@ -649,21 +666,31 @@ impl ApiHandler {
// falls outside `connect-src`). Session-authenticated so only
// the logged-in node owner can spin up fetches.
(Method::GET, "/api/node-app-catalog") => {
- if !self.is_authenticated(&headers).await { return Ok(Self::unauthorized()); }
+ if !self.is_authenticated(&headers).await {
+ return Ok(Self::unauthorized());
+ }
let data_dir = self.config.data_dir.clone();
let result = tokio::task::spawn_blocking(move || {
crate::container::node_catalog::verified_body(&data_dir)
- }).await.unwrap_or_else(|error| Err(anyhow::anyhow!(error)));
+ })
+ .await
+ .unwrap_or_else(|error| Err(anyhow::anyhow!(error)));
let (status, body) = match result {
Ok(Some(body)) => (StatusCode::OK, body),
Ok(None) => (StatusCode::NOT_FOUND, "{}".to_owned()),
Err(error) => {
tracing::warn!("Node demo catalog rejected: {error}");
- (StatusCode::CONFLICT, "{\"error\":\"Node demo catalog is unavailable\"}".to_owned())
- },
+ (
+ StatusCode::CONFLICT,
+ "{\"error\":\"Node demo catalog is unavailable\"}".to_owned(),
+ )
+ }
};
- Ok(Response::builder().status(status).header("Content-Type", "application/json")
- .header("Cache-Control", "private, no-store").body(hyper::Body::from(body))?)
+ Ok(Response::builder()
+ .status(status)
+ .header("Content-Type", "application/json")
+ .header("Cache-Control", "private, no-store")
+ .body(hyper::Body::from(body))?)
}
(Method::GET, "/api/app-catalog") => {
diff --git a/core/archipelago/src/api/handler/purchase.rs b/core/archipelago/src/api/handler/purchase.rs
index 1a469584..74697a10 100644
--- a/core/archipelago/src/api/handler/purchase.rs
+++ b/core/archipelago/src/api/handler/purchase.rs
@@ -23,7 +23,8 @@ impl ApiHandler {
request.method() == Method::POST
&& matches!(
path.as_str(),
- protocol::OFFER_ROUTE
+ protocol::PREPARE_OFFER_ROUTE
+ | protocol::OFFER_ROUTE
| protocol::ACCEPT_ROUTE
| protocol::SETTLE_ROUTE
| protocol::STATUS_ROUTE
@@ -59,6 +60,32 @@ impl ApiHandler {
)?;
let data_dir = &self.config.data_dir;
let result = match path.as_str() {
+ protocol::PREPARE_OFFER_ROUTE => {
+ #[derive(Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct Prepare {
+ content_id: String,
+ #[serde(default)]
+ retry: bool,
+ }
+ let input: Prepare = serde_json::from_slice(&bytes)?;
+ let identity = std::sync::Arc::new(
+ NodeIdentity::load_existing(&data_dir.join("identity")).await?,
+ );
+ anyhow::ensure!(identity.did_key()? == audience, "Node identity changed");
+ let root = data_dir.clone();
+ serde_json::to_value(
+ tokio::task::spawn_blocking(move || {
+ crate::registered_media::prepare_registered(
+ root,
+ identity,
+ &input.content_id,
+ input.retry,
+ )
+ })
+ .await??,
+ )?
+ }
protocol::OFFER_ROUTE => {
let body: OfferRequest = serde_json::from_slice(&bytes)?;
anyhow::ensure!(
diff --git a/core/archipelago/src/api/handler/registered_media.rs b/core/archipelago/src/api/handler/registered_media.rs
index 09e36572..6a9a8b72 100644
--- a/core/archipelago/src/api/handler/registered_media.rs
+++ b/core/archipelago/src/api/handler/registered_media.rs
@@ -4,7 +4,6 @@ use crate::{content_server::ByteRange, identity::NodeIdentity, registered_media:
use anyhow::{Context, Result};
use hyper::{Body, HeaderMap, Response, StatusCode};
use std::sync::Arc;
-use tokio::io::{AsyncReadExt, AsyncSeekExt};
fn route(path: &str) -> Result<(&str, &str)> {
let (content, purchase) = path
@@ -46,6 +45,125 @@ fn denied(message: &'static str) -> Response
{
}
impl ApiHandler {
+ pub(super) async fn handle_rental_control(
+ &self,
+ mut request: hyper::Request,
+ ) -> Result> {
+ use hyper::body::HttpBody;
+ #[derive(serde::Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct Control {
+ capability: String,
+ ready_id: Option,
+ #[serde(default)]
+ retry: bool,
+ }
+ let path = request.uri().path().to_owned();
+ let (base, action) = path.rsplit_once('/').context("Invalid rental action")?;
+ anyhow::ensure!(
+ matches!(action, "prepare" | "start") && request.method() == hyper::Method::POST,
+ "Invalid rental action"
+ );
+ let (content, purchase) = route(base)?;
+ let bytes = tokio::time::timeout(std::time::Duration::from_secs(15), async {
+ let mut bytes = Vec::new();
+ while let Some(chunk) = request.body_mut().data().await {
+ let chunk = chunk?;
+ anyhow::ensure!(
+ bytes
+ .len()
+ .checked_add(chunk.len())
+ .is_some_and(|n| n <= 16 * 1024),
+ "Rental request too large"
+ );
+ bytes.extend_from_slice(&chunk);
+ }
+ Ok::<_, anyhow::Error>(bytes)
+ })
+ .await
+ .context("Rental request timed out")??;
+ let audience = crate::identity::did_key_from_pubkey_hex(&self.self_pubkey_hex)?;
+ let buyer = crate::content_auth::authenticate_request(
+ request.headers(),
+ &audience,
+ &hyper::Method::POST,
+ &path,
+ &bytes,
+ chrono::Utc::now().timestamp(),
+ )?;
+ let input: Control = serde_json::from_slice(&bytes)?;
+ let identity =
+ Arc::new(NodeIdentity::load_existing(&self.config.data_dir.join("identity")).await?);
+ anyhow::ensure!(identity.did_key()? == audience, "Node identity changed");
+ let result = if action == "prepare" {
+ anyhow::ensure!(input.ready_id.is_none(), "Prepare does not start a rental");
+ let prior = crate::registered_media::paid_window(
+ &self.config.data_dir,
+ &identity,
+ content,
+ purchase,
+ &buyer,
+ &input.capability,
+ )
+ .await?;
+ let metadata = crate::registered_media::registered_metadata(
+ &self.config.data_dir,
+ &identity,
+ content,
+ )?;
+ anyhow::ensure!(
+ prior
+ .as_ref()
+ .is_none_or(|window| clock() >= window.started_at),
+ "Rental clock moved backwards"
+ );
+ if let Some(window) = prior.as_ref().filter(|window| clock() >= window.expires_at) {
+ serde_json::json!({"state":"expired", "viewing_seconds":metadata.0.viewing_seconds,"started_at":window.started_at,"expires_at":window.expires_at})
+ } else {
+ let state = crate::registered_media::prepare_paid(
+ self.config.data_dir.clone(),
+ identity,
+ content.into(),
+ purchase.into(),
+ buyer,
+ input.capability,
+ input.retry,
+ )
+ .await?;
+ let mut result = serde_json::to_value(state)?;
+ result["viewing_seconds"] = serde_json::json!(metadata.0.viewing_seconds);
+ result["started_at"] =
+ serde_json::json!(prior.as_ref().map(|window| window.started_at));
+ result["expires_at"] =
+ serde_json::json!(prior.as_ref().map(|window| window.expires_at));
+ result
+ }
+ } else {
+ anyhow::ensure!(!input.retry, "Start cannot retry verification");
+ let ready_id = input
+ .ready_id
+ .context("Media must be ready before explicit Start")?;
+ let window = crate::registered_media::start_paid(
+ self.config.data_dir.clone(),
+ identity,
+ content.into(),
+ purchase.into(),
+ buyer,
+ input.capability,
+ ready_id,
+ )
+ .await?;
+ anyhow::ensure!(clock() >= window.started_at, "Rental clock moved backwards");
+ serde_json::json!({"state": if clock() >= window.expires_at {"expired"} else {"started"},
+ "started_at":window.started_at,"expires_at":window.expires_at})
+ };
+ Ok(build_response(
+ StatusCode::OK,
+ "application/json",
+ Body::from(serde_json::to_vec(&result)?),
+ ))
+ }
+
pub(super) async fn handle_registered_rental(
&self,
path: &str,
@@ -102,7 +220,7 @@ impl ApiHandler {
let selected = content.to_owned();
let key = identity.clone();
let metadata = tokio::task::spawn_blocking(move || {
- crate::registered_media::registered_terms(&data, &key, &selected)
+ crate::registered_media::registered_metadata(&data, &key, &selected)
})
.await?;
let (receipt, _) = match metadata {
@@ -160,11 +278,9 @@ async fn rental_response(
let started = opened.started_at;
let expires = opened.expires_at;
let (start, length) = range.map_or((0, total), |(start, end)| (start, end - start + 1));
- let mut file = tokio::fs::File::from_std(opened.file);
- file.seek(std::io::SeekFrom::Start(start)).await?;
let chunks = futures_util::stream::try_unfold(
- (file, length, now),
- move |(mut file, left, now)| async move {
+ (opened.file, opened.verification, start, length, now),
+ move |(mut file, verification, position, left, now)| async move {
if left == 0 {
return Ok::<_, std::io::Error>(None);
}
@@ -175,15 +291,20 @@ async fn rental_response(
"Rental window ended",
));
}
- let mut bytes = vec![0; left.min(64 * 1024) as usize];
- let count = tokio::time::timeout(
- std::time::Duration::from_secs(expires - instant),
- file.read(&mut bytes),
- )
- .await
- .map_err(|_| {
- std::io::Error::new(std::io::ErrorKind::TimedOut, "Rental window ended")
- })??;
+ let read = tokio::task::spawn_blocking(move || {
+ let bytes = verification
+ .index
+ .read_slice(&mut file, position, left.min(64 * 1024) as usize)
+ .map_err(std::io::Error::other)?;
+ Ok::<_, std::io::Error>((file, verification, bytes))
+ });
+ let (file, verification, bytes) =
+ tokio::time::timeout(std::time::Duration::from_secs(expires - instant), read)
+ .await
+ .map_err(|_| {
+ std::io::Error::new(std::io::ErrorKind::TimedOut, "Rental window ended")
+ })?
+ .map_err(std::io::Error::other)??;
let instant = now();
if instant < started || instant >= expires {
return Err(std::io::Error::new(
@@ -191,14 +312,11 @@ async fn rental_response(
"Rental window ended",
));
}
- if count == 0 {
- return Err(std::io::Error::new(
- std::io::ErrorKind::UnexpectedEof,
- "Registered snapshot ended early",
- ));
- }
- bytes.truncate(count);
- Ok(Some((bytes, (file, left - count as u64, now))))
+ let count = bytes.len() as u64;
+ Ok(Some((
+ bytes,
+ (file, verification, position + count, left - count, now),
+ )))
},
);
let mut response = Response::builder()
@@ -257,10 +375,19 @@ mod tests {
);
}
fn opened(size: u64) -> OpenedMedia {
- let file = tempfile::tempfile().unwrap();
+ use sha2::{Digest, Sha256};
+ let mut file = tempfile::tempfile().unwrap();
file.set_len(size).unwrap();
+ let binding = crate::rental_chunk_index::Binding {
+ content_id: format!("registered_{}", uuid::Uuid::new_v4()),
+ receipt_sha256: "ab".repeat(32),
+ full_sha256: hex::encode(Sha256::digest(vec![0; size as usize])),
+ size,
+ };
+ let index = crate::rental_chunk_index::Index::scan(&mut file, binding, |_| Ok(())).unwrap();
OpenedMedia {
file,
+ verification: crate::rental_readiness::Ready::fixture(index),
size_bytes: size,
mime_type: "video/mp4".into(),
started_at: 1000,
diff --git a/core/archipelago/src/api/rpc/dispatcher.rs b/core/archipelago/src/api/rpc/dispatcher.rs
index e04bf941..f7f227f5 100644
--- a/core/archipelago/src/api/rpc/dispatcher.rs
+++ b/core/archipelago/src/api/rpc/dispatcher.rs
@@ -333,6 +333,8 @@ impl RpcHandler {
"content.indeehub-projects" => self.handle_content_indeehub_projects().await,
"content.browse-all-peers" => self.handle_content_browse_all_peers().await,
"content.playback-handle" => self.handle_playback_handle(params, session_token).await,
+ "content.playback-prepare" => self.handle_playback_prepare(params, session_token).await,
+ "content.playback-start" => self.handle_playback_start(params, session_token).await,
"content.playback-status" => self.handle_playback_status(params, session_token).await,
"content.rental-purchase" => self.handle_content_rental_purchase(params).await,
"content.purchase" => self.handle_content_purchase(params).await,
diff --git a/core/archipelago/src/api/rpc/mod.rs b/core/archipelago/src/api/rpc/mod.rs
index a576b095..e5c04f17 100644
--- a/core/archipelago/src/api/rpc/mod.rs
+++ b/core/archipelago/src/api/rpc/mod.rs
@@ -20,8 +20,8 @@ mod interfaces;
pub(crate) mod lnd;
mod marketplace;
mod media_registration;
-mod purchase;
mod playback;
+mod purchase;
// pub(crate): 13-10's `assistant::backends::select_backend` reuses
// `mesh::assistant::detect_ollama()` (D-04) rather than re-probing —
// matches the existing `pub(crate) mod bitcoin_relay;`/`pub(crate) mod
@@ -113,6 +113,8 @@ fn native_consent_origin_allowed(method: &str, headers: &hyper::HeaderMap, dev_m
| "content.cancel-purchase"
| "content.playback-handle"
| "content.playback-status"
+ | "content.playback-prepare"
+ | "content.playback-start"
) || nostr_signing_origin_allowed(headers, dev_mode)
}
@@ -831,6 +833,8 @@ mod nostr_signing_origin_tests {
"content.cancel-purchase",
"content.playback-handle",
"content.playback-status",
+ "content.playback-prepare",
+ "content.playback-start",
] {
assert!(!native_consent_origin_allowed(
method,
diff --git a/core/archipelago/src/api/rpc/playback.rs b/core/archipelago/src/api/rpc/playback.rs
index 010da04e..d8d69ceb 100644
--- a/core/archipelago/src/api/rpc/playback.rs
+++ b/core/archipelago/src/api/rpc/playback.rs
@@ -55,9 +55,7 @@ impl RpcHandler {
let expires = self
.playback_handles()
.expiry(&handle, &session, &context)?;
- Ok(
- serde_json::json!({"playback_url":format!("/api/rental-playback/{handle}"), "expires_at":expires}),
- )
+ Ok(serde_json::json!({"handle":handle,"expires_at":expires}))
}
pub(super) async fn handle_playback_status(
&self,
@@ -72,3 +70,179 @@ impl RpcHandler {
Ok(serde_json::json!({"expires_at":expires}))
}
}
+
+#[derive(Deserialize)]
+#[serde(deny_unknown_fields)]
+struct Prepare {
+ handle: String,
+ #[serde(default)]
+ retry: bool,
+}
+#[derive(Deserialize)]
+#[serde(deny_unknown_fields)]
+struct Start {
+ handle: String,
+ ready_id: String,
+}
+impl RpcHandler {
+ async fn playback_control(
+ &self,
+ handle: &str,
+ ready_id: Option<&str>,
+ retry: bool,
+ session: &Option,
+ ) -> Result {
+ let (session, context) = self.playback_context(session).await?;
+ let binding = self.playback_handles().lookup(handle, &session, &context)?;
+ let (capability, duration) = {
+ let journal = crate::content_purchase::Journal::open(&self.config.data_dir).await?;
+ let record = journal
+ .buyer(&binding.contract.id)
+ .await?
+ .context("Original purchase missing")?;
+ anyhow::ensure!(
+ record.contract == binding.contract,
+ "Original purchase changed"
+ );
+ let capability = record
+ .receipt()
+ .context("Original payment is not settled")?
+ .capability
+ .clone();
+ let envelope = journal
+ .protocol_envelope("buyer", &binding.contract.id)
+ .await?
+ .context("Original rental terms missing")?;
+ (
+ capability,
+ envelope
+ .offer
+ .viewing_seconds
+ .context("Purchase is not a timed rental")?,
+ )
+ };
+ let peer = crate::federation::load_unique_payment_peer(
+ &self.config.data_dir,
+ &binding.seller_onion,
+ )
+ .await?;
+ anyhow::ensure!(
+ peer.did == binding.contract.seller_did,
+ "Purchased seller identity changed"
+ );
+ let mesh = peer.fips_npub.context("Seller mesh binding unavailable")?;
+ let action = if ready_id.is_some() {
+ "start"
+ } else {
+ "prepare"
+ };
+ let path = format!(
+ "/content/{}/rental/{}/{action}",
+ binding.contract.content_id, binding.contract.id
+ );
+ let body = serde_json::json!({"capability":capability,"ready_id":ready_id,"retry":retry});
+ let (mut response, _) =
+ crate::fips::dial::PeerRequest::new(Some(&mesh), &binding.seller_onion, &path)
+ .require_fips()
+ .single_delivery()
+ .timeout(std::time::Duration::from_secs(20))
+ .send_content_json(&self.config.data_dir, &binding.contract.seller_did, &body)
+ .await?;
+ anyhow::ensure!(
+ response.status().is_success(),
+ "Original rental is unavailable; retry this purchase without paying again"
+ );
+ let mut bytes = Vec::new();
+ while let Some(chunk) = response.chunk().await? {
+ anyhow::ensure!(
+ bytes
+ .len()
+ .checked_add(chunk.len())
+ .is_some_and(|n| n <= 16 * 1024),
+ "Rental control response too large"
+ );
+ bytes.extend_from_slice(&chunk);
+ }
+ let remote: serde_json::Value = serde_json::from_slice(&bytes)?;
+ let state = remote["state"].as_str().context("Missing rental state")?;
+ let window = match (remote["started_at"].as_u64(), remote["expires_at"].as_u64()) {
+ (None, None) if remote["started_at"].is_null() && remote["expires_at"].is_null() => {
+ None
+ }
+ (Some(started), Some(expires))
+ if started > 0 && started.checked_add(duration) == Some(expires) =>
+ {
+ Some((started, expires))
+ }
+ _ => anyhow::bail!("Seller changed the original rental window"),
+ };
+ let (started_at, expires_at) =
+ window.map_or((None, None), |(start, end)| (Some(start), Some(end)));
+ let result = match state {
+ "preparing" if ready_id.is_none() => {
+ let completed = remote["completed_bytes"]
+ .as_u64()
+ .context("Invalid verification progress")?;
+ anyhow::ensure!(
+ remote["total_bytes"].as_u64() == Some(binding.contract.content_size)
+ && completed <= binding.contract.content_size
+ && remote["viewing_seconds"].as_u64() == Some(duration),
+ "Rental preparation terms changed"
+ );
+ serde_json::json!({"state":"preparing","completed_bytes":completed,"total_bytes":binding.contract.content_size,
+ "viewing_seconds":duration,"started_at":started_at,"expires_at":expires_at})
+ }
+ "ready" if ready_id.is_none() => {
+ let id = remote["ready_id"]
+ .as_str()
+ .context("Missing readiness identifier")?;
+ let parsed = uuid::Uuid::parse_str(id)?;
+ anyhow::ensure!(
+ parsed.to_string() == id
+ && parsed.get_version_num() == 4
+ && remote["total_bytes"].as_u64() == Some(binding.contract.content_size)
+ && remote["viewing_seconds"].as_u64() == Some(duration),
+ "Rental readiness changed"
+ );
+ serde_json::json!({"state":"ready","ready_id":id,"handle":handle,"viewing_seconds":duration,
+ "started_at":started_at,"expires_at":expires_at})
+ }
+ "unavailable" if ready_id.is_none() => {
+ serde_json::json!({"state":"unavailable","expires_at":expires_at})
+ }
+ "expired" if window.is_some() => {
+ serde_json::json!({"state":"expired","started_at":started_at,"expires_at":expires_at})
+ }
+ "started" if ready_id.is_some() && window.is_some() => {
+ serde_json::json!({"state":"started","started_at":started_at,
+ "expires_at":expires_at,"playback_url":format!("/api/rental-playback/{handle}")})
+ }
+ _ => {
+ anyhow::bail!("Unexpected rental state; recover this purchase without paying again")
+ }
+ };
+ if let Some((_, expires)) = window {
+ self.playback_handles()
+ .note_expiry(handle, &binding, expires)?;
+ }
+ Ok(result)
+ }
+ pub(super) async fn handle_playback_prepare(
+ &self,
+ params: Option,
+ session: &Option,
+ ) -> Result {
+ let input: Prepare = serde_json::from_value(params.context("Missing playback handle")?)?;
+ self.playback_control(&input.handle, None, input.retry, session)
+ .await
+ }
+ pub(super) async fn handle_playback_start(
+ &self,
+ params: Option,
+ session: &Option,
+ ) -> Result {
+ let input: Start = serde_json::from_value(params.context("Missing readiness identifier")?)?;
+ self.playback_control(&input.handle, Some(&input.ready_id), false, session)
+ .await
+ }
+}
diff --git a/core/archipelago/src/api/rpc/purchase.rs b/core/archipelago/src/api/rpc/purchase.rs
index 1a613d47..da51aea9 100644
--- a/core/archipelago/src/api/rpc/purchase.rs
+++ b/core/archipelago/src/api/rpc/purchase.rs
@@ -51,6 +51,9 @@ impl RpcHandler {
)
.await?;
match result {
+ ReadyPurchase::Preparing { .. } => {
+ anyhow::bail!("Ordinary purchase cannot prepare a timed rental")
+ }
ReadyPurchase::AwaitingConfirmation {
operation_id,
envelope_sha256,
@@ -127,6 +130,8 @@ impl RpcHandler {
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct RentalParams {
+ #[serde(default)]
+ retry_preparation: bool,
seller_did: String,
content_id: String,
expected_sha256: String,
@@ -168,8 +173,9 @@ impl RpcHandler {
¶ms.seller_did,
)
.await?;
- let transport =
- FipsPurchaseTransport::load(self.config.data_dir.clone(), onion.clone()).await?;
+ let transport = FipsPurchaseTransport::load(self.config.data_dir.clone(), onion.clone())
+ .await?
+ .retry_preparation(params.retry_preparation);
let expected = caller::ExpectedRental {
seller_did: params.seller_did,
content_id: params.content_id.clone(),
@@ -189,6 +195,11 @@ impl RpcHandler {
)
.await?
{
+ ReadyPurchase::Preparing {
+ completed_bytes,
+ total_bytes,
+ } => Ok(serde_json::json!({
+ "state":"preparing", "completed_bytes":completed_bytes,"total_bytes":total_bytes})),
ReadyPurchase::AwaitingConfirmation {
operation_id,
envelope_sha256,
diff --git a/core/archipelago/src/content_purchase_caller.rs b/core/archipelago/src/content_purchase_caller.rs
index 09d85154..d8249b62 100644
--- a/core/archipelago/src/content_purchase_caller.rs
+++ b/core/archipelago/src/content_purchase_caller.rs
@@ -12,6 +12,12 @@ pub(crate) trait PurchaseTransport: Send + Sync {
/// This identity is the independently verified peer binding, not response JSON.
fn seller_did(&self) -> &str;
fn seller_onion(&self) -> &str;
+ fn prepare_offer(
+ &self,
+ _content_id: &str,
+ ) -> impl Future