Fix paid-file recovery, app lifecycle regressions and wallet controls
Demo images / Build & push demo images (push) Failing after 1m10s
Demo images / Build & push demo images (push) Failing after 1m10s
This commit is contained in:
@@ -74,7 +74,7 @@ impl ApiHandler {
|
||||
let invoice_hash = headers
|
||||
.get("x-invoice-hash")
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.map(|s| s.to_string())
|
||||
.map(|s| s.to_ascii_lowercase())
|
||||
.or_else(|| {
|
||||
headers
|
||||
.get("x-onchain-address")
|
||||
@@ -98,6 +98,46 @@ impl ApiHandler {
|
||||
None => false,
|
||||
};
|
||||
|
||||
// Payment settlement is verified on the seller even when no status
|
||||
// poll preceded this download (e.g. direct payment from another node).
|
||||
let requires_payment = if !owner_session && headers.contains_key("x-invoice-hash") {
|
||||
content_server::load_catalog(&config.data_dir)
|
||||
.await?
|
||||
.items
|
||||
.iter()
|
||||
.any(|item| {
|
||||
item.id == content_id
|
||||
&& matches!(item.access, content_server::AccessControl::Paid { .. })
|
||||
})
|
||||
} else {
|
||||
false
|
||||
};
|
||||
if requires_payment {
|
||||
if let Some(hash) = headers.get("x-invoice-hash").and_then(|v| v.to_str().ok()) {
|
||||
if hash.len() != 64 || !hash.bytes().all(|c| c.is_ascii_hexdigit()) {
|
||||
return Ok(build_response(
|
||||
StatusCode::BAD_REQUEST,
|
||||
"text/plain",
|
||||
hyper::Body::from("Invalid payment hash"),
|
||||
));
|
||||
}
|
||||
if let Err(error) = self
|
||||
.rpc_handler
|
||||
.settle_content_invoice(hash, content_id)
|
||||
.await
|
||||
{
|
||||
tracing::warn!("Cannot verify peer-file invoice settlement: {error:#}");
|
||||
return Ok(build_response(
|
||||
StatusCode::SERVICE_UNAVAILABLE,
|
||||
"application/json",
|
||||
hyper::Body::from(
|
||||
r#"{"error":"Payment verification is temporarily unavailable. Retry the download without paying again."}"#,
|
||||
),
|
||||
));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Parse Range header for streaming support
|
||||
let range = headers
|
||||
.get("range")
|
||||
@@ -249,7 +289,13 @@ impl ApiHandler {
|
||||
.await
|
||||
{
|
||||
Ok((bolt11, payment_hash)) if !payment_hash.is_empty() => {
|
||||
crate::content_invoice::record_pending(&payment_hash, content_id, price_sats).await;
|
||||
crate::content_invoice::record_pending(
|
||||
&self.config.data_dir,
|
||||
&payment_hash,
|
||||
content_id,
|
||||
price_sats,
|
||||
)
|
||||
.await?;
|
||||
let body = serde_json::json!({
|
||||
"bolt11": bolt11,
|
||||
"payment_hash": payment_hash,
|
||||
@@ -309,26 +355,10 @@ impl ApiHandler {
|
||||
));
|
||||
}
|
||||
|
||||
// The hash must be one we issued for exactly this content item.
|
||||
match crate::content_invoice::lookup(payment_hash).await {
|
||||
Some((cid, _)) if cid == content_id => {}
|
||||
_ => {
|
||||
return Ok(build_response(
|
||||
StatusCode::NOT_FOUND,
|
||||
"application/json",
|
||||
hyper::Body::from(r#"{"error":"Unknown invoice"}"#),
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
// Already paid? Otherwise ask our LND and persist the result.
|
||||
let mut paid = crate::content_invoice::is_paid_for(payment_hash, content_id).await;
|
||||
if !paid {
|
||||
if let Ok(true) = self.rpc_handler.invoice_is_settled(payment_hash).await {
|
||||
crate::content_invoice::mark_paid(payment_hash).await;
|
||||
paid = true;
|
||||
}
|
||||
}
|
||||
let paid = self
|
||||
.rpc_handler
|
||||
.settle_content_invoice(payment_hash, content_id)
|
||||
.await?;
|
||||
|
||||
let body = serde_json::json!({ "paid": paid });
|
||||
Ok(build_response(
|
||||
@@ -389,7 +419,13 @@ impl ApiHandler {
|
||||
|
||||
match self.rpc_handler.new_onchain_address().await {
|
||||
Ok(address) if !address.is_empty() => {
|
||||
crate::content_invoice::record_pending(&address, content_id, price_sats).await;
|
||||
crate::content_invoice::record_pending(
|
||||
&self.config.data_dir,
|
||||
&address,
|
||||
content_id,
|
||||
price_sats,
|
||||
)
|
||||
.await?;
|
||||
let body = serde_json::json!({
|
||||
"address": address,
|
||||
"amount_sats": price_sats,
|
||||
@@ -439,7 +475,7 @@ impl ApiHandler {
|
||||
));
|
||||
}
|
||||
// The address must be one we issued for exactly this content item.
|
||||
let price = match crate::content_invoice::lookup(address).await {
|
||||
let price = match crate::content_invoice::lookup(&self.config.data_dir, address).await? {
|
||||
Some((cid, price)) if cid == content_id => price,
|
||||
_ => {
|
||||
return Ok(build_response(
|
||||
@@ -450,10 +486,11 @@ impl ApiHandler {
|
||||
}
|
||||
};
|
||||
|
||||
let mut paid = crate::content_invoice::is_paid_for(address, content_id).await;
|
||||
let mut paid =
|
||||
crate::content_invoice::is_paid_for(&self.config.data_dir, address, content_id).await;
|
||||
if !paid {
|
||||
if let Ok(true) = self.rpc_handler.onchain_received(address, price).await {
|
||||
crate::content_invoice::mark_paid(address).await;
|
||||
crate::content_invoice::mark_paid(&self.config.data_dir, address).await?;
|
||||
paid = true;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -73,6 +73,17 @@ fn paid_content_response(bytes: &[u8], mime: &str, paid_sats: u64) -> serde_json
|
||||
})
|
||||
}
|
||||
|
||||
// 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)
|
||||
}
|
||||
}
|
||||
|
||||
/// File purchases through an atomic no-clobber write in Files' own namespace.
|
||||
async fn file_purchase_in_files(
|
||||
data_dir: &std::path::Path,
|
||||
@@ -870,10 +881,29 @@ impl RpcHandler {
|
||||
if !is_valid_v3_onion(onion) {
|
||||
return Err(anyhow::anyhow!("Invalid v3 onion address"));
|
||||
}
|
||||
if payment_hash.is_empty() || !payment_hash.chars().all(|c| c.is_ascii_hexdigit()) {
|
||||
if payment_hash.len() != 64 || !payment_hash.chars().all(|c| c.is_ascii_hexdigit()) {
|
||||
return Err(anyhow::anyhow!("Invalid payment_hash"));
|
||||
}
|
||||
|
||||
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
|
||||
{
|
||||
return Ok(invoice_download_response(&bytes, &mime, cache_only));
|
||||
}
|
||||
// Older sellers only mark settlement during status polling. Always
|
||||
// perform that handshake before requesting bytes; retries never pay.
|
||||
// The download gate remains authoritative: a file may have become
|
||||
// free, and newer sellers verify directly if status polling fails.
|
||||
let _ = self
|
||||
.handle_content_invoice_status(Some(serde_json::json!({
|
||||
"onion": onion, "content_id": content_id, "payment_hash": payment_hash,
|
||||
})))
|
||||
.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;
|
||||
@@ -912,7 +942,7 @@ impl RpcHandler {
|
||||
|
||||
if response.status() == reqwest::StatusCode::PAYMENT_REQUIRED {
|
||||
return Ok(serde_json::json!({
|
||||
"error": "Seller has not registered this payment yet — wait for settlement and retry."
|
||||
"error": "The seller has not confirmed access yet. Retry the download without paying again."
|
||||
}));
|
||||
}
|
||||
if !response.status().is_success() {
|
||||
@@ -921,16 +951,45 @@ impl RpcHandler {
|
||||
}));
|
||||
}
|
||||
|
||||
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 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(),
|
||||
}))
|
||||
.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(
|
||||
&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(),
|
||||
)
|
||||
.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
|
||||
{
|
||||
tracing::warn!("Lightning purchase cached; optional Files copy failed: {error:#}");
|
||||
}
|
||||
Ok(invoice_download_response(&bytes, &mime, cache_only))
|
||||
}
|
||||
|
||||
/// Buyer side (#46): ask the seller for a fresh on-chain address to pay.
|
||||
@@ -1405,3 +1464,19 @@ impl RpcHandler {
|
||||
#[cfg(test)]
|
||||
#[path = "content_tests.rs"]
|
||||
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);
|
||||
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);
|
||||
assert_eq!(legacy["data"], "cGFpZCBieXRlcw==");
|
||||
assert_eq!(legacy["data"], legacy["data_base64"]);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -473,6 +473,7 @@ impl RpcHandler {
|
||||
));
|
||||
}
|
||||
|
||||
let fee_query = close_channel_fee_query(¶ms)?;
|
||||
let force = params
|
||||
.get("force")
|
||||
.and_then(|v| v.as_bool())
|
||||
@@ -498,13 +499,11 @@ impl RpcHandler {
|
||||
.build()
|
||||
.context("Failed to create streaming HTTP client")?;
|
||||
|
||||
let url = format!(
|
||||
"{LND_REST_BASE_URL}/v1/channels/{}/{}?force={}",
|
||||
parts[0], parts[1], force
|
||||
);
|
||||
let url = format!("{LND_REST_BASE_URL}/v1/channels/{}/{}", parts[0], parts[1]);
|
||||
|
||||
let mut resp = client
|
||||
.delete(&url)
|
||||
.query(&fee_query)
|
||||
.header("Grpc-Metadata-macaroon", &macaroon_hex)
|
||||
.send()
|
||||
.await
|
||||
@@ -572,3 +571,101 @@ impl RpcHandler {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// LND's CloseChannel REST endpoint takes fee selection as query parameters.
|
||||
/// With neither parameter LND uses a lax target; keep legacy clients on our
|
||||
/// explicit Standard target rather than silently accepting that default.
|
||||
fn close_channel_fee_query(params: &serde_json::Value) -> Result<Vec<(&'static str, String)>> {
|
||||
let force = match params.get("force") {
|
||||
None | Some(serde_json::Value::Null) => false,
|
||||
Some(value) => value
|
||||
.as_bool()
|
||||
.ok_or_else(|| anyhow::anyhow!("force must be a boolean"))?,
|
||||
};
|
||||
let integer = |key: &str, max: u64| -> Result<Option<u64>> {
|
||||
match params.get(key) {
|
||||
None | Some(serde_json::Value::Null) => Ok(None),
|
||||
Some(value) => {
|
||||
let n = value
|
||||
.as_u64()
|
||||
.ok_or_else(|| anyhow::anyhow!("{key} must be a positive whole number"))?;
|
||||
anyhow::ensure!((1..=max).contains(&n), "{key} must be between 1 and {max}");
|
||||
Ok(Some(n))
|
||||
}
|
||||
}
|
||||
};
|
||||
let target = integer("target_conf", 1008)?;
|
||||
let rate = integer("sat_per_vbyte", 5000)?;
|
||||
anyhow::ensure!(
|
||||
target.is_none() || rate.is_none(),
|
||||
"Specify either target_conf or sat_per_vbyte, not both"
|
||||
);
|
||||
anyhow::ensure!(
|
||||
!force || (target.is_none() && rate.is_none()),
|
||||
"Closing fee selection requires a cooperative close"
|
||||
);
|
||||
let mut query = vec![("force", force.to_string())];
|
||||
if !force {
|
||||
if let Some(rate) = rate {
|
||||
query.push(("sat_per_vbyte", rate.to_string()));
|
||||
} else {
|
||||
query.push(("target_conf", target.unwrap_or(6).to_string()));
|
||||
}
|
||||
}
|
||||
Ok(query)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod close_fee_tests {
|
||||
use super::*;
|
||||
#[test]
|
||||
fn close_fee_query_forwards_presets_custom_and_legacy_default() {
|
||||
for target in [1, 3, 6, 1008] {
|
||||
assert_eq!(
|
||||
close_channel_fee_query(&serde_json::json!({"target_conf":target})).unwrap(),
|
||||
vec![
|
||||
("force", "false".into()),
|
||||
("target_conf", target.to_string())
|
||||
]
|
||||
);
|
||||
}
|
||||
for rate in [1, 25, 5000] {
|
||||
let query =
|
||||
close_channel_fee_query(&serde_json::json!({"sat_per_vbyte":rate})).unwrap();
|
||||
let request = reqwest::Client::new()
|
||||
.delete("http://localhost/v1/channels/test/0")
|
||||
.query(&query)
|
||||
.build()
|
||||
.unwrap();
|
||||
assert_eq!(request.method(), reqwest::Method::DELETE);
|
||||
assert_eq!(
|
||||
request.url().query(),
|
||||
Some(format!("force=false&sat_per_vbyte={rate}").as_str())
|
||||
);
|
||||
}
|
||||
assert_eq!(
|
||||
close_channel_fee_query(&serde_json::json!({})).unwrap(),
|
||||
vec![("force", "false".into()), ("target_conf", "6".into())]
|
||||
);
|
||||
assert_eq!(
|
||||
close_channel_fee_query(&serde_json::json!({"force":true})).unwrap(),
|
||||
vec![("force", "true".into())]
|
||||
);
|
||||
}
|
||||
#[test]
|
||||
fn malformed_or_conflicting_close_fees_fail_before_wallet_access() {
|
||||
for params in [
|
||||
serde_json::json!({"target_conf":1,"sat_per_vbyte":2}),
|
||||
serde_json::json!({"force":true,"target_conf":1}),
|
||||
serde_json::json!({"force":"false"}),
|
||||
serde_json::json!({"target_conf":0}),
|
||||
serde_json::json!({"target_conf":1009}),
|
||||
serde_json::json!({"sat_per_vbyte":5001}),
|
||||
serde_json::json!({"sat_per_vbyte":-1}),
|
||||
serde_json::json!({"sat_per_vbyte":1.5}),
|
||||
serde_json::json!({"sat_per_vbyte":"25"}),
|
||||
] {
|
||||
assert!(close_channel_fee_query(¶ms).is_err(), "{params}");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -453,6 +453,56 @@ impl RpcHandler {
|
||||
Ok(settled)
|
||||
}
|
||||
|
||||
/// Verify against LND at download time, rather than relying on a browser
|
||||
/// having polled first. The memo/amount also recover pre-upgrade in-memory
|
||||
/// entitlements after restart; unrelated invoices never unlock a file.
|
||||
pub(crate) async fn settle_content_invoice(
|
||||
&self,
|
||||
hash: &str,
|
||||
content_id: &str,
|
||||
) -> Result<bool> {
|
||||
anyhow::ensure!(
|
||||
hash.len() == 64 && hash.bytes().all(|c| c.is_ascii_hexdigit()),
|
||||
"Invalid payment hash"
|
||||
);
|
||||
let hash = hash.to_ascii_lowercase();
|
||||
let existing = crate::content_invoice::lookup(&self.config.data_dir, &hash).await?;
|
||||
if let Some((id, _)) = &existing {
|
||||
if id != content_id {
|
||||
return Ok(false);
|
||||
}
|
||||
}
|
||||
if crate::content_invoice::is_paid_for(&self.config.data_dir, &hash, content_id).await {
|
||||
return Ok(true);
|
||||
}
|
||||
let (client, macaroon_hex) = self.lnd_client().await?;
|
||||
let response = client
|
||||
.get(format!("{LND_REST_BASE_URL}/v1/invoice/{hash}"))
|
||||
.header("Grpc-Metadata-macaroon", &macaroon_hex)
|
||||
.send()
|
||||
.await?;
|
||||
if response.status() == reqwest::StatusCode::NOT_FOUND {
|
||||
return Ok(false);
|
||||
}
|
||||
let body: serde_json::Value = response.error_for_status()?.json().await?;
|
||||
let Some(price) = content_invoice_amount(&body, content_id) else {
|
||||
return Ok(false);
|
||||
};
|
||||
if existing
|
||||
.as_ref()
|
||||
.is_some_and(|(_, expected)| *expected != price)
|
||||
{
|
||||
return Ok(false);
|
||||
}
|
||||
crate::content_invoice::record_pending(&self.config.data_dir, &hash, content_id, price)
|
||||
.await?;
|
||||
let settled = content_invoice_fully_settled(&body, price);
|
||||
if settled {
|
||||
crate::content_invoice::mark_paid(&self.config.data_dir, &hash).await?;
|
||||
}
|
||||
Ok(settled)
|
||||
}
|
||||
|
||||
/// Generate a fresh on-chain receive address (seller side, #46).
|
||||
pub(crate) async fn new_onchain_address(&self) -> Result<String> {
|
||||
let (client, macaroon_hex) = self.lnd_client().await?;
|
||||
@@ -1444,3 +1494,81 @@ mod tests {
|
||||
assert!(s.contains("[LND_REST_UNREACHABLE]"), "got: {s}");
|
||||
}
|
||||
}
|
||||
|
||||
// LND REST uses decimal strings for int64 fields. Match the complete seller
|
||||
// memo, not a substring supplied by a buyer or an arbitrary settled invoice.
|
||||
fn json_u64(value: &serde_json::Value) -> Option<u64> {
|
||||
value.as_u64().or_else(|| value.as_str()?.parse().ok())
|
||||
}
|
||||
fn content_invoice_fully_settled(body: &serde_json::Value, price: u64) -> bool {
|
||||
let settled = match body.get("state").and_then(|v| v.as_str()) {
|
||||
Some(state) => state == "SETTLED",
|
||||
None => body.get("settled").and_then(|v| v.as_bool()) == Some(true),
|
||||
};
|
||||
settled
|
||||
&& price > 0
|
||||
&& body
|
||||
.get("amt_paid_sat")
|
||||
.and_then(json_u64)
|
||||
.is_some_and(|paid| paid >= price)
|
||||
}
|
||||
fn content_invoice_amount(body: &serde_json::Value, content_id: &str) -> Option<u64> {
|
||||
if body.get("memo")?.as_str()? != format!("Archipelago peer file {content_id}") {
|
||||
return None;
|
||||
}
|
||||
body.get("value").and_then(json_u64).filter(|v| *v > 0)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod peer_file_invoice_tests {
|
||||
use super::*;
|
||||
#[test]
|
||||
fn settlement_requires_terminal_state_and_full_amount() {
|
||||
for state in ["OPEN", "ACCEPTED", "CANCELED", "unknown"] {
|
||||
assert!(!content_invoice_fully_settled(
|
||||
&serde_json::json!({"state":state,"settled":true,"amt_paid_sat":"100"}),
|
||||
7
|
||||
));
|
||||
}
|
||||
for amount in [
|
||||
serde_json::json!(6),
|
||||
serde_json::json!("-1"),
|
||||
serde_json::json!(null),
|
||||
serde_json::json!("bad"),
|
||||
] {
|
||||
assert!(!content_invoice_fully_settled(
|
||||
&serde_json::json!({"state":"SETTLED","amt_paid_sat":amount}),
|
||||
7
|
||||
));
|
||||
}
|
||||
for amount in [serde_json::json!(7), serde_json::json!("8")] {
|
||||
assert!(content_invoice_fully_settled(
|
||||
&serde_json::json!({"state":"SETTLED","amt_paid_sat":amount}),
|
||||
7
|
||||
));
|
||||
}
|
||||
assert!(content_invoice_fully_settled(
|
||||
&serde_json::json!({"settled":true,"amt_paid_sat":"7"}),
|
||||
7
|
||||
));
|
||||
assert!(!content_invoice_fully_settled(
|
||||
&serde_json::json!({"state":"SETTLED","amt_paid_sat":"7"}),
|
||||
0
|
||||
));
|
||||
}
|
||||
#[test]
|
||||
fn legacy_recovery_requires_exact_file_memo_and_positive_amount() {
|
||||
let invoice = serde_json::json!({"memo":"Archipelago peer file file-1", "value":"7"});
|
||||
assert_eq!(content_invoice_amount(&invoice, "file-1"), Some(7));
|
||||
assert_eq!(content_invoice_amount(&invoice, "file-2"), None);
|
||||
for value in [
|
||||
serde_json::json!("-1"),
|
||||
serde_json::json!(0),
|
||||
serde_json::json!("bad"),
|
||||
] {
|
||||
let mut invalid = invoice.clone();
|
||||
invalid["value"] = value;
|
||||
assert_eq!(content_invoice_amount(&invalid, "file-1"), None);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -20,7 +20,9 @@ use crate::data_model::{
|
||||
/// stopped-app restoration path in agreement with live-container discovery.
|
||||
fn canonical_package_id(name: &str) -> &str {
|
||||
match name.strip_prefix("archy-").unwrap_or(name) {
|
||||
"immich_server" => "immich",
|
||||
"immich_server" | "immich-server" => "immich",
|
||||
"immich-postgres" => "immich_postgres",
|
||||
"immich-redis" => "immich_redis",
|
||||
"mempool-web" | "mempool-frontend" => "mempool",
|
||||
name => name,
|
||||
}
|
||||
@@ -443,6 +445,19 @@ mod lifecycle_regression_tests {
|
||||
use super::*;
|
||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||
|
||||
#[test]
|
||||
fn immich_dependency_aliases_share_the_hidden_component_ids() {
|
||||
for id in [
|
||||
"immich-postgres",
|
||||
"immich_postgres",
|
||||
"archy-immich-postgres",
|
||||
] {
|
||||
assert_eq!(canonical_package_id(id), "immich_postgres");
|
||||
}
|
||||
assert_eq!(canonical_package_id("immich-redis"), "immich_redis");
|
||||
assert_eq!(canonical_package_id("immich-server"), "immich");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn registry_survives_empty_runtime_and_deduplicates_aliases() {
|
||||
let installed = ["archy-gitea", "gitea", "immich_server", "archy-removed"]
|
||||
|
||||
@@ -3153,7 +3153,12 @@ impl ProdContainerOrchestrator {
|
||||
let restart_for_exec_change = quadlet::exec_changed(&old_body, &new_body);
|
||||
let restart_for_health_change = quadlet::health_cmd_changed(&old_body, &new_body);
|
||||
let restart_for_security_change = quadlet::security_changed(&old_body, &new_body);
|
||||
let restart_for_managed_override =
|
||||
quadlet::redundant_managed_network_override(&unit, &unit_dir)
|
||||
.await?
|
||||
.is_some();
|
||||
let needs_restart = restart_required
|
||||
|| restart_for_managed_override
|
||||
|| restart_for_port_change
|
||||
|| restart_for_network_alias_change
|
||||
|| restart_for_exec_change
|
||||
@@ -4881,6 +4886,26 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
|
||||
/// here (production volumes live under `/var/lib/archipelago` — removal is a
|
||||
/// separate operation owned by the data layer, not this orchestrator).
|
||||
async fn remove(&self, app_id: &str, _preserve_data: bool) -> Result<()> {
|
||||
// A removed catalog entry must remain uninstallable. The RPC caller
|
||||
// still removes legacy containers and persists uninstall intent after
|
||||
// confirming they are gone; do not block it on a missing manifest.
|
||||
if !self.state.read().await.manifests.contains_key(app_id) {
|
||||
anyhow::ensure!(
|
||||
!app_id.is_empty()
|
||||
&& app_id.len() <= 128
|
||||
&& app_id
|
||||
.bytes()
|
||||
.all(|c| c.is_ascii_alphanumeric() || matches!(c, b'-' | b'_')),
|
||||
"Invalid app id"
|
||||
);
|
||||
let lock = self.app_lock(app_id).await;
|
||||
let _guard = lock.lock().await;
|
||||
for name in [app_id.to_string(), format!("archy-{app_id}")] {
|
||||
self.remove_quadlet_unit_if_present(&name).await?;
|
||||
}
|
||||
self.state.write().await.disabled.insert(app_id.to_string());
|
||||
return Ok(());
|
||||
}
|
||||
let lm = self.loaded(app_id).await?;
|
||||
let lock = self.app_lock(app_id).await;
|
||||
let _guard = lock.lock().await;
|
||||
@@ -7435,6 +7460,19 @@ app:
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn removed_catalog_entry_does_not_block_legacy_uninstall() {
|
||||
let rt = Arc::new(MockRuntime::default());
|
||||
let orch = orch_with(rt.clone()).await;
|
||||
orch.remove("cryptpad", true).await.unwrap();
|
||||
assert!(orch.state.read().await.disabled.contains("cryptpad"));
|
||||
assert!(
|
||||
rt.calls().is_empty(),
|
||||
"legacy RPC teardown owns the actual containers"
|
||||
);
|
||||
assert!(orch.remove("../other", true).await.is_err());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn remove_disables_manifest_so_reconcile_does_not_reinstall() {
|
||||
let rt = Arc::new(MockRuntime::default());
|
||||
|
||||
@@ -713,29 +713,76 @@ pub async fn unit_dir() -> Result<PathBuf> {
|
||||
Ok(dir)
|
||||
}
|
||||
|
||||
/// Atomically write `unit` into `dir/<name>.container` if the bytes
|
||||
/// differ from what's already there. Returns true if the file changed.
|
||||
/// The early same-node Portainer repair used a managed Quadlet drop-in. Once
|
||||
/// the manifest supplies slirp, the two Network= entries are additive and
|
||||
/// Podman rejects startup. Retire only that exact redundant managed override;
|
||||
/// arbitrary operator settings must survive reconciliation.
|
||||
pub async fn redundant_managed_network_override(
|
||||
unit: &QuadletUnit,
|
||||
dir: &Path,
|
||||
) -> Result<Option<PathBuf>> {
|
||||
if unit.name != "portainer" || !matches!(unit.network, NetworkMode::Slirp4netns) {
|
||||
return Ok(None);
|
||||
}
|
||||
let path = dir.join("portainer.container.d/archy-same-node-network.conf");
|
||||
let body = match fs::read_to_string(&path).await {
|
||||
Ok(body) => body,
|
||||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
|
||||
Err(error) => return Err(error).context("read managed Portainer network override"),
|
||||
};
|
||||
let lines: Vec<&str> = body
|
||||
.lines()
|
||||
.map(str::trim)
|
||||
.filter(|line| !line.is_empty() && !line.starts_with(['#', ';']))
|
||||
.collect();
|
||||
Ok((lines == ["[Container]", "Network=slirp4netns"]).then_some(path))
|
||||
}
|
||||
|
||||
async fn retire_managed_network_override(path: &Path) -> Result<()> {
|
||||
let backup = path.with_extension("conf.retired");
|
||||
match fs::hard_link(path, &backup).await {
|
||||
Ok(()) => {}
|
||||
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {
|
||||
anyhow::ensure!(
|
||||
fs::read(path).await? == fs::read(&backup).await?,
|
||||
"Existing Portainer override backup differs; preserve both for operator review"
|
||||
);
|
||||
}
|
||||
Err(error) => return Err(error).context("back up managed Portainer network override"),
|
||||
}
|
||||
fs::remove_file(path)
|
||||
.await
|
||||
.context("retire redundant Portainer network override")?;
|
||||
tracing::info!("Retired redundant managed Portainer network override; backup retained");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Atomically write the manifest unit and retire known redundant managed
|
||||
/// overrides. Returns true whenever systemd needs a daemon-reload.
|
||||
pub async fn write_if_changed(unit: &QuadletUnit, dir: &Path) -> Result<bool> {
|
||||
let path = dir.join(unit.unit_filename());
|
||||
let new_bytes = unit.render();
|
||||
|
||||
if let Ok(old) = fs::read_to_string(&path).await {
|
||||
if old == new_bytes {
|
||||
return Ok(false);
|
||||
}
|
||||
let redundant = redundant_managed_network_override(unit, dir).await?;
|
||||
let changed = fs::read_to_string(&path)
|
||||
.await
|
||||
.map(|old| old != new_bytes)
|
||||
.unwrap_or(true);
|
||||
if changed {
|
||||
fs::create_dir_all(dir)
|
||||
.await
|
||||
.with_context(|| format!("create_dir_all {}", dir.display()))?;
|
||||
let tmp = path.with_extension("container.tmp");
|
||||
fs::write(&tmp, new_bytes.as_bytes())
|
||||
.await
|
||||
.with_context(|| format!("write tmp {}", tmp.display()))?;
|
||||
fs::rename(&tmp, &path)
|
||||
.await
|
||||
.with_context(|| format!("rename {} -> {}", tmp.display(), path.display()))?;
|
||||
}
|
||||
|
||||
fs::create_dir_all(dir)
|
||||
.await
|
||||
.with_context(|| format!("create_dir_all {}", dir.display()))?;
|
||||
let tmp = path.with_extension("container.tmp");
|
||||
fs::write(&tmp, new_bytes.as_bytes())
|
||||
.await
|
||||
.with_context(|| format!("write tmp {}", tmp.display()))?;
|
||||
fs::rename(&tmp, &path)
|
||||
.await
|
||||
.with_context(|| format!("rename {} -> {}", tmp.display(), path.display()))?;
|
||||
Ok(true)
|
||||
if let Some(override_path) = &redundant {
|
||||
retire_managed_network_override(override_path).await?;
|
||||
}
|
||||
Ok(changed || redundant.is_some())
|
||||
}
|
||||
|
||||
/// Reload the user systemd manager. Required after any quadlet write
|
||||
@@ -1991,6 +2038,80 @@ app:
|
||||
assert!(!network_aliases_changed(new, new));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn redundant_portainer_override_is_backed_up_and_retired_even_when_base_matches() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let manifest =
|
||||
AppManifest::parse(include_str!("../../../../apps/portainer/manifest.yml")).unwrap();
|
||||
let unit = QuadletUnit::from_manifest(&manifest, "portainer");
|
||||
assert!(write_if_changed(&unit, dir.path()).await.unwrap());
|
||||
let path = dir
|
||||
.path()
|
||||
.join("portainer.container.d/archy-same-node-network.conf");
|
||||
fs::create_dir_all(path.parent().unwrap()).await.unwrap();
|
||||
let old = "[Container]\nNetwork=slirp4netns\n";
|
||||
fs::write(&path, old).await.unwrap();
|
||||
assert!(redundant_managed_network_override(&unit, dir.path())
|
||||
.await
|
||||
.unwrap()
|
||||
.is_some());
|
||||
assert!(write_if_changed(&unit, dir.path()).await.unwrap());
|
||||
assert!(!path.exists());
|
||||
assert_eq!(
|
||||
fs::read_to_string(path.with_extension("conf.retired"))
|
||||
.await
|
||||
.unwrap(),
|
||||
old
|
||||
);
|
||||
assert!(!write_if_changed(&unit, dir.path()).await.unwrap());
|
||||
assert_eq!(
|
||||
fs::read_to_string(dir.path().join("portainer.container"))
|
||||
.await
|
||||
.unwrap()
|
||||
.matches("Network=slirp4netns")
|
||||
.count(),
|
||||
1
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn network_override_migration_preserves_operator_customizations_and_failed_backups() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let manifest =
|
||||
AppManifest::parse(include_str!("../../../../apps/portainer/manifest.yml")).unwrap();
|
||||
let mut unit = QuadletUnit::from_manifest(&manifest, "portainer");
|
||||
let path = dir
|
||||
.path()
|
||||
.join("portainer.container.d/archy-same-node-network.conf");
|
||||
fs::create_dir_all(path.parent().unwrap()).await.unwrap();
|
||||
for custom in [
|
||||
"[Container]\nNetwork=custom-net\n",
|
||||
"[Container]\nNetwork=slirp4netns\nEnvironment=OPERATOR_SETTING=1\n",
|
||||
] {
|
||||
fs::write(&path, custom).await.unwrap();
|
||||
assert!(redundant_managed_network_override(&unit, dir.path())
|
||||
.await
|
||||
.unwrap()
|
||||
.is_none());
|
||||
write_if_changed(&unit, dir.path()).await.unwrap();
|
||||
assert_eq!(fs::read_to_string(&path).await.unwrap(), custom);
|
||||
}
|
||||
fs::write(&path, "[Container]\nNetwork=slirp4netns\n")
|
||||
.await
|
||||
.unwrap();
|
||||
unit.network = NetworkMode::Pasta;
|
||||
assert!(redundant_managed_network_override(&unit, dir.path())
|
||||
.await
|
||||
.unwrap()
|
||||
.is_none());
|
||||
unit.network = NetworkMode::Slirp4netns;
|
||||
fs::write(path.with_extension("conf.retired"), "different backup")
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(write_if_changed(&unit, dir.path()).await.is_err());
|
||||
assert!(path.exists(), "failure must preserve the active override");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn failed_runtime_change_remains_pending_when_unit_already_matches() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
|
||||
@@ -1,80 +1,155 @@
|
||||
//! Seller-side pending entitlements for Lightning-invoice peer-file sales (#46).
|
||||
//!
|
||||
//! When a buyer asks to pay for a paid catalog item with an external wallet (as
|
||||
//! opposed to the local-ecash fast path), the *selling* node mints a Lightning
|
||||
//! invoice on its own LND and records a pending entitlement here, keyed by the
|
||||
//! invoice's payment hash. The buyer pays the invoice from any wallet and polls
|
||||
//! for settlement; once the seller's LND confirms the invoice is settled we mark
|
||||
//! the entitlement paid, and the content gate (`content_server::serve_content`)
|
||||
//! then releases the file to anyone presenting that payment hash.
|
||||
//!
|
||||
//! State is in-memory and bounded by a TTL. If the seller restarts before the
|
||||
//! buyer pays, the buyer simply requests a fresh invoice — no value is lost
|
||||
//! because an unpaid invoice represents no money.
|
||||
//! Durable seller-side entitlements for peer-file invoices and on-chain sales.
|
||||
//! Payment records must outlive browser polling, process restarts and invoice
|
||||
//! expiry: an invoice can settle while the buyer is disconnected.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::sync::LazyLock;
|
||||
use std::time::{Duration, Instant};
|
||||
use tokio::sync::Mutex;
|
||||
use anyhow::{Context, Result};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use sha2::{Digest, Sha256};
|
||||
use std::path::{Path, PathBuf};
|
||||
use tokio::{fs, io::AsyncWriteExt, sync::Mutex};
|
||||
|
||||
/// How long a pending/paid entitlement is retained. Generous enough for a human
|
||||
/// to pay an invoice and download, short enough to keep the map small.
|
||||
const ENTITLEMENT_TTL: Duration = Duration::from_secs(3600); // 1 hour
|
||||
static WRITES: Mutex<()> = Mutex::const_new(());
|
||||
|
||||
#[derive(Clone)]
|
||||
#[derive(Clone, Serialize, Deserialize)]
|
||||
struct Entitlement {
|
||||
content_id: String,
|
||||
price_sats: u64,
|
||||
paid: bool,
|
||||
created_at: Instant,
|
||||
}
|
||||
|
||||
static ENTITLEMENTS: LazyLock<Mutex<HashMap<String, Entitlement>>> =
|
||||
LazyLock::new(|| Mutex::new(HashMap::new()));
|
||||
|
||||
/// Drop expired entries. Caller must hold the lock.
|
||||
fn prune(map: &mut HashMap<String, Entitlement>) {
|
||||
map.retain(|_, e| e.created_at.elapsed() < ENTITLEMENT_TTL);
|
||||
fn path(data_dir: &Path, token: &str) -> PathBuf {
|
||||
data_dir.join("content-entitlements").join(format!(
|
||||
"{}.json",
|
||||
hex::encode(Sha256::digest(token.as_bytes()))
|
||||
))
|
||||
}
|
||||
|
||||
/// Record a freshly-minted invoice as a pending (unpaid) entitlement.
|
||||
pub async fn record_pending(payment_hash: &str, content_id: &str, price_sats: u64) {
|
||||
let mut map = ENTITLEMENTS.lock().await;
|
||||
prune(&mut map);
|
||||
map.insert(
|
||||
payment_hash.to_string(),
|
||||
Entitlement {
|
||||
content_id: content_id.to_string(),
|
||||
price_sats,
|
||||
paid: false,
|
||||
created_at: Instant::now(),
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
/// Mark the entitlement for `payment_hash` paid. No-op if unknown/expired.
|
||||
pub async fn mark_paid(payment_hash: &str) {
|
||||
let mut map = ENTITLEMENTS.lock().await;
|
||||
prune(&mut map);
|
||||
if let Some(e) = map.get_mut(payment_hash) {
|
||||
e.paid = true;
|
||||
async fn read(data_dir: &Path, token: &str) -> Result<Option<Entitlement>> {
|
||||
match fs::read(path(data_dir, token)).await {
|
||||
Ok(bytes) => Ok(Some(
|
||||
serde_json::from_slice(&bytes).context("Invalid payment entitlement")?,
|
||||
)),
|
||||
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
|
||||
Err(e) => Err(e).context("Reading payment entitlement"),
|
||||
}
|
||||
}
|
||||
|
||||
/// The content_id + price an entitlement was issued for, if still live.
|
||||
pub async fn lookup(payment_hash: &str) -> Option<(String, u64)> {
|
||||
let mut map = ENTITLEMENTS.lock().await;
|
||||
prune(&mut map);
|
||||
map.get(payment_hash)
|
||||
.map(|e| (e.content_id.clone(), e.price_sats))
|
||||
async fn write(data_dir: &Path, token: &str, entry: &Entitlement) -> Result<()> {
|
||||
let target = path(data_dir, token);
|
||||
let dir = target.parent().unwrap();
|
||||
fs::create_dir_all(dir).await?;
|
||||
let tmp = target.with_extension("tmp");
|
||||
let mut file = fs::OpenOptions::new()
|
||||
.write(true)
|
||||
.create(true)
|
||||
.truncate(true)
|
||||
.mode(0o600)
|
||||
.open(&tmp)
|
||||
.await?;
|
||||
file.write_all(&serde_json::to_vec(entry)?).await?;
|
||||
file.sync_all().await?;
|
||||
drop(file);
|
||||
fs::rename(&tmp, &target).await?;
|
||||
fs::File::open(dir).await?.sync_all().await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// True if `payment_hash` is a paid entitlement for exactly `content_id`.
|
||||
/// This is the gate the content server consults to release a file.
|
||||
pub async fn is_paid_for(payment_hash: &str, content_id: &str) -> bool {
|
||||
let mut map = ENTITLEMENTS.lock().await;
|
||||
prune(&mut map);
|
||||
map.get(payment_hash)
|
||||
/// Save before exposing an invoice/address to the buyer. Never overwrite an
|
||||
/// existing payment or silently rebind its token to another item or price.
|
||||
pub async fn record_pending(
|
||||
data_dir: &Path,
|
||||
token: &str,
|
||||
content_id: &str,
|
||||
price_sats: u64,
|
||||
) -> Result<()> {
|
||||
let _lock = WRITES.lock().await;
|
||||
if let Some(existing) = read(data_dir, token).await? {
|
||||
anyhow::ensure!(
|
||||
existing.content_id == content_id && existing.price_sats == price_sats,
|
||||
"Payment entitlement mismatch"
|
||||
);
|
||||
return Ok(());
|
||||
}
|
||||
write(
|
||||
data_dir,
|
||||
token,
|
||||
&Entitlement {
|
||||
content_id: content_id.into(),
|
||||
price_sats,
|
||||
paid: false,
|
||||
},
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
pub async fn mark_paid(data_dir: &Path, token: &str) -> Result<()> {
|
||||
let _lock = WRITES.lock().await;
|
||||
let mut entry = read(data_dir, token)
|
||||
.await?
|
||||
.context("Unknown payment entitlement")?;
|
||||
entry.paid = true;
|
||||
write(data_dir, token, &entry).await
|
||||
}
|
||||
|
||||
pub async fn lookup(data_dir: &Path, token: &str) -> Result<Option<(String, u64)>> {
|
||||
Ok(read(data_dir, token)
|
||||
.await?
|
||||
.map(|e| (e.content_id, e.price_sats)))
|
||||
}
|
||||
|
||||
pub async fn is_paid_for(data_dir: &Path, token: &str, content_id: &str) -> bool {
|
||||
read(data_dir, token)
|
||||
.await
|
||||
.ok()
|
||||
.flatten()
|
||||
.map(|e| e.paid && e.content_id == content_id)
|
||||
.unwrap_or(false)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
#[tokio::test]
|
||||
async fn paid_entitlement_survives_reload_and_cannot_be_rebound() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
record_pending(dir.path(), "hash", "file", 12)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(!is_paid_for(dir.path(), "hash", "file").await);
|
||||
mark_paid(dir.path(), "hash").await.unwrap();
|
||||
// All reads reopen disk; no process-local entitlement map exists.
|
||||
assert!(is_paid_for(dir.path(), "hash", "file").await);
|
||||
assert!(!is_paid_for(dir.path(), "hash", "other").await);
|
||||
record_pending(dir.path(), "hash", "file", 12)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(is_paid_for(dir.path(), "hash", "file").await);
|
||||
assert!(record_pending(dir.path(), "hash", "other", 12)
|
||||
.await
|
||||
.is_err());
|
||||
assert!(record_pending(dir.path(), "hash", "file", 13)
|
||||
.await
|
||||
.is_err());
|
||||
let other = tempfile::tempdir().unwrap();
|
||||
assert!(!is_paid_for(other.path(), "hash", "file").await);
|
||||
assert!(mark_paid(dir.path(), "unknown").await.is_err());
|
||||
}
|
||||
#[tokio::test]
|
||||
async fn corrupt_or_unwritable_records_fail_closed() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
record_pending(dir.path(), "../../token", "file", 1)
|
||||
.await
|
||||
.unwrap();
|
||||
fs::write(path(dir.path(), "../../token"), b"broken")
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(lookup(dir.path(), "../../token").await.is_err());
|
||||
assert!(!is_paid_for(dir.path(), "../../token", "file").await);
|
||||
assert!(record_pending(dir.path(), "../../token", "file", 1)
|
||||
.await
|
||||
.is_err());
|
||||
let file = dir.path().join("not-directory");
|
||||
fs::write(&file, b"x").await.unwrap();
|
||||
assert!(record_pending(&file, "hash", "file", 1).await.is_err());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,7 +12,9 @@
|
||||
use anyhow::{Context, Result};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::path::{Path, PathBuf};
|
||||
use tokio::fs;
|
||||
use tokio::{fs, io::AsyncWriteExt, sync::Mutex};
|
||||
|
||||
static PURCHASE_WRITES: Mutex<()> = Mutex::const_new(());
|
||||
|
||||
const OWNED_DIR: &str = "purchased-content";
|
||||
const OWNED_INDEX: &str = "owned.json";
|
||||
@@ -66,20 +68,51 @@ fn bytes_path(data_dir: &Path, onion: &str, content_id: &str) -> PathBuf {
|
||||
.join(sanitize(content_id))
|
||||
}
|
||||
|
||||
async fn load_index(data_dir: &Path) -> OwnedIndex {
|
||||
async fn load_index_checked(data_dir: &Path) -> Result<OwnedIndex> {
|
||||
match fs::read_to_string(index_path(data_dir)).await {
|
||||
Ok(s) => serde_json::from_str(&s).unwrap_or_default(),
|
||||
Err(_) => OwnedIndex::default(),
|
||||
Ok(s) => serde_json::from_str(&s)
|
||||
.context("Invalid purchase index; existing records were preserved"),
|
||||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(OwnedIndex::default()),
|
||||
Err(error) => Err(error).context("Reading purchase index"),
|
||||
}
|
||||
}
|
||||
|
||||
async fn load_index(data_dir: &Path) -> OwnedIndex {
|
||||
load_index_checked(data_dir).await.unwrap_or_default()
|
||||
}
|
||||
|
||||
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 result = async {
|
||||
let mut file = fs::OpenOptions::new()
|
||||
.write(true)
|
||||
.create_new(true)
|
||||
.mode(0o600)
|
||||
.open(&temp)
|
||||
.await?;
|
||||
file.write_all(bytes).await?;
|
||||
file.sync_all().await?;
|
||||
drop(file);
|
||||
fs::rename(&temp, path).await?;
|
||||
fs::File::open(parent).await?.sync_all().await?;
|
||||
Ok::<_, anyhow::Error>(())
|
||||
}
|
||||
.await;
|
||||
if result.is_err() {
|
||||
let _ = fs::remove_file(&temp).await;
|
||||
}
|
||||
result
|
||||
}
|
||||
|
||||
async fn save_index(data_dir: &Path, index: &OwnedIndex) -> Result<()> {
|
||||
let root = owned_root(data_dir);
|
||||
fs::create_dir_all(&root)
|
||||
.await
|
||||
.with_context(|| format!("creating {}", root.display()))?;
|
||||
let content = serde_json::to_string_pretty(index).context("serializing owned index")?;
|
||||
fs::write(index_path(data_dir), content)
|
||||
atomic_write(&index_path(data_dir), content.as_bytes())
|
||||
.await
|
||||
.context("writing owned index")
|
||||
}
|
||||
@@ -98,17 +131,15 @@ pub async fn record_purchase(
|
||||
ecash_backend: &str,
|
||||
purchased_at: &str,
|
||||
) -> Result<()> {
|
||||
// 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;
|
||||
let mut index = load_index_checked(data_dir).await?;
|
||||
let path = bytes_path(data_dir, onion, content_id);
|
||||
if let Some(parent) = path.parent() {
|
||||
fs::create_dir_all(parent)
|
||||
.await
|
||||
.with_context(|| format!("creating {}", parent.display()))?;
|
||||
}
|
||||
fs::write(&path, bytes)
|
||||
atomic_write(&path, bytes)
|
||||
.await
|
||||
.with_context(|| format!("writing purchased bytes to {}", path.display()))?;
|
||||
|
||||
let mut index = load_index(data_dir).await;
|
||||
let entry = OwnedItem {
|
||||
onion: onion.to_string(),
|
||||
content_id: content_id.to_string(),
|
||||
@@ -165,3 +196,90 @@ pub async fn read_owned(
|
||||
.unwrap_or_else(|| "application/octet-stream".to_string());
|
||||
Some((mime, bytes))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
#[tokio::test]
|
||||
async fn concurrent_purchases_preserve_every_item_and_exact_bytes() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let mut jobs = tokio::task::JoinSet::new();
|
||||
for n in 0..24 {
|
||||
let root = dir.path().to_path_buf();
|
||||
jobs.spawn(async move {
|
||||
let id = format!("file-{n}");
|
||||
record_purchase(
|
||||
&root,
|
||||
"seller.onion",
|
||||
&id,
|
||||
&id,
|
||||
"text/plain",
|
||||
id.as_bytes(),
|
||||
5,
|
||||
"lightning",
|
||||
"now",
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
});
|
||||
}
|
||||
while let Some(result) = jobs.join_next().await {
|
||||
result.unwrap();
|
||||
}
|
||||
assert_eq!(list_owned(dir.path()).await.len(), 24);
|
||||
for n in 0..24 {
|
||||
let id = format!("file-{n}");
|
||||
assert!(is_owned(dir.path(), "seller.onion", &id).await);
|
||||
let (mime, bytes) = read_owned(dir.path(), "seller.onion", &id).await.unwrap();
|
||||
assert_eq!(mime, "text/plain");
|
||||
assert_eq!(bytes, id.as_bytes());
|
||||
}
|
||||
record_purchase(
|
||||
dir.path(),
|
||||
"seller.onion",
|
||||
"file-0",
|
||||
"file-0",
|
||||
"text/plain",
|
||||
b"updated",
|
||||
5,
|
||||
"lightning",
|
||||
"later",
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(list_owned(dir.path()).await.len(), 24);
|
||||
assert_eq!(
|
||||
read_owned(dir.path(), "seller.onion", "file-0")
|
||||
.await
|
||||
.unwrap()
|
||||
.1,
|
||||
b"updated"
|
||||
);
|
||||
}
|
||||
#[tokio::test]
|
||||
async fn damaged_index_is_preserved_instead_of_erasing_prior_ownership() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
fs::create_dir_all(owned_root(dir.path())).await.unwrap();
|
||||
fs::write(index_path(dir.path()), b"damaged but preserve me")
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(record_purchase(
|
||||
dir.path(),
|
||||
"seller.onion",
|
||||
"new",
|
||||
"new",
|
||||
"text/plain",
|
||||
b"bytes",
|
||||
5,
|
||||
"lightning",
|
||||
"now"
|
||||
)
|
||||
.await
|
||||
.is_err());
|
||||
assert_eq!(
|
||||
fs::read(index_path(dir.path())).await.unwrap(),
|
||||
b"damaged but preserve me"
|
||||
);
|
||||
assert!(!bytes_path(dir.path(), "seller.onion", "new").exists());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -400,7 +400,7 @@ where
|
||||
if !authorized {
|
||||
if let Some(hash) = invoice_hash {
|
||||
if method_accepted(&item.access, "lightning")
|
||||
&& crate::content_invoice::is_paid_for(hash, id).await
|
||||
&& crate::content_invoice::is_paid_for(data_dir, hash, id).await
|
||||
{
|
||||
authorized = true;
|
||||
}
|
||||
|
||||
@@ -664,6 +664,10 @@ pub async fn start_stopped_stack_containers(data_dir: &Path) -> RecoveryReport {
|
||||
start_stopped_app_stacks(data_dir).await
|
||||
}
|
||||
|
||||
fn stack_member_needs_recovery(state: Option<&str>, user_stopped: bool) -> bool {
|
||||
!user_stopped && matches!(state, Some("exited" | "stopped" | "created" | "configured"))
|
||||
}
|
||||
|
||||
async fn start_stopped_app_stacks(data_dir: &Path) -> RecoveryReport {
|
||||
let user_stopped = load_user_stopped(data_dir).await;
|
||||
let mut report = RecoveryReport {
|
||||
@@ -677,24 +681,27 @@ async fn start_stopped_app_stacks(data_dir: &Path) -> RecoveryReport {
|
||||
continue;
|
||||
}
|
||||
|
||||
info!(
|
||||
"Recovering stopped {} stack containers after boot",
|
||||
stack.name
|
||||
);
|
||||
// Healthy members must never acquire a restarting overlay merely
|
||||
// because the periodic recovery scan ran. Queue existing stopped
|
||||
// members only; recheck each immediately before starting below.
|
||||
let mut pending = Vec::new();
|
||||
for container in stack.containers {
|
||||
let state = container_state(container).await;
|
||||
if stack_member_needs_recovery(state.as_deref(), user_stopped.contains(*container)) {
|
||||
pending.push((*container).to_string());
|
||||
}
|
||||
}
|
||||
if pending.is_empty() {
|
||||
continue;
|
||||
}
|
||||
info!("Recovering stopped {} stack containers", stack.name);
|
||||
repair_stack_network_aliases(stack).await;
|
||||
|
||||
// Register the whole stack up front: the per-member dependency waits
|
||||
// below can take minutes, and the UI should say "Restarting", not
|
||||
// "Stopped", for members still queued behind them.
|
||||
pending_boot_starts_add(
|
||||
stack
|
||||
.containers
|
||||
.iter()
|
||||
.filter(|c| !user_stopped.contains(**c))
|
||||
.map(|c| (*c).to_string()),
|
||||
);
|
||||
pending_boot_starts_add(pending.iter().cloned());
|
||||
|
||||
for container in stack.containers {
|
||||
if !pending.iter().any(|name| name.as_str() == *container) {
|
||||
continue;
|
||||
}
|
||||
if user_stopped.contains(*container) {
|
||||
info!("Skipping user-stopped container: {}", container);
|
||||
continue;
|
||||
@@ -706,8 +713,8 @@ async fn start_stopped_app_stacks(data_dir: &Path) -> RecoveryReport {
|
||||
pending_boot_start_done(container);
|
||||
continue;
|
||||
}
|
||||
Some(_) => {}
|
||||
None => {
|
||||
Some(state) if stack_member_needs_recovery(Some(&state), false) => {}
|
||||
_ => {
|
||||
pending_boot_start_done(container);
|
||||
continue;
|
||||
}
|
||||
@@ -1534,3 +1541,24 @@ mod installed_concurrency_tests {
|
||||
assert!(!dir.path().join("installed-apps.json.tmp").exists());
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod stack_recovery_overlay_tests {
|
||||
use super::stack_member_needs_recovery;
|
||||
#[test]
|
||||
fn only_existing_stopped_members_receive_recovery_overlay() {
|
||||
for state in [
|
||||
None,
|
||||
Some("running"),
|
||||
Some("paused"),
|
||||
Some("restarting"),
|
||||
Some("removing"),
|
||||
] {
|
||||
assert!(!stack_member_needs_recovery(state, false));
|
||||
}
|
||||
for state in ["exited", "stopped", "created", "configured"] {
|
||||
assert!(stack_member_needs_recovery(Some(state), false));
|
||||
assert!(!stack_member_needs_recovery(Some(state), true));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user