Serialize wallet mutations and preserve network and seed recovery state

This commit is contained in:
archipelago
2026-10-06 15:39:44 -04:00
parent 9f0df2ac44
commit 9c95b8732f
11 changed files with 322 additions and 70 deletions
+1 -1
View File
@@ -40,7 +40,7 @@ impl RpcHandler {
/// balance sitting in the *other* one so the UI can say what switching
/// would reveal rather than appearing to lose money.
pub(super) async fn handle_wallet_ecash_network(&self) -> Result<serde_json::Value> {
let current = ecash::load_network(&self.config.data_dir).await;
let current = ecash::load_network(&self.config.data_dir).await?;
let wallet = ecash::load_wallet(&self.config.data_dir).await?;
Ok(serde_json::json!({
"network": current,
+2 -13
View File
@@ -172,19 +172,8 @@ async fn process_payment(
session::save_sessions(data_dir, &store).await?;
// Record the streaming revenue
let mut wallet = ecash::load_wallet(data_dir).await?;
wallet.record_tx(
ecash::TransactionType::StreamingRevenue,
received_sats,
&format!(
"Streaming payment: {} sats for {} from {}",
received_sats, service.service_id, peer_id
),
&wallet.mint_url.clone(),
peer_id,
);
ecash::save_wallet(data_dir, &wallet).await?;
// Serialize the history write with all other wallet mutations.
ecash::record_streaming_revenue(data_dir, received_sats, &service.service_id, peer_id).await?;
debug!(
"Gate: accepted {} sats from {} for {}, allotment={}",
+85 -21
View File
@@ -267,40 +267,39 @@ impl EcashNetwork {
/// Read the node's ecash network. Absent file = mainnet, so nodes that never
/// touch this setting behave exactly as before.
pub async fn load_network(data_dir: &Path) -> EcashNetwork {
pub async fn load_network(data_dir: &Path) -> Result<EcashNetwork> {
let path = data_dir.join(NETWORK_FILE);
let Ok(content) = fs::read_to_string(&path).await else {
return EcashNetwork::Mainnet;
};
serde_json::from_str::<NetworkConfig>(&content)
.map(|c| c.network)
.unwrap_or_default()
match fs::read_to_string(&path).await {
Ok(content) => Ok(serde_json::from_str::<NetworkConfig>(&content)
.context("Ecash network configuration is damaged; no wallet was selected")?
.network),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(EcashNetwork::Mainnet),
Err(error) => Err(error).context("Could not read ecash network configuration"),
}
}
/// Switch the node's ecash network. The other network's wallet is left on
/// disk untouched, so switching is reversible and loses nothing.
pub async fn save_network(data_dir: &Path, network: EcashNetwork) -> Result<()> {
let _mutation = super::mutation::guard(data_dir).await?;
let dir = data_dir.join("wallet");
fs::create_dir_all(&dir)
.await
.context("Failed to create wallet dir")?;
let content = serde_json::to_string_pretty(&NetworkConfig { network })
.context("Failed to serialize ecash network")?;
fs::write(data_dir.join(NETWORK_FILE), content)
.await
.context("Failed to write ecash network")?;
write_file_atomically(&data_dir.join(NETWORK_FILE), &content).await?;
Ok(())
}
#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize)]
struct NetworkConfig {
#[serde(default)]
network: EcashNetwork,
}
/// Load wallet state from disk.
pub async fn load_wallet(data_dir: &Path) -> Result<WalletState> {
let network = load_network(data_dir).await;
let network = load_network(data_dir).await?;
let path = data_dir.join(network.wallet_file());
let content = match fs::read_to_string(&path).await {
Ok(content) => content,
@@ -376,20 +375,42 @@ async fn write_file_atomically(path: &Path, content: &str) -> Result<()> {
}
/// Save wallet state to disk.
pub async fn save_wallet(data_dir: &Path, wallet: &WalletState) -> Result<()> {
pub(super) async fn save_wallet(data_dir: &Path, wallet: &WalletState) -> Result<()> {
let dir = data_dir.join("wallet");
fs::create_dir_all(&dir)
.await
.context("Failed to create wallet dir")?;
let path = data_dir.join(load_network(data_dir).await.wallet_file());
let path = data_dir.join(load_network(data_dir).await?.wallet_file());
let content = serde_json::to_string_pretty(wallet).context("Failed to serialize wallet")?;
write_file_atomically(&path, &content).await?;
Ok(())
}
/// Record revenue without overwriting a simultaneous receipt/send snapshot.
pub(crate) async fn record_streaming_revenue(
data_dir: &Path,
amount_sats: u64,
service_id: &str,
peer_id: &str,
) -> Result<()> {
let _mutation = super::mutation::guard(data_dir).await?;
let mut wallet = load_wallet(data_dir).await?;
wallet.record_tx(
TransactionType::StreamingRevenue,
amount_sats,
&format!(
"Streaming payment: {} sats for {} from {}",
amount_sats, service_id, peer_id
),
&wallet.mint_url.clone(),
peer_id,
);
save_wallet(data_dir, &wallet).await
}
/// Load accepted mints list.
pub async fn load_accepted_mints(data_dir: &Path) -> Result<AcceptedMints> {
let network = load_network(data_dir).await;
let network = load_network(data_dir).await?;
let path = data_dir.join(network.mints_file());
if !path.exists() {
return Ok(AcceptedMints {
@@ -415,11 +436,12 @@ pub async fn load_accepted_mints(data_dir: &Path) -> Result<AcceptedMints> {
/// Save accepted mints list.
pub async fn save_accepted_mints(data_dir: &Path, mints: &AcceptedMints) -> Result<()> {
let _mutation = super::mutation::guard(data_dir).await?;
let dir = data_dir.join("wallet");
fs::create_dir_all(&dir)
.await
.context("Failed to create wallet dir")?;
let path = data_dir.join(load_network(data_dir).await.mints_file());
let path = data_dir.join(load_network(data_dir).await?.mints_file());
let content =
serde_json::to_string_pretty(mints).context("Failed to serialize accepted mints")?;
write_file_atomically(&path, &content).await?;
@@ -451,6 +473,7 @@ pub async fn mint_quote(
/// Mint new ecash tokens after a Lightning invoice has been paid.
pub async fn mint_tokens(data_dir: &Path, quote_id: &str, amount_sats: u64) -> Result<u64> {
let _mutation = super::mutation::guard(data_dir).await?;
let mut wallet = load_wallet(data_dir).await?;
let mint_url = wallet.mint_url.clone();
let client = mint_client(data_dir, &mint_url).await?;
@@ -481,6 +504,7 @@ pub async fn melt_quote(data_dir: &Path, bolt11: &str) -> Result<super::mint_cli
/// Melt ecash tokens to pay a Lightning invoice.
pub async fn melt_tokens(data_dir: &Path, quote_id: &str, bolt11: &str) -> Result<u64> {
let _mutation = super::mutation::guard(data_dir).await?;
let mut wallet = load_wallet(data_dir).await?;
let mint_url = wallet.mint_url.clone();
let client = mint_client(data_dir, &mint_url).await?;
@@ -591,6 +615,17 @@ pub async fn swap_between_mints(
to_mint: &str,
amount_sats: u64,
max_fee_sats: u64,
) -> Result<u64> {
let _mutation = super::mutation::guard(data_dir).await?;
swap_between_mints_locked(data_dir, from_mint, to_mint, amount_sats, max_fee_sats).await
}
async fn swap_between_mints_locked(
data_dir: &Path,
from_mint: &str,
to_mint: &str,
amount_sats: u64,
max_fee_sats: u64,
) -> Result<u64> {
if amount_sats == 0 {
anyhow::bail!("swap amount must be greater than zero");
@@ -747,8 +782,9 @@ async fn wait_for_mint_quote_paid(client: &MintClient, quote_id: &str) -> Result
/// Create an ecash token string to send to a peer, drawing from the home mint.
pub async fn send_token(data_dir: &Path, amount_sats: u64) -> Result<String> {
let _mutation = super::mutation::guard(data_dir).await?;
let mint_url = load_wallet(data_dir).await?.mint_url;
send_token_at(data_dir, &mint_url, amount_sats).await
send_token_at_locked(data_dir, &mint_url, amount_sats).await
}
/// Create an ecash token denominated in a specific mint's tokens.
@@ -757,6 +793,11 @@ pub async fn send_token(data_dir: &Path, amount_sats: u64) -> Result<String> {
/// on the seeder's accepted mint, we send a token from *that* mint so the seeder
/// only ever receives its own mint's proofs (see plan §2a, payer-side swap).
pub async fn send_token_at(data_dir: &Path, mint_url: &str, amount_sats: u64) -> Result<String> {
let _mutation = super::mutation::guard(data_dir).await?;
send_token_at_locked(data_dir, mint_url, amount_sats).await
}
async fn send_token_at_locked(data_dir: &Path, mint_url: &str, amount_sats: u64) -> Result<String> {
let mut wallet = load_wallet(data_dir).await?;
let mint_url = mint_url.to_string();
@@ -947,6 +988,7 @@ pub async fn build_payment_token(
amount_sats: u64,
max_fee_sats: u64,
) -> Result<String> {
let _mutation = super::mutation::guard(data_dir).await?;
if amount_sats == 0 {
anyhow::bail!("payment amount must be greater than zero");
}
@@ -975,15 +1017,16 @@ pub async fn build_payment_token(
"Payment plan: direct from {} for {} sats",
mint_url, amount_sats
);
send_token_at(data_dir, &mint_url, amount_sats).await
send_token_at_locked(data_dir, &mint_url, amount_sats).await
}
PaymentPlan::Swap { from_mint, to_mint } => {
debug!(
"Payment plan: swap {}→{} then pay {} sats (fee cap {})",
from_mint, to_mint, amount_sats, max_fee_sats
);
swap_between_mints(data_dir, &from_mint, &to_mint, amount_sats, max_fee_sats).await?;
send_token_at(data_dir, &to_mint, amount_sats).await
swap_between_mints_locked(data_dir, &from_mint, &to_mint, amount_sats, max_fee_sats)
.await?;
send_token_at_locked(data_dir, &to_mint, amount_sats).await
}
PaymentPlan::Insufficient => anyhow::bail!(
"cannot pay {} sats: no accepted mint covers it within balance/trust",
@@ -1051,6 +1094,7 @@ async fn remove_pending_swap(data_dir: &Path, mint_quote_id: &str) -> Result<()>
/// Returns the total sats reclaimed. Safe to call repeatedly (idempotent): a
/// quote is only minted once, and `ISSUED` quotes are never re-claimed.
pub async fn resume_pending_swaps(data_dir: &Path) -> Result<u64> {
let _mutation = super::mutation::guard(data_dir).await?;
let pending = load_pending_swaps(data_dir).await?;
let mut reclaimed = 0u64;
for swap in pending {
@@ -1189,6 +1233,7 @@ fn target_liquidity_score(liq: &SwapLiquidity, to_mint: &str) -> i64 {
/// Receive a Cashu token from a peer — swaps proofs at the mint for fresh ones.
pub async fn receive_token(data_dir: &Path, token_str: &str) -> Result<u64> {
let _mutation = super::mutation::guard(data_dir).await?;
// Handle legacy format for backwards compatibility
if token_str.starts_with("cashuSend_") {
return receive_legacy_token(data_dir, token_str).await;
@@ -1314,6 +1359,7 @@ pub async fn verify_and_receive_payment(
token_str: &str,
required_sats: u64,
) -> Result<u64> {
let _mutation = super::mutation::guard(data_dir).await?;
let token_str = token_str.trim();
// Synthetic legacy balances are not cryptographic proof of payment.
if token_str.starts_with("cashuSend_") {
@@ -1434,6 +1480,7 @@ pub struct RestoreOutcome {
/// coins or resurrecting spent ones, which matters because the most likely
/// time to press this button is when something already looks wrong.
pub async fn restore_from_seed(data_dir: &Path, mint_url: &str) -> Result<RestoreOutcome> {
let _mutation = super::mutation::guard(data_dir).await?;
let recovery = RecoverySource::load(data_dir).await?.ok_or_else(|| {
anyhow::anyhow!(
"This wallet has no backup phrase yet, so there is nothing to restore from. \
@@ -2297,7 +2344,7 @@ mod tests {
async fn ecash_network_defaults_to_mainnet_and_leaves_files_alone() {
let tmp = TempDir::new().unwrap();
let dir = tmp.path();
assert_eq!(load_network(dir).await, EcashNetwork::Mainnet);
assert_eq!(load_network(dir).await.unwrap(), EcashNetwork::Mainnet);
// A node that never touches this setting has no new file.
assert!(!dir.join(NETWORK_FILE).exists());
assert_eq!(
@@ -2391,6 +2438,23 @@ mod tests {
assert_eq!(std::fs::read_to_string(&path).unwrap(), damaged);
}
#[tokio::test]
async fn damaged_network_selection_never_falls_back_to_real_funds() {
let tmp = TempDir::new().unwrap();
let path = tmp.path().join(NETWORK_FILE);
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
for bytes in ["", "{}", "{broken", r#"{"network":"unknown"}"#] {
std::fs::write(&path, bytes).unwrap();
assert!(load_network(tmp.path()).await.is_err());
assert!(load_wallet(tmp.path()).await.is_err());
assert!(save_wallet(tmp.path(), &WalletState::default())
.await
.is_err());
assert!(!tmp.path().join("wallet/ecash.json").exists());
assert_eq!(std::fs::read_to_string(&path).unwrap(), bytes);
}
}
#[tokio::test]
async fn an_empty_wallet_file_is_preserved_as_damaged() {
let tmp = TempDir::new().unwrap();
+2 -2
View File
@@ -559,7 +559,7 @@ async fn ecash_phrase(data_dir: &Path) -> Result<(String, [u8; 64])> {
/// so a node that restores the same ecash phrase recovers the same address.
pub async fn lnaddress(data_dir: &Path) -> Result<serde_json::Value> {
let _state_guard = MINIBITS_STATE_LOCK.lock().await;
let network = ecash::load_network(data_dir).await;
let network = ecash::load_network(data_dir).await?;
if network == EcashNetwork::Testnet {
return Err(anyhow!(
"Minibits Lightning addresses are mainnet-only — switch the ecash network to mainnet to set one up"
@@ -901,7 +901,7 @@ async fn ensure_mint_accepted(data_dir: &Path, mint_url: &str) -> Result<()> {
/// so it isn't purely a log-line event.
pub async fn claim_and_redeem(data_dir: &Path) -> Result<ClaimOutcome> {
let _state_guard = MINIBITS_STATE_LOCK.lock().await;
let network = ecash::load_network(data_dir).await;
let network = ecash::load_network(data_dir).await?;
if network == EcashNetwork::Testnet {
return Ok(NO_CLAIMS);
}
+1
View File
@@ -8,5 +8,6 @@ pub mod ecash;
pub mod fedimint_client;
pub mod minibits;
pub mod mint_client;
mod mutation;
pub mod nut13;
pub mod profits;
+73
View File
@@ -0,0 +1,73 @@
//! Serialize complete wallet mutations, independently for each node data directory.
//! This prevents lost updates; durable recovery of interrupted remote operations
//! requires the separate operation journal.
use anyhow::{Context, Result};
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex as RegistryMutex, OnceLock, Weak};
use tokio::sync::{Mutex, OwnedMutexGuard};
type Registry = HashMap<PathBuf, Weak<Mutex<()>>>;
static LOCKS: OnceLock<RegistryMutex<Registry>> = OnceLock::new();
async fn lock_for(data_dir: &Path) -> Result<Arc<Mutex<()>>> {
tokio::fs::create_dir_all(data_dir)
.await
.context("Could not prepare wallet data directory")?;
let canonical = tokio::fs::canonicalize(data_dir)
.await
.context("Could not resolve wallet data directory")?;
let mut registry = LOCKS
.get_or_init(|| RegistryMutex::new(HashMap::new()))
.lock()
.map_err(|_| anyhow::anyhow!("Wallet mutation lock registry is unavailable"))?;
registry.retain(|_, lock| lock.strong_count() > 0);
if let Some(lock) = registry.get(&canonical).and_then(Weak::upgrade) {
return Ok(lock);
}
let lock = Arc::new(Mutex::new(()));
registry.insert(canonical, Arc::downgrade(&lock));
Ok(lock)
}
pub(super) async fn guard(data_dir: &Path) -> Result<OwnedMutexGuard<()>> {
Ok(lock_for(data_dir).await?.lock_owned().await)
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn aliases_share_a_lock_but_different_nodes_do_not() {
let first = tempfile::tempdir().unwrap();
let second = tempfile::tempdir().unwrap();
let held = guard(first.path()).await.unwrap();
let alias = lock_for(&first.path().join(".")).await.unwrap();
assert!(alias.try_lock().is_err());
let independent = lock_for(second.path()).await.unwrap();
assert!(independent.try_lock().is_ok());
drop(held);
assert!(alias.try_lock().is_ok());
}
#[tokio::test]
async fn symlink_aliases_cannot_bypass_serialization() {
let parent = tempfile::tempdir().unwrap();
let actual = parent.path().join("node");
std::fs::create_dir(&actual).unwrap();
let alias = parent.path().join("alias");
std::os::unix::fs::symlink(&actual, &alias).unwrap();
let _held = guard(&actual).await.unwrap();
assert!(lock_for(&alias).await.unwrap().try_lock().is_err());
}
#[tokio::test]
async fn invalid_directory_fails_without_replacing_existing_data() {
let parent = tempfile::tempdir().unwrap();
let file = parent.path().join("node");
std::fs::write(&file, b"preserve").unwrap();
assert!(guard(&file).await.is_err());
assert_eq!(std::fs::read(file).unwrap(), b"preserve");
}
}
+66 -33
View File
@@ -223,6 +223,7 @@ pub async fn establish_from_master(
data_dir: &Path,
master: &crate::seed::MasterSeed,
) -> Result<EcashSeed> {
let _mutation = super::mutation::guard(data_dir).await?;
let derived = crate::seed::derive_cashu_mnemonic(master)?;
if let Some(existing) = load_seed(data_dir).await? {
@@ -259,6 +260,7 @@ pub async fn establish_from_master(
/// silent default when derivation was possible; [`establish_from_master`] is
/// what a node with a master seed gets.
pub async fn establish_independent(data_dir: &Path) -> Result<EcashSeed> {
let _mutation = super::mutation::guard(data_dir).await?;
if let Some(existing) = load_seed(data_dir).await? {
return Ok(existing);
}
@@ -294,6 +296,7 @@ pub async fn establish_independent(data_dir: &Path) -> Result<EcashSeed> {
/// dangerous choice if the imported phrase turned out to be the one already
/// in use.
pub async fn import_mnemonic(data_dir: &Path, words: &str, confirm: bool) -> Result<EcashSeed> {
let _mutation = super::mutation::guard(data_dir).await?;
let mnemonic: bip39::Mnemonic = words
.split_whitespace()
.collect::<Vec<_>>()
@@ -330,24 +333,24 @@ pub async fn import_mnemonic(data_dir: &Path, words: &str, confirm: bool) -> Res
Ok(EcashSeed::from_mnemonic(mnemonic, SeedSource::Imported))
}
/// Move the current seed file aside, timestamped, before it is replaced.
/// Durably copy the current seed before replacing it; keep the live seed on failure.
///
/// Never deleted and never overwritten: this file may be the last copy of the
/// words a balance was minted under, and the whole point of the module is that
/// such a thing is not casually destroyed.
async fn archive_seed(data_dir: &Path) -> Result<()> {
let from = seed_path(data_dir);
if !from.exists() {
return Ok(());
}
let content = fs::read(&from)
.await
.context("Could not read the previous ecash phrase for backup")?;
let stamp = chrono::Utc::now().format("%Y%m%dT%H%M%SZ");
let to = data_dir.join(format!("wallet/cashu_seed.replaced-{stamp}.json"));
fs::rename(&from, &to).await.with_context(|| {
format!(
"Could not archive the previous ecash phrase to {}",
to.display()
)
})?;
let to = data_dir.join(format!(
"wallet/cashu_seed.replaced-{stamp}-{}.json",
uuid::Uuid::new_v4()
));
persist_private_file(&to, &content)
.await
.context("Could not durably archive the previous ecash phrase")?;
warn!("Previous ecash phrase archived to {}", to.display());
Ok(())
}
@@ -367,17 +370,7 @@ async fn write_seed(data_dir: &Path, mnemonic: &bip39::Mnemonic, source: SeedSou
};
let content =
serde_json::to_string_pretty(&stored).context("Failed to serialize the ecash seed")?;
fs::write(&path, content)
.await
.context("Failed to write the ecash seed")?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
fs::set_permissions(&path, std::fs::Permissions::from_mode(0o600))
.await
.context("Failed to restrict permissions on the ecash seed")?;
}
persist_private_file(&path, content.as_bytes()).await?;
Ok(())
}
@@ -420,46 +413,48 @@ pub async fn reserve_counters(data_dir: &Path, keyset_id: &str, count: usize) ->
}
let content =
serde_json::to_string_pretty(&state).context("Failed to serialize ecash counters")?;
persist_counters(&path, content.as_bytes()).await?;
persist_private_file(&path, content.as_bytes()).await?;
Ok(start)
}
/// Never truncate the active reservation file. A reservation is not usable
/// until both its replacement file and directory entry have reached storage.
async fn persist_counters(path: &Path, content: &[u8]) -> Result<()> {
async fn persist_private_file(path: &Path, content: &[u8]) -> Result<()> {
use tokio::io::AsyncWriteExt;
struct PendingCounterFile(PathBuf);
impl Drop for PendingCounterFile {
struct PendingPrivateFile(PathBuf);
impl Drop for PendingPrivateFile {
fn drop(&mut self) {
let _ = std::fs::remove_file(&self.0);
}
}
let parent = path.parent().context("Counter file has no directory")?;
let parent = path
.parent()
.context("Private wallet file has no directory")?;
let temporary =
PendingCounterFile(parent.join(format!(".cashu-counters-{}.tmp", uuid::Uuid::new_v4())));
PendingPrivateFile(parent.join(format!(".cashu-private-{}.tmp", uuid::Uuid::new_v4())));
let mut file = fs::OpenOptions::new()
.write(true)
.create_new(true)
.mode(0o600)
.open(&temporary.0)
.await
.context("Could not create the ecash counter reservation")?;
.context("Could not create the private ecash state")?;
file.write_all(content)
.await
.context("Could not write the ecash counter reservation")?;
.context("Could not write the private ecash state")?;
file.sync_all()
.await
.context("Could not flush the ecash counter reservation")?;
.context("Could not flush the private ecash state")?;
drop(file);
fs::rename(&temporary.0, path)
.await
.context("Could not replace the ecash counter reservation")?;
.context("Could not replace the private ecash state")?;
fs::File::open(parent)
.await?
.sync_all()
.await
.context("Could not flush the ecash counter directory")?;
.context("Could not flush the private ecash directory")?;
Ok(())
}
@@ -909,4 +904,42 @@ mod tests {
assert!(source.next_outputs("01fc0ec0e59cd6fa", 1).await.is_err());
assert!(!dir.path().join(COUNTER_FILE).exists());
}
#[tokio::test]
async fn archiving_keeps_the_live_seed_and_never_overwrites_an_earlier_backup() {
use std::os::unix::fs::PermissionsExt;
let dir = tempfile::tempdir().unwrap();
establish_independent(dir.path()).await.unwrap();
let live = fs::read(seed_path(dir.path())).await.unwrap();
archive_seed(dir.path()).await.unwrap();
archive_seed(dir.path()).await.unwrap();
assert_eq!(fs::read(seed_path(dir.path())).await.unwrap(), live);
let backups: Vec<_> = std::fs::read_dir(dir.path().join("wallet"))
.unwrap()
.map(|entry| entry.unwrap().path())
.filter(|path| {
path.file_name()
.unwrap()
.to_string_lossy()
.starts_with("cashu_seed.replaced-")
})
.collect();
assert_eq!(backups.len(), 2);
for backup in backups {
assert_eq!(std::fs::read(&backup).unwrap(), live);
assert_eq!(
std::fs::metadata(backup).unwrap().permissions().mode() & 0o777,
0o600
);
}
}
#[tokio::test]
async fn simultaneous_seed_establishment_keeps_one_identity() {
let dir = tempfile::tempdir().unwrap();
let (first, second) = tokio::join!(
establish_independent(dir.path()),
establish_independent(dir.path())
);
assert_eq!(first.unwrap().phrase(), second.unwrap().phrase());
}
}
@@ -426,3 +426,37 @@ async fn paid_file_gate_delivers_bytes_only_after_payment_and_does_not_charge_mi
}
}
}
#[tokio::test]
async fn simultaneous_receipts_and_revenue_writes_preserve_every_payment() {
let mint = Mint::start(0, None).await;
let wallet_dir = mint.wallet().await;
let mut tasks = tokio::task::JoinSet::new();
for exponent in 0..8 {
let amount = 1u64 << exponent;
let token = CashuToken::new(&mint.url, vec![proof(ACTIVE, amount)])
.serialize()
.unwrap();
let data_dir = wallet_dir.path().to_path_buf();
tasks.spawn(async move {
verify_and_receive_payment(&data_dir, &token, amount)
.await
.unwrap();
});
}
for index in 0..16 {
let data_dir = wallet_dir.path().to_path_buf();
tasks.spawn(async move {
record_streaming_revenue(&data_dir, 1, "fixture-service", &format!("peer-{index}"))
.await
.unwrap();
});
}
while let Some(result) = tasks.join_next().await {
result.unwrap();
}
let wallet = load_wallet(wallet_dir.path()).await.unwrap();
assert_eq!(wallet.balance(), 255);
assert_eq!(wallet.transactions.len(), 24);
assert_eq!(mint.requests.lock().unwrap().len(), 8);
}