use std::{ fmt, path::{Path, PathBuf}, sync::Arc, time::{Duration, SystemTime, UNIX_EPOCH}, }; use serde::{Deserialize, Serialize}; use sha2::{Digest, Sha256}; use spacetime_client::{SpacetimeClient, SpacetimeClientError}; use tokio::{ fs::{self, File, OpenOptions}, io::{AsyncReadExt, AsyncWriteExt}, sync::{Mutex, Notify}, time::sleep, }; use tracing::{debug, warn}; use crate::config::AppConfig; const PENDING_FILE_PREFIX: &str = "refund-"; const CORRUPT_FILE_PREFIX: &str = "corrupt-"; const TEMP_FILE_PREFIX: &str = "tmp-"; const OUTBOX_FILE_EXTENSION: &str = ".json"; #[derive(Clone)] pub struct WalletRefundOutbox { dir: PathBuf, batch_size: usize, flush_interval: Duration, max_bytes: u64, spacetime_client: SpacetimeClient, enqueue_lock: Arc>, flush_notify: Arc, } #[derive(Clone, Debug, Deserialize, Serialize)] pub(crate) struct WalletRefundOutboxRecord { pub owner_user_id: String, pub amount: u64, pub ledger_id: String, pub created_at_micros: i64, pub asset_kind: String, pub asset_id: String, #[serde(default, skip_serializing_if = "Option::is_none")] pub external_generation_job_id: Option, } #[derive(Debug)] pub enum WalletRefundOutboxEnqueueOutcome { Enqueued, Dropped { reason: &'static str }, } #[derive(Debug)] pub enum WalletRefundOutboxError { Io(std::io::Error), Json(serde_json::Error), Spacetime(SpacetimeClientError), } impl WalletRefundOutbox { pub fn from_config(config: &AppConfig, spacetime_client: SpacetimeClient) -> Option> { if !config.wallet_refund_outbox_enabled { return None; } Some(Arc::new(Self { dir: config.wallet_refund_outbox_dir.clone(), batch_size: config.wallet_refund_outbox_batch_size.max(1), flush_interval: config.wallet_refund_outbox_flush_interval, max_bytes: config.wallet_refund_outbox_max_bytes, spacetime_client, enqueue_lock: Arc::new(Mutex::new(())), flush_notify: Arc::new(Notify::new()), })) } pub async fn enqueue( &self, record: WalletRefundOutboxRecord, ) -> Result { let _guard = self.enqueue_lock.lock().await; fs::create_dir_all(&self.dir).await?; let pending_path = self.pending_path_for_ledger(&record.ledger_id); if fs::metadata(&pending_path).await.is_ok() { self.flush_notify.notify_one(); return Ok(WalletRefundOutboxEnqueueOutcome::Enqueued); } let bytes = serde_json::to_vec(&record)?; let line_bytes = bytes.len().min(u64::MAX as usize) as u64; let current_bytes = directory_size_if_exists(&self.dir).unwrap_or(0); if current_bytes.saturating_add(line_bytes) > self.max_bytes { return Ok(WalletRefundOutboxEnqueueOutcome::Dropped { reason: "max_bytes", }); } let temp_path = self.temp_path(); let mut file = OpenOptions::new() .create_new(true) .write(true) .open(&temp_path) .await?; file.write_all(&bytes).await?; file.flush().await?; file.sync_data().await?; drop(file); if fs::metadata(&pending_path).await.is_ok() { let _ = fs::remove_file(&temp_path).await; self.flush_notify.notify_one(); return Ok(WalletRefundOutboxEnqueueOutcome::Enqueued); } fs::rename(&temp_path, &pending_path).await?; sync_directory_metadata(&self.dir).await?; self.flush_notify.notify_one(); Ok(WalletRefundOutboxEnqueueOutcome::Enqueued) } pub fn spawn_worker(self: Arc) { tokio::spawn(async move { loop { tokio::select! { _ = sleep(self.flush_interval) => { if let Err(error) = self.flush_pending_files_once().await { warn!(error = %error, "wallet refund outbox 重放退款失败,将保留文件等待重试"); } } _ = self.flush_notify.notified() => { if let Err(error) = self.flush_pending_files_once().await { warn!(error = %error, "wallet refund outbox 主动重放退款失败,将保留文件等待重试"); } } } } }); } pub async fn flush_for_shutdown(&self) -> Result<(), WalletRefundOutboxError> { self.flush_pending_files_once().await } async fn flush_pending_files_once(&self) -> Result<(), WalletRefundOutboxError> { fs::create_dir_all(&self.dir).await?; let pending_files = self.list_pending_files().await?; for path in pending_files.into_iter().take(self.batch_size) { let record = match read_refund_record(&path).await { Ok(record) => record, Err(error) if error.is_data_corruption() => { let corrupt_path = self.corrupt_path_for(&path); fs::rename(&path, &corrupt_path).await?; sync_directory_metadata(&self.dir).await?; warn!( error = %error, source = %path.display(), target = %corrupt_path.display(), "wallet refund outbox 文件无法解析,已隔离" ); continue; } Err(error) => return Err(error), }; match self .spacetime_client .refund_profile_wallet_points_with_metadata( record.owner_user_id.clone(), record.amount, record.ledger_id.clone(), record.created_at_micros, refund_metadata_json(record.external_generation_job_id.as_deref()), ) .await { Ok(_) => { match fs::remove_file(&path).await { Ok(()) => {} Err(error) if error.kind() == std::io::ErrorKind::NotFound => {} Err(error) => return Err(error.into()), } sync_directory_metadata(&self.dir).await?; debug!( ledger_id = %record.ledger_id, owner_user_id = %record.owner_user_id, asset_kind = %record.asset_kind, asset_id = %record.asset_id, external_generation_job_id = ?record.external_generation_job_id, path = %path.display(), "wallet refund outbox 退款已重放并删除文件" ); } Err(error) => return Err(WalletRefundOutboxError::Spacetime(error)), } } Ok(()) } async fn list_pending_files(&self) -> Result, WalletRefundOutboxError> { let mut entries = fs::read_dir(&self.dir).await?; let mut files = Vec::new(); while let Some(entry) = entries.next_entry().await? { let path = entry.path(); let Some(name) = path.file_name().and_then(|value| value.to_str()) else { continue; }; if name.starts_with(PENDING_FILE_PREFIX) && name.ends_with(OUTBOX_FILE_EXTENSION) { files.push(path); } } files.sort(); Ok(files) } fn pending_path_for_ledger(&self, ledger_id: &str) -> PathBuf { self.dir.join(format!( "{PENDING_FILE_PREFIX}{}{OUTBOX_FILE_EXTENSION}", ledger_id_hash(ledger_id) )) } fn temp_path(&self) -> PathBuf { self.dir.join(format!( "{TEMP_FILE_PREFIX}{}-{uuid}{OUTBOX_FILE_EXTENSION}", current_unix_micros(), uuid = uuid::Uuid::new_v4() )) } fn corrupt_path_for(&self, path: &Path) -> PathBuf { let name = path .file_name() .and_then(|value| value.to_str()) .unwrap_or("unknown.json"); self.dir.join(format!( "{CORRUPT_FILE_PREFIX}{}-{uuid}-{name}", current_unix_micros(), uuid = uuid::Uuid::new_v4() )) } } fn refund_metadata_json(external_generation_job_id: Option<&str>) -> String { let Some(external_generation_job_id) = external_generation_job_id .map(str::trim) .filter(|value| !value.is_empty()) else { return module_runtime::PROFILE_INVITE_CODE_METADATA_DEFAULT_JSON.to_string(); }; serde_json::json!({ "externalGenerationJobId": external_generation_job_id, }) .to_string() } impl fmt::Debug for WalletRefundOutbox { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { f.debug_struct("WalletRefundOutbox") .field("dir", &self.dir) .field("batch_size", &self.batch_size) .field("flush_interval", &self.flush_interval) .field("max_bytes", &self.max_bytes) .finish() } } impl fmt::Display for WalletRefundOutboxError { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { match self { Self::Io(error) => write!(f, "{error}"), Self::Json(error) => write!(f, "{error}"), Self::Spacetime(error) => write!(f, "{error}"), } } } impl From for WalletRefundOutboxError { fn from(value: std::io::Error) -> Self { Self::Io(value) } } impl From for WalletRefundOutboxError { fn from(value: serde_json::Error) -> Self { Self::Json(value) } } impl WalletRefundOutboxError { fn is_data_corruption(&self) -> bool { matches!(self, Self::Json(_)) } } async fn read_refund_record( path: &Path, ) -> Result { let mut file = File::open(path).await?; let mut bytes = Vec::new(); file.read_to_end(&mut bytes).await?; Ok(serde_json::from_slice::(&bytes)?) } fn directory_size_if_exists(path: &Path) -> Result { if !path.is_dir() { return Ok(0); } let mut total = 0u64; for entry in std::fs::read_dir(path)? { let entry = entry?; if !is_pending_outbox_file_name(&entry.file_name()) { continue; } let metadata = entry.metadata()?; if metadata.is_file() { total = total.saturating_add(metadata.len()); } } Ok(total) } fn current_unix_micros() -> u128 { SystemTime::now() .duration_since(UNIX_EPOCH) .unwrap_or_default() .as_micros() } fn ledger_id_hash(ledger_id: &str) -> String { hex::encode(Sha256::digest(ledger_id.as_bytes())) } fn is_pending_outbox_file_name(name: &std::ffi::OsStr) -> bool { name.to_str().is_some_and(|value| { value.starts_with(PENDING_FILE_PREFIX) && value.ends_with(OUTBOX_FILE_EXTENSION) }) } async fn sync_directory_metadata(path: &Path) -> Result<(), WalletRefundOutboxError> { let path = path.to_path_buf(); tokio::task::spawn_blocking(move || { let dir = std::fs::File::open(path)?; dir.sync_all() }) .await .map_err(|error| std::io::Error::new(std::io::ErrorKind::Other, error.to_string()))??; Ok(()) } #[cfg(test)] mod tests { use super::*; fn sample_record(ledger_id: &str) -> WalletRefundOutboxRecord { WalletRefundOutboxRecord { owner_user_id: "user-1".to_string(), amount: 2, ledger_id: ledger_id.to_string(), created_at_micros: 1_713_680_000_000_000, asset_kind: "puzzle_initial_image".to_string(), asset_id: "asset-1".to_string(), external_generation_job_id: Some("extgen-test".to_string()), } } fn test_dir(name: &str) -> PathBuf { let dir = std::env::temp_dir().join(format!( "genarrative-wallet-refund-outbox-{name}-{}", current_unix_micros() )); let _ = std::fs::remove_dir_all(&dir); dir } fn test_outbox(dir: PathBuf, max_bytes: u64) -> Arc { let config = AppConfig { wallet_refund_outbox_dir: dir, wallet_refund_outbox_batch_size: 500, wallet_refund_outbox_flush_interval: Duration::from_secs(60), wallet_refund_outbox_max_bytes: max_bytes, ..AppConfig::default() }; WalletRefundOutbox::from_config( &config, SpacetimeClient::new(spacetime_client::SpacetimeClientConfig { server_url: "http://127.0.0.1:1".to_string(), database: "missing".to_string(), token: None, pool_size: 1, procedure_timeout: Duration::from_millis(10), subscribe_cached_read_models: false, }), ) .expect("outbox should be enabled") } #[tokio::test] async fn enqueue_is_idempotent_per_ledger_id() { let dir = test_dir("idempotent"); let outbox = test_outbox(dir.clone(), 1024 * 1024); outbox.enqueue(sample_record("ledger-1")).await.unwrap(); outbox.enqueue(sample_record("ledger-1")).await.unwrap(); let pending_count = std::fs::read_dir(&dir) .unwrap() .filter_map(Result::ok) .filter(|entry| is_pending_outbox_file_name(&entry.file_name())) .count(); assert_eq!(pending_count, 1); let _ = std::fs::remove_dir_all(dir); } #[tokio::test] async fn enqueue_drops_when_outbox_exceeds_max_bytes() { let dir = test_dir("max-bytes"); let outbox = test_outbox(dir.clone(), 1); let outcome = outbox.enqueue(sample_record("ledger-1")).await.unwrap(); assert!(matches!( outcome, WalletRefundOutboxEnqueueOutcome::Dropped { reason: "max_bytes" } )); assert!(!dir.exists() || std::fs::read_dir(&dir).unwrap().next().is_none()); let _ = std::fs::remove_dir_all(dir); } #[tokio::test] async fn flush_quarantines_corrupt_file() { let dir = test_dir("corrupt"); std::fs::create_dir_all(&dir).unwrap(); let pending_path = dir.join(format!("{PENDING_FILE_PREFIX}bad{OUTBOX_FILE_EXTENSION}")); std::fs::write(&pending_path, b"{not-json}").unwrap(); let outbox = test_outbox(dir.clone(), 1024 * 1024); outbox.flush_pending_files_once().await.unwrap(); assert!(!pending_path.exists()); let corrupt_count = std::fs::read_dir(&dir) .unwrap() .filter_map(Result::ok) .filter(|entry| { entry .file_name() .to_str() .is_some_and(|name| name.starts_with(CORRUPT_FILE_PREFIX)) }) .count(); assert_eq!(corrupt_count, 1); let _ = std::fs::remove_dir_all(dir); } #[tokio::test] async fn shutdown_flush_keeps_file_when_spacetime_is_unavailable() { let dir = test_dir("shutdown"); let outbox = test_outbox(dir.clone(), 1024 * 1024); outbox.enqueue(sample_record("ledger-1")).await.unwrap(); let result = outbox.flush_for_shutdown().await; assert!( matches!(result, Err(WalletRefundOutboxError::Spacetime(_))), "missing test SpacetimeDB should keep refund file for retry" ); let pending_count = std::fs::read_dir(&dir) .unwrap() .filter_map(Result::ok) .filter(|entry| is_pending_outbox_file_name(&entry.file_name())) .count(); assert_eq!(pending_count, 1); let _ = std::fs::remove_dir_all(dir); } }