f1ca3e79c1
worker 和 controller 改为订阅 external_generation_job 队列变更来唤醒 claim 与扩缩容评估 非 HTTP 角色关闭 API 读模型订阅并将 SpacetimeDB 连接池收敛为 1 更新 worker/controller 环境模板和运维记忆,明确 poll interval 只作兜底
484 lines
16 KiB
Rust
484 lines
16 KiB
Rust
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<Mutex<()>>,
|
|
flush_notify: Arc<Notify>,
|
|
}
|
|
|
|
#[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<String>,
|
|
}
|
|
|
|
#[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<Arc<Self>> {
|
|
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<WalletRefundOutboxEnqueueOutcome, WalletRefundOutboxError> {
|
|
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<Self>) {
|
|
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<Vec<PathBuf>, 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<std::io::Error> for WalletRefundOutboxError {
|
|
fn from(value: std::io::Error) -> Self {
|
|
Self::Io(value)
|
|
}
|
|
}
|
|
|
|
impl From<serde_json::Error> 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<WalletRefundOutboxRecord, WalletRefundOutboxError> {
|
|
let mut file = File::open(path).await?;
|
|
let mut bytes = Vec::new();
|
|
file.read_to_end(&mut bytes).await?;
|
|
Ok(serde_json::from_slice::<WalletRefundOutboxRecord>(&bytes)?)
|
|
}
|
|
|
|
fn directory_size_if_exists(path: &Path) -> Result<u64, std::io::Error> {
|
|
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<WalletRefundOutbox> {
|
|
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);
|
|
}
|
|
}
|