Publish content with an atomic price and visibility policy
This commit is contained in:
@@ -6,11 +6,12 @@
|
||||
use anyhow::{Context, Result};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::path::{Path, PathBuf};
|
||||
use tokio::fs;
|
||||
use tokio::{fs, io::AsyncWriteExt, sync::Mutex};
|
||||
use tracing::{debug, warn};
|
||||
|
||||
const CATALOG_FILE: &str = "content/catalog.json";
|
||||
const CONTENT_DIR: &str = "content/files";
|
||||
static CATALOG_WRITES: Mutex<()> = Mutex::const_new(());
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct ContentItem {
|
||||
@@ -78,27 +79,39 @@ pub struct ContentCatalog {
|
||||
/// Load the content catalog from disk.
|
||||
pub async fn load_catalog(data_dir: &Path) -> Result<ContentCatalog> {
|
||||
let path = data_dir.join(CATALOG_FILE);
|
||||
if !path.exists() {
|
||||
return Ok(ContentCatalog::default());
|
||||
}
|
||||
let content = fs::read_to_string(&path)
|
||||
.await
|
||||
.context("Failed to read content catalog")?;
|
||||
let catalog: ContentCatalog = serde_json::from_str(&content).unwrap_or_default();
|
||||
Ok(catalog)
|
||||
let content = match fs::read(&path).await {
|
||||
Ok(bytes) => bytes,
|
||||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
|
||||
return Ok(ContentCatalog::default())
|
||||
}
|
||||
Err(error) => return Err(error).context("Failed to read content catalog"),
|
||||
};
|
||||
serde_json::from_slice(&content)
|
||||
.context("Invalid content catalog; existing shares were preserved")
|
||||
}
|
||||
|
||||
/// Save the content catalog to disk.
|
||||
pub async fn save_catalog(data_dir: &Path, catalog: &ContentCatalog) -> Result<()> {
|
||||
let _lock = CATALOG_WRITES.lock().await;
|
||||
save_catalog_unlocked(data_dir, catalog).await
|
||||
}
|
||||
|
||||
async fn save_catalog_unlocked(data_dir: &Path, catalog: &ContentCatalog) -> Result<()> {
|
||||
let dir = data_dir.join("content");
|
||||
fs::create_dir_all(&dir)
|
||||
.await
|
||||
.context("Failed to create content dir")?;
|
||||
let path = data_dir.join(CATALOG_FILE);
|
||||
let content = serde_json::to_string_pretty(catalog).context("Failed to serialize catalog")?;
|
||||
fs::write(&path, content)
|
||||
.await
|
||||
.context("Failed to write catalog")?;
|
||||
let (file, temporary) = tempfile::NamedTempFile::new_in(&dir)?.into_parts();
|
||||
let mut file = fs::File::from_std(file);
|
||||
file.write_all(content.as_bytes()).await?;
|
||||
file.sync_all().await?;
|
||||
drop(file);
|
||||
// Complete the rename synchronously while holding the catalog lock: an
|
||||
// aborted async caller must not leave a late rename racing the next writer.
|
||||
std::fs::rename(&temporary, &path).context("Failed to commit content catalog")?;
|
||||
fs::File::open(&dir).await?.sync_all().await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -106,13 +119,14 @@ pub async fn save_catalog(data_dir: &Path, catalog: &ContentCatalog) -> Result<(
|
||||
/// means the entry gets pruned again next time it's requested, so errors are
|
||||
/// logged rather than propagated.
|
||||
async fn prune_missing_content_entry(data_dir: &Path, id: &str) {
|
||||
let _lock = CATALOG_WRITES.lock().await;
|
||||
let Ok(mut catalog) = load_catalog(data_dir).await else {
|
||||
return;
|
||||
};
|
||||
let before = catalog.items.len();
|
||||
catalog.items.retain(|i| i.id != id);
|
||||
if catalog.items.len() != before {
|
||||
if let Err(e) = save_catalog(data_dir, &catalog).await {
|
||||
if let Err(e) = save_catalog_unlocked(data_dir, &catalog).await {
|
||||
warn!(error = %e, content_id = %id, "failed to save catalog after pruning missing content entry");
|
||||
}
|
||||
}
|
||||
@@ -149,6 +163,7 @@ pub fn content_file_path(data_dir: &Path, item: &ContentItem) -> PathBuf {
|
||||
/// (2026-07-22). Same filename → update the existing entry in place and
|
||||
/// keep its id, so existing buyers' owned records stay valid.
|
||||
pub async fn add_item(data_dir: &Path, item: ContentItem) -> Result<ContentCatalog> {
|
||||
let _lock = CATALOG_WRITES.lock().await;
|
||||
let mut catalog = load_catalog(data_dir).await?;
|
||||
if catalog.items.iter().any(|i| i.id == item.id) {
|
||||
return Err(anyhow::anyhow!("Content item '{}' already exists", item.id));
|
||||
@@ -165,24 +180,26 @@ pub async fn add_item(data_dir: &Path, item: ContentItem) -> Result<ContentCatal
|
||||
} else {
|
||||
catalog.items.push(item);
|
||||
}
|
||||
save_catalog(data_dir, &catalog).await?;
|
||||
save_catalog_unlocked(data_dir, &catalog).await?;
|
||||
Ok(catalog)
|
||||
}
|
||||
|
||||
/// Remove a content item from the catalog.
|
||||
pub async fn remove_item(data_dir: &Path, id: &str) -> Result<ContentCatalog> {
|
||||
let _lock = CATALOG_WRITES.lock().await;
|
||||
let mut catalog = load_catalog(data_dir).await?;
|
||||
catalog.items.retain(|i| i.id != id);
|
||||
save_catalog(data_dir, &catalog).await?;
|
||||
save_catalog_unlocked(data_dir, &catalog).await?;
|
||||
Ok(catalog)
|
||||
}
|
||||
|
||||
/// Update access control for a content item.
|
||||
pub async fn set_access(data_dir: &Path, id: &str, access: AccessControl) -> Result<()> {
|
||||
let _lock = CATALOG_WRITES.lock().await;
|
||||
let mut catalog = load_catalog(data_dir).await?;
|
||||
if let Some(item) = catalog.items.iter_mut().find(|i| i.id == id) {
|
||||
item.access = access;
|
||||
save_catalog(data_dir, &catalog).await?;
|
||||
save_catalog_unlocked(data_dir, &catalog).await?;
|
||||
Ok(())
|
||||
} else {
|
||||
Err(anyhow::anyhow!("Content item '{}' not found", id))
|
||||
@@ -191,16 +208,36 @@ pub async fn set_access(data_dir: &Path, id: &str, access: AccessControl) -> Res
|
||||
|
||||
/// Update availability for a content item.
|
||||
pub async fn set_availability(data_dir: &Path, id: &str, availability: Availability) -> Result<()> {
|
||||
let _lock = CATALOG_WRITES.lock().await;
|
||||
let mut catalog = load_catalog(data_dir).await?;
|
||||
if let Some(item) = catalog.items.iter_mut().find(|i| i.id == id) {
|
||||
item.availability = availability;
|
||||
save_catalog(data_dir, &catalog).await?;
|
||||
save_catalog_unlocked(data_dir, &catalog).await?;
|
||||
Ok(())
|
||||
} else {
|
||||
Err(anyhow::anyhow!("Content item '{}' not found", id))
|
||||
}
|
||||
}
|
||||
|
||||
/// Change price and visibility in one durable catalog transaction.
|
||||
pub async fn configure_item(
|
||||
data_dir: &Path,
|
||||
id: &str,
|
||||
access: AccessControl,
|
||||
availability: Availability,
|
||||
) -> Result<()> {
|
||||
let _lock = CATALOG_WRITES.lock().await;
|
||||
let mut catalog = load_catalog(data_dir).await?;
|
||||
let item = catalog
|
||||
.items
|
||||
.iter_mut()
|
||||
.find(|item| item.id == id)
|
||||
.context("Content item not found")?;
|
||||
item.access = access;
|
||||
item.availability = availability;
|
||||
save_catalog_unlocked(data_dir, &catalog).await
|
||||
}
|
||||
|
||||
/// A byte range request (start, optional end).
|
||||
pub enum ByteRange {
|
||||
From { start: u64, end: Option<u64> },
|
||||
@@ -1066,6 +1103,87 @@ mod paid_read_order_tests {
|
||||
dir
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn corrupt_catalog_is_not_overwritten_by_a_mutation() {
|
||||
let dir = fixture(b"file").await;
|
||||
let item = load_catalog(dir.path()).await.unwrap().items.remove(0);
|
||||
fs::write(dir.path().join(CATALOG_FILE), b"corrupt")
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(add_item(dir.path(), item).await.is_err());
|
||||
assert!(remove_item(dir.path(), "paid").await.is_err());
|
||||
assert_eq!(
|
||||
fs::read(dir.path().join(CATALOG_FILE)).await.unwrap(),
|
||||
b"corrupt"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn concurrent_catalog_updates_preserve_all_items_and_sharing_fields() {
|
||||
let dir = fixture(b"file").await;
|
||||
let template = load_catalog(dir.path()).await.unwrap().items.remove(0);
|
||||
let mut tasks = Vec::new();
|
||||
for i in 0..16 {
|
||||
let root = dir.path().to_owned();
|
||||
let mut item = template.clone();
|
||||
item.id = format!("item-{i}");
|
||||
item.filename = format!("file-{i}");
|
||||
tasks.push(tokio::spawn(async move {
|
||||
add_item(&root, item).await.unwrap();
|
||||
}));
|
||||
}
|
||||
for task in tasks {
|
||||
task.await.unwrap();
|
||||
}
|
||||
assert_eq!(load_catalog(dir.path()).await.unwrap().items.len(), 17);
|
||||
let (price, visibility) = tokio::join!(
|
||||
set_access(
|
||||
dir.path(),
|
||||
"paid",
|
||||
AccessControl::Paid {
|
||||
price_sats: 7,
|
||||
accepted: vec!["ecash".into()]
|
||||
}
|
||||
),
|
||||
set_availability(dir.path(), "paid", Availability::Nobody),
|
||||
);
|
||||
price.unwrap();
|
||||
visibility.unwrap();
|
||||
let item = load_catalog(dir.path()).await.unwrap().items.remove(0);
|
||||
assert!(matches!(
|
||||
item.access,
|
||||
AccessControl::Paid { price_sats: 7, .. }
|
||||
));
|
||||
assert!(matches!(item.availability, Availability::Nobody));
|
||||
configure_item(
|
||||
dir.path(),
|
||||
"paid",
|
||||
AccessControl::Paid {
|
||||
price_sats: 9,
|
||||
accepted: vec!["ecash".into()],
|
||||
},
|
||||
Availability::AllPeers,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let item = load_catalog(dir.path()).await.unwrap().items.remove(0);
|
||||
assert!(matches!(
|
||||
item.access,
|
||||
AccessControl::Paid { price_sats: 9, .. }
|
||||
));
|
||||
assert!(matches!(item.availability, Availability::AllPeers));
|
||||
use std::os::unix::fs::PermissionsExt;
|
||||
assert_eq!(
|
||||
fs::metadata(dir.path().join(CATALOG_FILE))
|
||||
.await
|
||||
.unwrap()
|
||||
.permissions()
|
||||
.mode()
|
||||
& 0o777,
|
||||
0o600
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn all_read_failures_precede_redemption_even_as_root() {
|
||||
for kind in [
|
||||
|
||||
Reference in New Issue
Block a user