From e25e240d16b68eaf9ca68b280fcf1f47cc856e4a Mon Sep 17 00:00:00 2001 From: RWDai <27391645+RWDai@users.noreply.github.com> Date: Wed, 23 Sep 2026 11:30:33 +0800 Subject: [PATCH 1/9] feat(admin): batch adjust user wallet balances --- apps/aether-gateway/src/data/state/runtime.rs | 3 +- .../admin/billing/wallets/mutations/adjust.rs | 3 +- .../src/handlers/admin/request/billing.rs | 4 +- .../src/handlers/admin/users/batch.rs | 95 +++++++ .../state/runtime/wallet/balance_mutations.rs | 19 +- .../src/state/runtime/wallet/mutations.rs | 2 +- .../aether-gateway/src/tests/control/admin.rs | 1 + .../src/tests/control/admin/users_batch.rs | 232 ++++++++++++++++++ .../adapters/postgres/src/wallet.rs | 158 ++++++++++-- .../contracts/src/repository/wallet/types.rs | 7 +- .../runtime/src/repository/wallet/memory.rs | 9 +- frontend/src/api/__tests__/users.spec.ts | 24 +- frontend/src/api/users.ts | 30 ++- .../components/UserBatchActionDialog.vue | 83 ++++++- .../components/user-management-config.ts | 7 + frontend/src/i18n/messages.ts | 11 + 16 files changed, 655 insertions(+), 33 deletions(-) create mode 100644 apps/aether-gateway/src/tests/control/admin/users_batch.rs diff --git a/apps/aether-gateway/src/data/state/runtime.rs b/apps/aether-gateway/src/data/state/runtime.rs index 901563379..b06449a80 100644 --- a/apps/aether-gateway/src/data/state/runtime.rs +++ b/apps/aether-gateway/src/data/state/runtime.rs @@ -1066,7 +1066,8 @@ impl GatewayDataState { pub(crate) async fn adjust_wallet_balance( &self, input: AdjustWalletBalanceInput, - ) -> Result, DataLayerError> { + ) -> Result)>, DataLayerError> + { match &self.wallet_writer { Some(repository) => repository.adjust_wallet_balance(input).await, None => Ok(None), diff --git a/apps/aether-gateway/src/handlers/admin/billing/wallets/mutations/adjust.rs b/apps/aether-gateway/src/handlers/admin/billing/wallets/mutations/adjust.rs index 8efabaa0b..f3d2809af 100644 --- a/apps/aether-gateway/src/handlers/admin/billing/wallets/mutations/adjust.rs +++ b/apps/aether-gateway/src/handlers/admin/billing/wallets/mutations/adjust.rs @@ -63,13 +63,14 @@ pub(in super::super) async fn build_admin_wallet_adjust_response( } let operator_id = admin_wallet_operator_id(request_context); let has_wallet_writer = state.has_wallet_data_writer(); - let Some((wallet, transaction)) = state + let Some((wallet, Some(transaction))) = state .admin_adjust_wallet_balance( &wallet_id, amount_usd, &balance_type, operator_id.as_deref(), description.as_deref(), + false, ) .await? else { diff --git a/apps/aether-gateway/src/handlers/admin/request/billing.rs b/apps/aether-gateway/src/handlers/admin/request/billing.rs index 3ecde6481..38a7b9ddf 100644 --- a/apps/aether-gateway/src/handlers/admin/request/billing.rs +++ b/apps/aether-gateway/src/handlers/admin/request/billing.rs @@ -387,10 +387,11 @@ impl<'a> AdminAppState<'a> { balance_type: &str, operator_id: Option<&str>, description: Option<&str>, + clamp_deduction_to_available_balance: bool, ) -> Result< Option<( aether_data::repository::wallet::StoredWalletSnapshot, - crate::AdminWalletTransactionRecord, + Option, )>, GatewayError, > { @@ -401,6 +402,7 @@ impl<'a> AdminAppState<'a> { balance_type, operator_id, description, + clamp_deduction_to_available_balance, ) .await } diff --git a/apps/aether-gateway/src/handlers/admin/users/batch.rs b/apps/aether-gateway/src/handlers/admin/users/batch.rs index 75d945001..c6635ef6c 100644 --- a/apps/aether-gateway/src/handlers/admin/users/batch.rs +++ b/apps/aether-gateway/src/handlers/admin/users/batch.rs @@ -88,6 +88,7 @@ struct AdminUserBatchMutation { role: Option, is_active: Option, unlimited: Option, + wallet_balance_adjustment: Option, modified_fields: Vec<&'static str>, } @@ -97,6 +98,18 @@ impl AdminUserBatchMutation { } } +#[derive(Debug, Clone, Copy)] +struct AdminUserWalletBalanceAdjustment { + operation: AdminUserWalletBalanceOperation, + amount: f64, +} + +#[derive(Debug, Clone, Copy)] +enum AdminUserWalletBalanceOperation { + Add, + Deduct, +} + pub(in super::super) async fn build_admin_resolve_user_selection_response( state: &AdminAppState<'_>, _request_context: &AdminRequestContext<'_>, @@ -160,6 +173,11 @@ pub(in super::super) async fn build_admin_user_batch_action_response( "当前为只读模式,无法批量更新用户钱包", )); } + if mutation.wallet_balance_adjustment.is_some() && !state.has_auth_wallet_write_capability() { + return Ok(build_admin_users_read_only_response( + "当前为只读模式,无法批量调整用户钱包余额", + )); + } let active_admin_demotions = count_active_admin_demotions(&mutation, &resolved.items); let active_admin_count = if active_admin_demotions > 0 { state.count_active_admin_users().await? @@ -211,6 +229,23 @@ pub(in super::super) async fn build_admin_user_batch_action_response( } } + if let Some(adjustment) = mutation.wallet_balance_adjustment { + if !apply_batch_user_wallet_balance_adjustment( + state, + &item.user_id, + adjustment, + current_admin_user_id, + ) + .await? + { + failures.push(json!({ + "user_id": item.user_id, + "reason": "用户钱包不可用", + })); + continue; + } + } + if mutation.has_auth_user_fields() && state .update_local_auth_user_admin_fields( @@ -604,10 +639,37 @@ fn parse_batch_mutation( }), "update_access_control" => parse_access_control_mutation(payload), "update_role" => parse_role_mutation(payload), + "adjust_wallet_balance" => parse_wallet_balance_adjustment_mutation(payload), _ => Err("不支持的批量操作".to_string()), } } +fn parse_wallet_balance_adjustment_mutation( + payload: Option, +) -> Result { + let Some(Value::Object(payload)) = payload else { + return Err("payload 必须是对象".to_string()); + }; + let operation = match payload.get("operation").and_then(Value::as_str) { + Some("add") => AdminUserWalletBalanceOperation::Add, + Some("deduct") => AdminUserWalletBalanceOperation::Deduct, + _ => return Err("operation 必须为 add 或 deduct".to_string()), + }; + let amount = payload + .get("amount") + .and_then(Value::as_f64) + .ok_or_else(|| "amount 必须为大于 0 的有限数字".to_string())?; + if !amount.is_finite() || amount <= 0.0 { + return Err("amount 必须为大于 0 的有限数字".to_string()); + } + + Ok(AdminUserBatchMutation { + wallet_balance_adjustment: Some(AdminUserWalletBalanceAdjustment { operation, amount }), + modified_fields: vec!["wallet_balance"], + ..AdminUserBatchMutation::default() + }) +} + fn parse_role_mutation(payload: Option) -> Result { let Some(Value::Object(payload)) = payload else { return Err("payload 必须是对象".to_string()); @@ -724,6 +786,39 @@ async fn apply_batch_user_wallet_limit_mode( } } +async fn apply_batch_user_wallet_balance_adjustment( + state: &AdminAppState<'_>, + user_id: &str, + adjustment: AdminUserWalletBalanceAdjustment, + operator_id: Option<&str>, +) -> Result { + // Resolve only the wallet ID; the repository clamps the deduction under its row lock. + let Some(wallet) = state + .find_wallet(aether_data::repository::wallet::WalletLookupKey::UserId( + user_id, + )) + .await? + else { + return Ok(false); + }; + let amount = match adjustment.operation { + AdminUserWalletBalanceOperation::Add => adjustment.amount, + AdminUserWalletBalanceOperation::Deduct => -adjustment.amount, + }; + + Ok(state + .admin_adjust_wallet_balance( + &wallet.id, + amount, + "recharge", + operator_id, + Some("管理员批量调整用户余额"), + true, + ) + .await? + .is_some()) +} + fn build_admin_user_batch_bad_request_response(detail: String) -> Response { if detail.as_str() == "缺少 user_id" { return build_admin_users_bad_request_response("缺少 user_id"); diff --git a/apps/aether-gateway/src/state/runtime/wallet/balance_mutations.rs b/apps/aether-gateway/src/state/runtime/wallet/balance_mutations.rs index 8501b155a..f1e170669 100644 --- a/apps/aether-gateway/src/state/runtime/wallet/balance_mutations.rs +++ b/apps/aether-gateway/src/state/runtime/wallet/balance_mutations.rs @@ -10,10 +10,11 @@ impl AppState { balance_type: &str, operator_id: Option<&str>, description: Option<&str>, + clamp_deduction_to_available_balance: bool, ) -> Result< Option<( aether_data::repository::wallet::StoredWalletSnapshot, - AdminWalletTransactionRecord, + Option, )>, GatewayError, > { @@ -27,6 +28,14 @@ impl AppState { let before_recharge = wallet.balance; let before_gift = wallet.gift_balance; let before_total = before_recharge + before_gift; + let amount_usd = if clamp_deduction_to_available_balance && amount_usd < 0.0 { + -(-amount_usd).min(before_total.max(0.0)) + } else { + amount_usd + }; + if amount_usd == 0.0 { + return Ok(Some((wallet.clone(), None))); + } let mut after_recharge = before_recharge; let mut after_gift = before_gift; @@ -90,7 +99,7 @@ impl AppState { let updated_wallet = wallet.clone(); drop(guard); self.invalidate_auth_context_cache(); - return Ok(Some((updated_wallet, transaction))); + return Ok(Some((updated_wallet, Some(transaction)))); } Ok(self @@ -100,10 +109,14 @@ impl AppState { balance_type: balance_type.to_string(), operator_id: operator_id.map(ToOwned::to_owned), description: description.map(ToOwned::to_owned), + clamp_deduction_to_available_balance, }) .await? .map(|(wallet, transaction)| { - (wallet, stored_wallet_transaction_to_gateway(transaction)) + ( + wallet, + transaction.map(stored_wallet_transaction_to_gateway), + ) })) } diff --git a/apps/aether-gateway/src/state/runtime/wallet/mutations.rs b/apps/aether-gateway/src/state/runtime/wallet/mutations.rs index 995ba626f..40b4960a8 100644 --- a/apps/aether-gateway/src/state/runtime/wallet/mutations.rs +++ b/apps/aether-gateway/src/state/runtime/wallet/mutations.rs @@ -112,7 +112,7 @@ impl AppState { ) -> Result< Option<( aether_data::repository::wallet::StoredWalletSnapshot, - aether_data::repository::wallet::StoredAdminWalletTransaction, + Option, )>, GatewayError, > { diff --git a/apps/aether-gateway/src/tests/control/admin.rs b/apps/aether-gateway/src/tests/control/admin.rs index 442d422e7..ab54ff946 100644 --- a/apps/aether-gateway/src/tests/control/admin.rs +++ b/apps/aether-gateway/src/tests/control/admin.rs @@ -21,5 +21,6 @@ mod system; mod system_import; mod usage; mod users; +mod users_batch; mod video_tasks; mod wallets; diff --git a/apps/aether-gateway/src/tests/control/admin/users_batch.rs b/apps/aether-gateway/src/tests/control/admin/users_batch.rs new file mode 100644 index 000000000..7abb4d7bc --- /dev/null +++ b/apps/aether-gateway/src/tests/control/admin/users_batch.rs @@ -0,0 +1,232 @@ +use aether_data::repository::users::StoredUserAuthRecord; +use aether_data::repository::wallet::StoredWalletSnapshot; +use axum::http::StatusCode; +use chrono::Utc; +use reqwest::{Client, RequestBuilder, Response}; +use serde_json::{json, Value}; + +use super::super::{build_router_with_state, start_server, AppState}; + +fn admin_headers(request: RequestBuilder) -> RequestBuilder { + request + .header(crate::constants::GATEWAY_HEADER, "rust-phase3b") + .header(crate::constants::TRUSTED_ADMIN_USER_ID_HEADER, "admin-user") + .header(crate::constants::TRUSTED_ADMIN_USER_ROLE_HEADER, "admin") + .header( + crate::constants::TRUSTED_ADMIN_SESSION_ID_HEADER, + "session-admin", + ) +} + +fn sample_user(user_id: &str) -> StoredUserAuthRecord { + StoredUserAuthRecord::new( + user_id.to_string(), + Some(format!("{user_id}@example.com")), + true, + user_id.to_string(), + Some("hash".to_string()), + "user".to_string(), + "local".to_string(), + Some(json!(["openai"])), + Some(json!(["openai:chat"])), + Some(json!(["gpt-4.1"])), + true, + false, + Some(Utc::now()), + Some(Utc::now()), + ) + .expect("test user should build") +} + +fn sample_wallet(user_id: &str, balance: f64, gift_balance: f64) -> StoredWalletSnapshot { + StoredWalletSnapshot::new( + format!("wallet-{user_id}"), + Some(user_id.to_string()), + None, + balance, + gift_balance, + "finite".to_string(), + "USD".to_string(), + "active".to_string(), + balance.max(0.0), + 0.0, + 0.0, + 0.0, + 1_710_000_000, + ) + .expect("test wallet should build") +} + +async fn post_batch_action(client: &Client, gateway_url: &str, payload: Value) -> Response { + admin_headers(client.post(format!("{gateway_url}/api/admin/users/batch-action"))) + .json(&payload) + .send() + .await + .expect("batch request should complete") +} + +async fn wallet_detail(client: &Client, gateway_url: &str, user_id: &str) -> Value { + admin_headers(client.get(format!("{gateway_url}/api/admin/wallets/wallet-{user_id}"))) + .send() + .await + .expect("wallet lookup should complete") + .json() + .await + .expect("wallet response should parse") +} + +#[tokio::test] +async fn gateway_batches_wallet_addition_deduction_and_clamped_deduction_per_user() { + let state = AppState::new() + .expect("gateway should build") + .with_auth_users_for_tests([sample_user("user-1"), sample_user("user-2")]) + .with_auth_wallets_for_tests([ + sample_wallet("user-1", 10.0, 3.0), + sample_wallet("user-2", 2.0, 1.0), + ]); + let (gateway_url, gateway_handle) = start_server(build_router_with_state(state)).await; + let client = Client::new(); + let selection = json!({ "user_ids": ["user-1", "user-2"] }); + + let add_response = post_batch_action( + &client, + &gateway_url, + json!({ + "selection": selection.clone(), + "action": "adjust_wallet_balance", + "payload": { "operation": "add", "amount": 5.0 } + }), + ) + .await; + assert_eq!(add_response.status(), StatusCode::OK); + let add_result: Value = add_response.json().await.expect("response should parse"); + assert_eq!(add_result["success"], 2); + assert_eq!(add_result["failed"], 0); + assert_eq!(add_result["modified_fields"], json!(["wallet_balance"])); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 18.0 + ); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-2").await["balance"], + 8.0 + ); + + let deduct_response = post_batch_action( + &client, + &gateway_url, + json!({ + "selection": selection.clone(), + "action": "adjust_wallet_balance", + "payload": { "operation": "deduct", "amount": 4.0 } + }), + ) + .await; + assert_eq!(deduct_response.status(), StatusCode::OK); + let deduct_result: Value = deduct_response.json().await.expect("response should parse"); + assert_eq!(deduct_result["success"], 2); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 14.0 + ); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-2").await["balance"], + 4.0 + ); + + let over_deduct_response = post_batch_action( + &client, + &gateway_url, + json!({ + "selection": selection, + "action": "adjust_wallet_balance", + "payload": { "operation": "deduct", "amount": 100.0 } + }), + ) + .await; + assert_eq!(over_deduct_response.status(), StatusCode::OK); + let over_deduct_result: Value = over_deduct_response + .json() + .await + .expect("response should parse"); + assert_eq!(over_deduct_result["success"], 2); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 0.0 + ); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-2").await["balance"], + 0.0 + ); + + gateway_handle.abort(); +} + +#[tokio::test] +async fn gateway_reports_missing_wallet_and_skips_zero_delta_for_non_positive_balance() { + let state = AppState::new() + .expect("gateway should build") + .with_auth_users_for_tests([sample_user("user-negative"), sample_user("user-no-wallet")]) + .with_auth_wallets_for_tests([sample_wallet("user-negative", -2.0, 1.0)]); + let (gateway_url, gateway_handle) = start_server(build_router_with_state(state)).await; + let client = Client::new(); + + let response = post_batch_action( + &client, + &gateway_url, + json!({ + "selection": { "user_ids": ["user-negative", "user-no-wallet"] }, + "action": "adjust_wallet_balance", + "payload": { "operation": "deduct", "amount": 10.0 } + }), + ) + .await; + assert_eq!(response.status(), StatusCode::OK); + let result: Value = response.json().await.expect("response should parse"); + assert_eq!(result["success"], 1); + assert_eq!(result["failed"], 1); + assert_eq!(result["failures"][0]["user_id"], "user-no-wallet"); + assert_eq!(result["failures"][0]["reason"], "用户钱包不可用"); + + let wallet = wallet_detail(&client, &gateway_url, "user-negative").await; + assert_eq!(wallet["balance"], -1.0); + assert_eq!(wallet["total_adjusted"], 0.0); + + gateway_handle.abort(); +} + +#[tokio::test] +async fn gateway_rejects_zero_and_non_finite_batch_wallet_adjustments() { + let state = AppState::new() + .expect("gateway should build") + .with_auth_users_for_tests([sample_user("user-1")]) + .with_auth_wallets_for_tests([sample_wallet("user-1", 10.0, 0.0)]); + let (gateway_url, gateway_handle) = start_server(build_router_with_state(state)).await; + let client = Client::new(); + + let zero_response = post_batch_action( + &client, + &gateway_url, + json!({ + "selection": { "user_ids": ["user-1"] }, + "action": "adjust_wallet_balance", + "payload": { "operation": "add", "amount": 0.0 } + }), + ) + .await; + assert_eq!(zero_response.status(), StatusCode::BAD_REQUEST); + + let non_finite_response = admin_headers( + client.post(format!("{gateway_url}/api/admin/users/batch-action")), + ) + .header(reqwest::header::CONTENT_TYPE, "application/json") + .body( + r#"{"selection":{"user_ids":["user-1"]},"action":"adjust_wallet_balance","payload":{"operation":"add","amount":1e999}}"#, + ) + .send() + .await + .expect("non-finite amount request should complete"); + assert_eq!(non_finite_response.status(), StatusCode::BAD_REQUEST); + + gateway_handle.abort(); +} diff --git a/crates/aether-data/adapters/postgres/src/wallet.rs b/crates/aether-data/adapters/postgres/src/wallet.rs index eb9311f8f..4d26e8b35 100644 --- a/crates/aether-data/adapters/postgres/src/wallet.rs +++ b/crates/aether-data/adapters/postgres/src/wallet.rs @@ -807,6 +807,14 @@ impl SqlxWalletRepository { } } +fn effective_wallet_adjustment_amount(input: &AdjustWalletBalanceInput, before_total: f64) -> f64 { + if input.clamp_deduction_to_available_balance && input.amount_usd < 0.0 { + -(-input.amount_usd).min(before_total.max(0.0)) + } else { + input.amount_usd + } +} + #[async_trait] impl WalletReadRepository for SqlxWalletRepository { async fn find( @@ -4262,7 +4270,8 @@ RETURNING async fn adjust_wallet_balance( &self, input: AdjustWalletBalanceInput, - ) -> Result, DataLayerError> { + ) -> Result)>, DataLayerError> + { if !input.amount_usd.is_finite() || input.amount_usd == 0.0 { return Err(DataLayerError::InvalidInput( "adjustment amount must be finite and non-zero".to_string(), @@ -4285,7 +4294,8 @@ SELECT CAST(total_recharged AS DOUBLE PRECISION) AS total_recharged, CAST(total_consumed AS DOUBLE PRECISION) AS total_consumed, CAST(total_refunded AS DOUBLE PRECISION) AS total_refunded, - CAST(total_adjusted AS DOUBLE PRECISION) AS total_adjusted + CAST(total_adjusted AS DOUBLE PRECISION) AS total_adjusted, + CAST(EXTRACT(EPOCH FROM updated_at) AS BIGINT) AS updated_at_unix_secs FROM wallets WHERE id = $1 FOR UPDATE @@ -4312,17 +4322,21 @@ FOR UPDATE "wallet balance is invalid".to_string(), )); } + let amount_usd = effective_wallet_adjustment_amount(&input, before_total); + if amount_usd == 0.0 { + return Ok(Some((map_wallet_row(&row)?, None))); + } let mut after_recharge = before_recharge; let mut after_gift = before_gift; - if input.amount_usd > 0.0 { + if amount_usd > 0.0 { if input.balance_type.eq_ignore_ascii_case("gift") { - after_gift += input.amount_usd; + after_gift += amount_usd; } else { - after_recharge += input.amount_usd; + after_recharge += amount_usd; } } else { - let mut remaining = -input.amount_usd; + let mut remaining = -amount_usd; let consume_positive_bucket = |balance: &mut f64, to_consume: &mut f64| { if *to_consume <= 0.0 { return; @@ -4344,7 +4358,7 @@ FOR UPDATE } } let after_total = after_recharge + after_gift; - let after_total_adjusted = before_total_adjusted + input.amount_usd; + let after_total_adjusted = before_total_adjusted + amount_usd; if !after_recharge.is_finite() || !after_gift.is_finite() || !after_total.is_finite() @@ -4383,7 +4397,7 @@ RETURNING .bind(&input.wallet_id) .bind(after_recharge) .bind(after_gift) - .bind(input.amount_usd) + .bind(amount_usd) .fetch_one(&mut **tx) .await .map_postgres_err()?; @@ -4439,7 +4453,7 @@ VALUES ( ) .bind(&transaction_id) .bind(&input.wallet_id) - .bind(input.amount_usd) + .bind(amount_usd) .bind(before_total) .bind(after_total) .bind(before_recharge) @@ -4455,12 +4469,12 @@ VALUES ( Ok(Some(( wallet, - StoredAdminWalletTransaction { + Some(StoredAdminWalletTransaction { id: transaction_id, wallet_id: input.wallet_id, category: "adjust".to_string(), reason_code: "adjust_admin".to_string(), - amount: input.amount_usd, + amount: amount_usd, balance_before: before_total, balance_after: after_total, recharge_balance_before: before_recharge, @@ -4474,7 +4488,7 @@ VALUES ( operator_email: None, description: Some(description), created_at_unix_ms: Some(created_at), - }, + }), ))) }) }) @@ -8489,13 +8503,14 @@ VALUES ($1, $2, 'gift', 'gift_initial', $3, 0, $3, 0, 0, 0, $3, 'system_task', $ #[cfg(test)] mod tests { use aether_data_contracts::repository::wallet::{ - CreateManualWalletRechargeInput, CreditAdminPaymentOrderInput, ProcessPaymentCallbackInput, - ProcessPaymentCallbackOutcome, RedeemWalletCodeInput, RedeemWalletCodeOutcome, - WalletLookupKey, WalletMutationOutcome, WalletReadRepository, WalletWriteRepository, + AdjustWalletBalanceInput, CreateManualWalletRechargeInput, CreditAdminPaymentOrderInput, + ProcessPaymentCallbackInput, ProcessPaymentCallbackOutcome, RedeemWalletCodeInput, + RedeemWalletCodeOutcome, WalletLookupKey, WalletMutationOutcome, WalletReadRepository, + WalletWriteRepository, }; use sqlx::Row; - use super::SqlxWalletRepository; + use super::{effective_wallet_adjustment_amount, SqlxWalletRepository}; use crate::{PostgresPoolConfig, PostgresPoolFactory}; #[test] @@ -9221,6 +9236,117 @@ mod tests { pool.close().await; } + #[test] + fn bulk_adjustment_clamp_is_opt_in_and_uses_available_total() { + let input = AdjustWalletBalanceInput { + wallet_id: "wallet-1".to_string(), + amount_usd: -100.0, + balance_type: "recharge".to_string(), + operator_id: None, + description: None, + clamp_deduction_to_available_balance: true, + }; + assert_eq!(effective_wallet_adjustment_amount(&input, 13.0), -13.0); + assert_eq!(effective_wallet_adjustment_amount(&input, 0.0), -0.0); + assert_eq!(effective_wallet_adjustment_amount(&input, -1.0), -0.0); + + let legacy_input = AdjustWalletBalanceInput { + clamp_deduction_to_available_balance: false, + ..input + }; + assert_eq!( + effective_wallet_adjustment_amount(&legacy_input, 13.0), + -100.0 + ); + } + + #[tokio::test] + #[ignore = "requires AETHER_TEST_DATABASE_URL and PostgreSQL bootstrap schema"] + async fn live_bulk_wallet_adjustment_persists_actual_delta_and_skips_zero_ledger() { + let pool = isolated_wallet_test_pool().await; + let (wallet_id, _) = seed_wallet(&pool).await; + let repository = SqlxWalletRepository::new(pool.clone()); + + let (wallet, transaction) = repository + .adjust_wallet_balance(AdjustWalletBalanceInput { + wallet_id: wallet_id.clone(), + amount_usd: -100.0, + balance_type: "recharge".to_string(), + operator_id: Some("admin-user".to_string()), + description: Some("bulk deduction".to_string()), + clamp_deduction_to_available_balance: true, + }) + .await + .expect("bulk adjustment should succeed") + .expect("wallet should exist"); + let transaction = + transaction.expect("positive available balance should create a ledger row"); + assert_eq!(transaction.amount, -13.0); + assert_eq!(transaction.balance_before, 13.0); + assert_eq!(transaction.balance_after, 0.0); + assert_eq!(wallet.balance + wallet.gift_balance, 0.0); + let persisted_amount: f64 = + sqlx::query_scalar("SELECT amount FROM wallet_transactions WHERE id = $1") + .bind(&transaction.id) + .fetch_one(&pool) + .await + .expect("ledger should store the effective deduction"); + assert_eq!(persisted_amount, -13.0); + + let (wallet, transaction) = repository + .adjust_wallet_balance(AdjustWalletBalanceInput { + wallet_id: wallet_id.clone(), + amount_usd: -1.0, + balance_type: "recharge".to_string(), + operator_id: Some("admin-user".to_string()), + description: Some("bulk deduction at zero".to_string()), + clamp_deduction_to_available_balance: true, + }) + .await + .expect("zero-balance adjustment should succeed") + .expect("wallet should still exist"); + assert_eq!(wallet.balance + wallet.gift_balance, 0.0); + assert!( + transaction.is_none(), + "zero effective delta must not create a ledger row" + ); + let transaction_count: i64 = + sqlx::query_scalar("SELECT COUNT(*) FROM wallet_transactions WHERE wallet_id = $1") + .bind(&wallet_id) + .fetch_one(&pool) + .await + .expect("ledger row count should be readable"); + assert_eq!(transaction_count, 1); + + sqlx::query("UPDATE wallets SET balance = -2, gift_balance = 1 WHERE id = $1") + .bind(&wallet_id) + .execute(&pool) + .await + .expect("legacy negative wallet balance should be seeded"); + let (wallet, transaction) = repository + .adjust_wallet_balance(AdjustWalletBalanceInput { + wallet_id: wallet_id.clone(), + amount_usd: -1.0, + balance_type: "recharge".to_string(), + operator_id: Some("admin-user".to_string()), + description: Some("bulk deduction from negative balance".to_string()), + clamp_deduction_to_available_balance: true, + }) + .await + .expect("legacy negative wallet should remain usable") + .expect("wallet should still exist"); + assert_eq!(wallet.balance + wallet.gift_balance, -1.0); + assert!(transaction.is_none()); + let transaction_count: i64 = + sqlx::query_scalar("SELECT COUNT(*) FROM wallet_transactions WHERE wallet_id = $1") + .bind(&wallet_id) + .fetch_one(&pool) + .await + .expect("ledger row count should be readable"); + assert_eq!(transaction_count, 1); + pool.close().await; + } + #[tokio::test] async fn repository_constructs_from_lazy_pool() { let factory = PostgresPoolFactory::new(PostgresPoolConfig { diff --git a/crates/aether-data/contracts/src/repository/wallet/types.rs b/crates/aether-data/contracts/src/repository/wallet/types.rs index 25a7e8ea5..f391d49d2 100644 --- a/crates/aether-data/contracts/src/repository/wallet/types.rs +++ b/crates/aether-data/contracts/src/repository/wallet/types.rs @@ -2462,6 +2462,8 @@ pub struct AdjustWalletBalanceInput { pub balance_type: String, pub operator_id: Option, pub description: Option, + #[serde(default)] + pub clamp_deduction_to_available_balance: bool, } #[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] @@ -2893,7 +2895,10 @@ pub trait WalletWriteRepository: Send + Sync { async fn adjust_wallet_balance( &self, input: AdjustWalletBalanceInput, - ) -> Result, crate::DataLayerError>; + ) -> Result< + Option<(StoredWalletSnapshot, Option)>, + crate::DataLayerError, + >; async fn create_manual_wallet_recharge( &self, diff --git a/crates/aether-data/runtime/src/repository/wallet/memory.rs b/crates/aether-data/runtime/src/repository/wallet/memory.rs index 08df3efc8..73a47063a 100644 --- a/crates/aether-data/runtime/src/repository/wallet/memory.rs +++ b/crates/aether-data/runtime/src/repository/wallet/memory.rs @@ -2322,8 +2322,13 @@ impl WalletWriteRepository for InMemoryWalletRepository { async fn adjust_wallet_balance( &self, _input: AdjustWalletBalanceInput, - ) -> Result, DataLayerError> - { + ) -> Result< + Option<( + StoredWalletSnapshot, + Option, + )>, + DataLayerError, + > { Ok(None) } diff --git a/frontend/src/api/__tests__/users.spec.ts b/frontend/src/api/__tests__/users.spec.ts index 3f2f8d486..f4dc72e52 100644 --- a/frontend/src/api/__tests__/users.spec.ts +++ b/frontend/src/api/__tests__/users.spec.ts @@ -17,7 +17,7 @@ vi.mock('@/utils/cache', () => ({ cachedRequest: cachedRequestMock, })) -import { usersApi } from '@/api/users' +import { buildUserBatchBalanceAdjustmentPayload, usersApi } from '@/api/users' describe('usersApi admin list query', () => { beforeEach(() => { @@ -104,3 +104,25 @@ describe('usersApi admin list query', () => { expect(getMock).toHaveBeenCalledWith('/api/admin/users/target-user/api-keys') }) }) + +describe('user batch wallet balance payload', () => { + it('builds an addition payload from a valid positive amount', () => { + expect(buildUserBatchBalanceAdjustmentPayload('add', '12.5')).toEqual({ + operation: 'add', + amount: 12.5, + }) + }) + + it('builds a deduction payload from a valid positive amount', () => { + expect(buildUserBatchBalanceAdjustmentPayload('deduct', 4)).toEqual({ + operation: 'deduct', + amount: 4, + }) + }) + + it('rejects zero, blank, and non-finite amounts', () => { + expect(buildUserBatchBalanceAdjustmentPayload('add', '0')).toBeNull() + expect(buildUserBatchBalanceAdjustmentPayload('deduct', '')).toBeNull() + expect(buildUserBatchBalanceAdjustmentPayload('deduct', '1e999')).toBeNull() + }) +}) diff --git a/frontend/src/api/users.ts b/frontend/src/api/users.ts index 8a9b8d6a7..ee2121c16 100644 --- a/frontend/src/api/users.ts +++ b/frontend/src/api/users.ts @@ -120,9 +120,28 @@ export interface UserBatchRolePayload { role: UserRole } -export type UserBatchAction = 'enable' | 'disable' | 'update_access_control' | 'update_role' +export type UserBatchBalanceOperation = 'add' | 'deduct' -export type UserBatchActionPayload = UserBatchAccessControlPayload | UserBatchRolePayload +export interface UserBatchBalanceAdjustmentPayload { + operation: UserBatchBalanceOperation + amount: number +} + +export function buildUserBatchBalanceAdjustmentPayload( + operation: UserBatchBalanceOperation, + amountInput: string | number, +): UserBatchBalanceAdjustmentPayload | null { + const amount = typeof amountInput === 'number' ? amountInput : Number(amountInput.trim()) + if (!Number.isFinite(amount) || amount <= 0) return null + return { operation, amount } +} + +export type UserBatchAction = 'enable' | 'disable' | 'update_access_control' | 'update_role' | 'adjust_wallet_balance' + +export type UserBatchActionPayload = + | UserBatchAccessControlPayload + | UserBatchRolePayload + | UserBatchBalanceAdjustmentPayload export interface UserBatchToggleActionRequest { selection: UserBatchSelection @@ -142,10 +161,17 @@ export interface UserBatchRoleActionRequest { payload: UserBatchRolePayload } +export interface UserBatchBalanceActionRequest { + selection: UserBatchSelection + action: 'adjust_wallet_balance' + payload: UserBatchBalanceAdjustmentPayload +} + export type UserBatchActionRequest = | UserBatchToggleActionRequest | UserBatchAccessControlActionRequest | UserBatchRoleActionRequest + | UserBatchBalanceActionRequest export interface UserBatchActionFailure { user_id: string diff --git a/frontend/src/features/users/components/UserBatchActionDialog.vue b/frontend/src/features/users/components/UserBatchActionDialog.vue index fa87b7748..282e51e9b 100644 --- a/frontend/src/features/users/components/UserBatchActionDialog.vue +++ b/frontend/src/features/users/components/UserBatchActionDialog.vue @@ -2,7 +2,7 @@ + + + + {{ legacyT('调整金额 (USD)') }} + + + + + {{ legacyT('增加') }} + + + + {{ legacyT('扣减') }} + + + + + + {{ legacyT('请输入大于 0 的有限金额') }} + + + {{ legacyT('扣减超过单个用户可用余额时,该用户余额将归零。') }} + + + ('enable') const targetRole = ref('user') const quotaMode = ref('skip') +const balanceOperation = ref('add') +const balanceAmount = ref('') const selectedGroupIds = ref([]) const previewLoading = ref(false) const previewItems = ref([]) @@ -125,7 +181,16 @@ const lastResult = ref(null) const hasAnyTarget = computed(() => props.selectedCount > 0 || selectedGroupIds.value.length > 0) const impactCount = computed(() => resolvedTotal.value ?? props.selectedCount) -const canExecute = computed(() => hasAnyTarget.value && !previewLoading.value && !executing.value) +const balancePayload = computed(() => buildUserBatchBalanceAdjustmentPayload( + balanceOperation.value, + balanceAmount.value, +)) +const canExecute = computed(() => ( + hasAnyTarget.value + && !previewLoading.value + && !executing.value + && (selectedAction.value !== 'adjust_wallet_balance' || balancePayload.value !== null) +)) const selectedActionLabel = computed(() => ( USER_BATCH_ACTION_OPTIONS.find((action) => action.value === selectedAction.value)?.label ?? '批量操作' )) @@ -149,7 +214,9 @@ const lastResultLabel = computed(() => { }) const lastResultFailuresLabel = computed(() => { if (!lastResult.value || lastResult.value.failures.length === 0) return '' - const failures = lastResult.value.failures.slice(0, 3).map((item) => `${item.user_id} ${item.reason}`).join(locale.value === 'en-US' ? '; ' : ';') + const failures = lastResult.value.failures.slice(0, 3) + .map((item) => `${item.user_id} ${legacyT(item.reason)}`) + .join(locale.value === 'en-US' ? '; ' : ';') return locale.value === 'en-US' ? `: ${failures}` : `:${failures}` }) @@ -177,6 +244,8 @@ function resetLocalState(): void { selectedAction.value = 'enable' targetRole.value = 'user' quotaMode.value = 'skip' + balanceOperation.value = 'add' + balanceAmount.value = '' selectedGroupIds.value = [] lastResult.value = null } @@ -235,6 +304,12 @@ async function executeBatchAction(): Promise { return } request = { selection, action: 'update_access_control', payload } + } else if (selectedAction.value === 'adjust_wallet_balance') { + if (balancePayload.value === null) { + warning(legacyT('请输入大于 0 的有限金额')) + return + } + request = { selection, action: 'adjust_wallet_balance', payload: balancePayload.value } } else if (selectedAction.value === 'update_role') { request = { selection, action: 'update_role', payload: buildRolePayload() } } else { @@ -253,7 +328,7 @@ async function executeBatchAction(): Promise { } emit('completed', result) } catch (err) { - error(parseApiError(err, '批量操作失败'), legacyT('批量操作失败')) + error(legacyT(parseApiError(err, '批量操作失败')), legacyT('批量操作失败')) } finally { executing.value = false } diff --git a/frontend/src/features/users/components/user-management-config.ts b/frontend/src/features/users/components/user-management-config.ts index f08293046..a52c742d0 100644 --- a/frontend/src/features/users/components/user-management-config.ts +++ b/frontend/src/features/users/components/user-management-config.ts @@ -3,6 +3,7 @@ import { CheckCircle2, ShieldCheck, UserCog, + Wallet, } from 'lucide-vue-next' import type { Component } from 'vue' import type { UserBatchAction, UserRole } from '@/api/users' @@ -64,6 +65,12 @@ export const USER_BATCH_ACTION_OPTIONS: UserBatchActionOption[] = [ description: '批量设为普通用户或管理员', icon: UserCog, }, + { + value: 'adjust_wallet_balance', + label: '调整余额', + description: '批量增加或扣减钱包余额', + icon: Wallet, + }, ] export function formatUserRoleLabel(role: UserRole | string): string { diff --git a/frontend/src/i18n/messages.ts b/frontend/src/i18n/messages.ts index 9aa571175..fdb2f26f8 100644 --- a/frontend/src/i18n/messages.ts +++ b/frontend/src/i18n/messages.ts @@ -1453,6 +1453,17 @@ const legacyExactEnglishMessages: Record = { '暂无选项': 'No options', '用户批量操作': 'User batch actions', '按当前选择批量调整用户状态、角色和额度': 'Batch update user status, role, and quota for the current selection', + '按当前选择批量调整用户状态、角色、额度和钱包余额': 'Batch update user status, role, quota, and wallet balances for the current selection', + '调整余额': 'Adjust balance', + '批量增加或扣减钱包余额': 'Add or deduct wallet balances in bulk', + '调整金额 (USD)': 'Adjustment amount (USD)', + '余额调整方式': 'Balance adjustment mode', + '增加': 'Add', + '扣减': 'Deduct', + '请输入大于 0 的有限金额': 'Enter a finite amount greater than 0', + '扣减超过单个用户可用余额时,该用户余额将归零。': 'Deductions above an individual user’s available balance will be clamped to zero.', + '用户钱包不可用': 'User wallet is unavailable', + '当前为只读模式,无法批量调整用户钱包余额': 'Cannot adjust user wallet balances in read-only mode', '影响用户:': 'Affected users:', '目标为当前筛选条件匹配的全部用户,执行前后端会重新解析。': 'Targets all users matching the current filters; the backend will resolve the selection again before execution.', '目标为当前已勾选的用户,重复 ID 会自动去重。': 'Targets the currently selected users; duplicate IDs are deduplicated automatically.', From 30bb0c31301c84d14512c894c19b016faa625359 Mon Sep 17 00:00:00 2001 From: RWDai <27391645+RWDai@users.noreply.github.com> Date: Wed, 23 Sep 2026 16:29:47 +0800 Subject: [PATCH 2/9] fix: harden bulk wallet balance adjustment --- .../src/handlers/admin/users/batch.rs | 168 +++++++++++++++--- .../src/handlers/admin/users/mod.rs | 9 +- .../src/handlers/admin/users/shared.rs | 12 ++ apps/aether-gateway/src/state/app.rs | 2 + apps/aether-gateway/src/state/core.rs | 2 + .../state/runtime/wallet/balance_mutations.rs | 7 + apps/aether-gateway/src/state/testing.rs | 8 + .../src/tests/control/admin/users_batch.rs | 157 +++++++++++++++- frontend/src/api/users.ts | 4 + .../components/UserBatchActionDialog.vue | 10 ++ .../components/UserBatchResultSummary.vue | 15 ++ frontend/src/i18n/messages.ts | 11 ++ 12 files changed, 370 insertions(+), 35 deletions(-) diff --git a/apps/aether-gateway/src/handlers/admin/users/batch.rs b/apps/aether-gateway/src/handlers/admin/users/batch.rs index c6635ef6c..7dd7f5485 100644 --- a/apps/aether-gateway/src/handlers/admin/users/batch.rs +++ b/apps/aether-gateway/src/handlers/admin/users/batch.rs @@ -1,6 +1,7 @@ use super::{ build_admin_users_bad_request_response, build_admin_users_permission_denied_response, build_admin_users_read_only_response, disabled_user_policy_detail, disabled_user_policy_field, + management_token_may_adjust_admin_wallet_balance, management_token_may_administer_user_accounts, normalize_admin_user_role, }; use crate::handlers::admin::request::{AdminAppState, AdminRequestContext}; @@ -173,6 +174,13 @@ pub(in super::super) async fn build_admin_user_batch_action_response( "当前为只读模式,无法批量更新用户钱包", )); } + if mutation.wallet_balance_adjustment.is_some() + && !management_token_may_adjust_admin_wallet_balance(request_context) + { + return Ok(build_admin_users_permission_denied_response( + request_context, + )); + } if mutation.wallet_balance_adjustment.is_some() && !state.has_auth_wallet_write_capability() { return Ok(build_admin_users_read_only_response( "当前为只读模式,无法批量调整用户钱包余额", @@ -195,9 +203,29 @@ pub(in super::super) async fn build_admin_user_batch_action_response( .iter() .map(|user_id| json!({ "user_id": user_id, "reason": "用户不存在或已删除" })) .collect::>(); + let mut completed_user_ids = Vec::new(); + let mut uncertain_user_ids = Vec::new(); + let mut unprocessed_user_ids = Vec::new(); + let mut interrupted = false; - for item in &resolved.items { - if state.find_user_auth_by_id(&item.user_id).await?.is_none() { + for (item_index, item) in resolved.items.iter().enumerate() { + let user = match state.find_user_auth_by_id(&item.user_id).await { + Ok(user) => user, + Err(_) => { + record_batch_action_interruption( + &resolved.items, + item_index, + false, + "读取用户状态失败,批次已中止,该用户未执行", + &mut failures, + &mut uncertain_user_ids, + &mut unprocessed_user_ids, + ); + interrupted = true; + break; + } + }; + if user.is_none() { failures.push(json!({ "user_id": item.user_id, "reason": "用户不存在或已删除", @@ -220,34 +248,66 @@ pub(in super::super) async fn build_admin_user_batch_action_response( } if let Some(unlimited) = mutation.unlimited { - if !apply_batch_user_wallet_limit_mode(state, &item.user_id, unlimited).await? { - failures.push(json!({ - "user_id": item.user_id, - "reason": "用户钱包不可用", - })); - continue; + match apply_batch_user_wallet_limit_mode(state, &item.user_id, unlimited).await { + Ok(true) => {} + Ok(false) => { + failures.push(json!({ + "user_id": item.user_id, + "reason": "用户钱包不可用", + })); + continue; + } + Err(_) => { + record_batch_action_interruption( + &resolved.items, + item_index, + true, + "用户钱包更新结果未确认,批次已中止,请核对钱包后再重试", + &mut failures, + &mut uncertain_user_ids, + &mut unprocessed_user_ids, + ); + interrupted = true; + break; + } } } if let Some(adjustment) = mutation.wallet_balance_adjustment { - if !apply_batch_user_wallet_balance_adjustment( + match apply_batch_user_wallet_balance_adjustment( state, &item.user_id, adjustment, current_admin_user_id, ) - .await? + .await { - failures.push(json!({ - "user_id": item.user_id, - "reason": "用户钱包不可用", - })); - continue; + Ok(true) => {} + Ok(false) => { + failures.push(json!({ + "user_id": item.user_id, + "reason": "用户钱包不可用", + })); + continue; + } + Err(_) => { + record_batch_action_interruption( + &resolved.items, + item_index, + true, + "余额调整结果未确认,批次已中止,请核对钱包后再重试", + &mut failures, + &mut uncertain_user_ids, + &mut unprocessed_user_ids, + ); + interrupted = true; + break; + } } } - if mutation.has_auth_user_fields() - && state + if mutation.has_auth_user_fields() { + let updated_user = match state .update_local_auth_user_admin_fields( &item.user_id, mutation.role.clone(), @@ -261,22 +321,39 @@ pub(in super::super) async fn build_admin_user_batch_action_response( None, mutation.is_active, ) - .await? - .is_none() - { - failures.push(json!({ - "user_id": item.user_id, - "reason": "用户不存在或已删除", - })); - continue; + .await + { + Ok(user) => user, + Err(_) => { + record_batch_action_interruption( + &resolved.items, + item_index, + true, + "用户更新结果未确认,批次已中止,请核对后再重试", + &mut failures, + &mut uncertain_user_ids, + &mut unprocessed_user_ids, + ); + interrupted = true; + break; + } + }; + if updated_user.is_none() { + failures.push(json!({ + "user_id": item.user_id, + "reason": "用户不存在或已删除", + })); + continue; + } } success += 1; + completed_user_ids.push(item.user_id.clone()); } let failed = failures.len(); let total = success + failed; - let response = Json(json!({ + let mut response_payload = json!({ "total": total, "success": success, "failed": failed, @@ -284,8 +361,14 @@ pub(in super::super) async fn build_admin_user_batch_action_response( "warnings": resolved.warnings, "action": request.action.trim().to_ascii_lowercase(), "modified_fields": mutation.modified_fields, - })) - .into_response(); + "interrupted": interrupted, + }); + if interrupted { + response_payload["completed_user_ids"] = json!(completed_user_ids); + response_payload["uncertain_user_ids"] = json!(uncertain_user_ids); + response_payload["unprocessed_user_ids"] = json!(unprocessed_user_ids); + } + let response = Json(response_payload).into_response(); Ok(attach_admin_audit_response( response, @@ -296,6 +379,35 @@ pub(in super::super) async fn build_admin_user_batch_action_response( )) } +fn record_batch_action_interruption( + items: &[AdminUserSelectionItem], + item_index: usize, + current_result_uncertain: bool, + reason: &str, + failures: &mut Vec, + uncertain_user_ids: &mut Vec, + unprocessed_user_ids: &mut Vec, +) { + let current_item = &items[item_index]; + failures.push(json!({ + "user_id": current_item.user_id, + "reason": reason, + })); + if current_result_uncertain { + uncertain_user_ids.push(current_item.user_id.clone()); + } else { + unprocessed_user_ids.push(current_item.user_id.clone()); + } + + for item in items.iter().skip(item_index + 1) { + failures.push(json!({ + "user_id": item.user_id, + "reason": "因前序错误未执行", + })); + unprocessed_user_ids.push(item.user_id.clone()); + } +} + fn parse_resolve_selection_request( request_body: Option<&Bytes>, ) -> Result { diff --git a/apps/aether-gateway/src/handlers/admin/users/mod.rs b/apps/aether-gateway/src/handlers/admin/users/mod.rs index e5c854748..1ee823993 100644 --- a/apps/aether-gateway/src/handlers/admin/users/mod.rs +++ b/apps/aether-gateway/src/handlers/admin/users/mod.rs @@ -51,10 +51,11 @@ use self::shared::{ build_admin_users_data_unavailable_response, build_admin_users_permission_denied_response, build_admin_users_read_only_response, disabled_user_policy_detail, disabled_user_policy_field, format_optional_datetime_iso8601, legacy_admin_list_policy_mode, - legacy_admin_rate_limit_policy_mode, management_token_may_administer_user_accounts, - normalize_admin_optional_user_email, normalize_admin_user_group_ids, normalize_admin_user_role, - normalize_admin_username, validate_admin_user_password, AdminCreateUserApiKeyRequest, - AdminCreateUserRequest, AdminToggleUserApiKeyLockRequest, AdminUpdateUserApiKeyRequest, + legacy_admin_rate_limit_policy_mode, management_token_may_adjust_admin_wallet_balance, + management_token_may_administer_user_accounts, normalize_admin_optional_user_email, + normalize_admin_user_group_ids, normalize_admin_user_role, normalize_admin_username, + validate_admin_user_password, AdminCreateUserApiKeyRequest, AdminCreateUserRequest, + AdminToggleUserApiKeyLockRequest, AdminUpdateUserApiKeyRequest, }; pub(crate) use self::shared::{ normalize_admin_list_policy_mode, normalize_admin_rate_limit_policy_mode, diff --git a/apps/aether-gateway/src/handlers/admin/users/shared.rs b/apps/aether-gateway/src/handlers/admin/users/shared.rs index b304932c9..df322a49f 100644 --- a/apps/aether-gateway/src/handlers/admin/users/shared.rs +++ b/apps/aether-gateway/src/handlers/admin/users/shared.rs @@ -165,6 +165,18 @@ pub(super) fn management_token_may_administer_user_accounts( }) } +pub(super) fn management_token_may_adjust_admin_wallet_balance( + request_context: &crate::handlers::admin::request::AdminRequestContext<'_>, +) -> bool { + request_context.decision().is_some_and(|decision| { + crate::control::management_token_principal_has_permission(decision, "admin:wallets:write") + || crate::control::management_token_principal_has_permission( + decision, + "admin:wallets:admin", + ) + }) +} + pub(super) fn build_admin_users_permission_denied_response( request_context: &crate::handlers::admin::request::AdminRequestContext<'_>, ) -> Response { diff --git a/apps/aether-gateway/src/state/app.rs b/apps/aether-gateway/src/state/app.rs index 1851b7c0c..3f2c7657d 100644 --- a/apps/aether-gateway/src/state/app.rs +++ b/apps/aether-gateway/src/state/app.rs @@ -492,6 +492,8 @@ pub struct AppState { Arc>>, >, #[cfg(test)] + pub(crate) auth_wallet_adjustment_error_for_tests: Option, + #[cfg(test)] pub(crate) admin_wallet_payment_order_store: Option>>>, #[cfg(test)] diff --git a/apps/aether-gateway/src/state/core.rs b/apps/aether-gateway/src/state/core.rs index ab462ed90..7e3e060ae 100644 --- a/apps/aether-gateway/src/state/core.rs +++ b/apps/aether-gateway/src/state/core.rs @@ -466,6 +466,8 @@ impl AppState { #[cfg(test)] auth_wallet_store: Some(Arc::new(StdMutex::new(HashMap::new()))), #[cfg(test)] + auth_wallet_adjustment_error_for_tests: None, + #[cfg(test)] admin_wallet_payment_order_store: Some(Arc::new(StdMutex::new(HashMap::new()))), #[cfg(test)] admin_payment_callback_store: Some(Arc::new(StdMutex::new(HashMap::new()))), diff --git a/apps/aether-gateway/src/state/runtime/wallet/balance_mutations.rs b/apps/aether-gateway/src/state/runtime/wallet/balance_mutations.rs index f1e170669..212233803 100644 --- a/apps/aether-gateway/src/state/runtime/wallet/balance_mutations.rs +++ b/apps/aether-gateway/src/state/runtime/wallet/balance_mutations.rs @@ -18,6 +18,13 @@ impl AppState { )>, GatewayError, > { + #[cfg(test)] + if self.auth_wallet_adjustment_error_for_tests.as_deref() == Some(wallet_id) { + return Err(GatewayError::Internal( + "injected test wallet adjustment failure".to_string(), + )); + } + #[cfg(test)] if let Some(store) = self.auth_wallet_store.as_ref() { let mut guard = store.lock().expect("auth wallet store should lock"); diff --git a/apps/aether-gateway/src/state/testing.rs b/apps/aether-gateway/src/state/testing.rs index bdbf662e9..cfa990a32 100644 --- a/apps/aether-gateway/src/state/testing.rs +++ b/apps/aether-gateway/src/state/testing.rs @@ -472,6 +472,14 @@ impl AppState { self } + pub(crate) fn fail_auth_wallet_adjustment_for_tests( + mut self, + wallet_id: impl Into, + ) -> Self { + self.auth_wallet_adjustment_error_for_tests = Some(wallet_id.into()); + self + } + pub(crate) fn with_admin_wallet_payment_orders_for_tests(mut self, orders: I) -> Self where I: IntoIterator, diff --git a/apps/aether-gateway/src/tests/control/admin/users_batch.rs b/apps/aether-gateway/src/tests/control/admin/users_batch.rs index 7abb4d7bc..2a1cb6466 100644 --- a/apps/aether-gateway/src/tests/control/admin/users_batch.rs +++ b/apps/aether-gateway/src/tests/control/admin/users_batch.rs @@ -1,11 +1,17 @@ -use aether_data::repository::users::StoredUserAuthRecord; +use std::sync::Arc; + +use aether_data::repository::management_tokens::InMemoryManagementTokenRepository; +use aether_data::repository::users::{InMemoryUserReadRepository, StoredUserAuthRecord}; use aether_data::repository::wallet::StoredWalletSnapshot; use axum::http::StatusCode; use chrono::Utc; use reqwest::{Client, RequestBuilder, Response}; use serde_json::{json, Value}; -use super::super::{build_router_with_state, start_server, AppState}; +use super::super::{ + build_router_with_state, hash_management_token, sample_management_token, start_server, AppState, +}; +use crate::data::GatewayDataState; fn admin_headers(request: RequestBuilder) -> RequestBuilder { request @@ -19,13 +25,17 @@ fn admin_headers(request: RequestBuilder) -> RequestBuilder { } fn sample_user(user_id: &str) -> StoredUserAuthRecord { + sample_user_with_role(user_id, "user") +} + +fn sample_user_with_role(user_id: &str, role: &str) -> StoredUserAuthRecord { StoredUserAuthRecord::new( user_id.to_string(), Some(format!("{user_id}@example.com")), true, user_id.to_string(), Some("hash".to_string()), - "user".to_string(), + role.to_string(), "local".to_string(), Some(json!(["openai"])), Some(json!(["openai:chat"])), @@ -162,6 +172,147 @@ async fn gateway_batches_wallet_addition_deduction_and_clamped_deduction_per_use gateway_handle.abort(); } +#[tokio::test] +async fn gateway_requires_wallet_write_permission_for_batch_balance_adjustments() { + let users_write_token = "ae-batch-users-write-only"; + let wallet_write_token = "ae-batch-users-wallet-write"; + let token_owner = sample_user_with_role("token-owner", "admin"); + let target_user = sample_user("user-1"); + let mut users_only = sample_management_token( + "token-users-write-only", + &token_owner.id, + &token_owner.username, + true, + ); + users_only.token.allowed_ips = None; + users_only.token.permissions = Some(json!(["admin:users:write"])); + let mut users_and_wallets = sample_management_token( + "token-users-and-wallet-write", + &token_owner.id, + &token_owner.username, + true, + ); + users_and_wallets.token.allowed_ips = None; + users_and_wallets.token.permissions = Some(json!(["admin:users:write", "admin:wallets:write"])); + let token_repository = Arc::new(InMemoryManagementTokenRepository::seed_with_hashes( + vec![users_only, users_and_wallets], + vec![ + ( + hash_management_token(users_write_token), + "token-users-write-only".to_string(), + ), + ( + hash_management_token(wallet_write_token), + "token-users-and-wallet-write".to_string(), + ), + ], + )); + let user_repository = Arc::new(InMemoryUserReadRepository::seed_auth_users(vec![ + token_owner.clone(), + target_user.clone(), + ])); + let data = GatewayDataState::with_management_token_repository_for_tests(token_repository) + .with_user_reader(user_repository); + let state = AppState::new() + .expect("gateway should build") + .with_data_state_for_tests(data) + .with_auth_users_for_tests([token_owner, target_user]) + .with_auth_wallets_for_tests([sample_wallet("user-1", 10.0, 0.0)]); + let (gateway_url, gateway_handle) = start_server(build_router_with_state(state)).await; + let client = Client::new(); + let payload = json!({ + "selection": { "user_ids": ["user-1"] }, + "action": "adjust_wallet_balance", + "payload": { "operation": "add", "amount": 5.0 } + }); + + let denied = client + .post(format!("{gateway_url}/api/admin/users/batch-action")) + .header(crate::constants::GATEWAY_HEADER, "rust-phase3b") + .bearer_auth(users_write_token) + .json(&payload) + .send() + .await + .expect("users-only management token request should complete"); + assert_eq!(denied.status(), StatusCode::FORBIDDEN); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 10.0 + ); + + let allowed = client + .post(format!("{gateway_url}/api/admin/users/batch-action")) + .header(crate::constants::GATEWAY_HEADER, "rust-phase3b") + .bearer_auth(wallet_write_token) + .json(&payload) + .send() + .await + .expect("wallet-write management token request should complete"); + assert_eq!(allowed.status(), StatusCode::OK); + let result: Value = allowed.json().await.expect("response should parse"); + assert_eq!(result["success"], 1); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 15.0 + ); + + gateway_handle.abort(); +} + +#[tokio::test] +async fn gateway_reports_completed_uncertain_and_unprocessed_users_after_adjustment_error() { + let state = AppState::new() + .expect("gateway should build") + .with_auth_users_for_tests([ + sample_user("user-1"), + sample_user("user-2"), + sample_user("user-3"), + ]) + .with_auth_wallets_for_tests([ + sample_wallet("user-1", 10.0, 0.0), + sample_wallet("user-2", 20.0, 0.0), + sample_wallet("user-3", 30.0, 0.0), + ]) + .fail_auth_wallet_adjustment_for_tests("wallet-user-2"); + let (gateway_url, gateway_handle) = start_server(build_router_with_state(state)).await; + let client = Client::new(); + + let response = post_batch_action( + &client, + &gateway_url, + json!({ + "selection": { "user_ids": ["user-1", "user-2", "user-3"] }, + "action": "adjust_wallet_balance", + "payload": { "operation": "add", "amount": 5.0 } + }), + ) + .await; + assert_eq!(response.status(), StatusCode::OK); + let result: Value = response.json().await.expect("response should parse"); + assert_eq!(result["interrupted"], true); + assert_eq!(result["success"], 1); + assert_eq!(result["failed"], 2); + assert_eq!(result["completed_user_ids"], json!(["user-1"])); + assert_eq!(result["uncertain_user_ids"], json!(["user-2"])); + assert_eq!(result["unprocessed_user_ids"], json!(["user-3"])); + assert_eq!(result["failures"][0]["user_id"], "user-2"); + assert_eq!(result["failures"][1]["user_id"], "user-3"); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 15.0 + ); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-2").await["balance"], + 20.0 + ); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-3").await["balance"], + 30.0 + ); + + gateway_handle.abort(); +} + #[tokio::test] async fn gateway_reports_missing_wallet_and_skips_zero_delta_for_non_positive_balance() { let state = AppState::new() diff --git a/frontend/src/api/users.ts b/frontend/src/api/users.ts index ee2121c16..1f216e31f 100644 --- a/frontend/src/api/users.ts +++ b/frontend/src/api/users.ts @@ -183,6 +183,10 @@ export interface UserBatchActionResponse { success: number failed: number failures: UserBatchActionFailure[] + interrupted?: boolean + completed_user_ids?: string[] + uncertain_user_ids?: string[] + unprocessed_user_ids?: string[] warnings?: UserBatchSelectionWarning[] action?: string modified_fields?: string[] diff --git a/frontend/src/features/users/components/UserBatchActionDialog.vue b/frontend/src/features/users/components/UserBatchActionDialog.vue index 282e51e9b..4712016a3 100644 --- a/frontend/src/features/users/components/UserBatchActionDialog.vue +++ b/frontend/src/features/users/components/UserBatchActionDialog.vue @@ -210,6 +210,11 @@ const targetRoleWarning = computed(() => { const executeButtonLabel = computed(() => legacyT(`确认${selectedActionLabel.value}(${impactCount.value})`)) const lastResultLabel = computed(() => { if (!lastResult.value) return '' + if (lastResult.value.interrupted) { + return legacyT( + `批量操作中断:成功 ${lastResult.value.success} 个,结果待确认 ${lastResult.value.uncertain_user_ids?.length ?? 0} 个,尚未执行 ${lastResult.value.unprocessed_user_ids?.length ?? 0} 个`, + ) + } return legacyT(`成功 ${lastResult.value.success} 个,失败 ${lastResult.value.failed} 个`) }) const lastResultFailuresLabel = computed(() => { @@ -320,6 +325,11 @@ async function executeBatchAction(): Promise { try { const result = await usersStore.batchAction(request) lastResult.value = result + if (result.interrupted) { + warning(`${lastResultLabel.value};${legacyT('请核对余额后再重试,勿直接重试整批')}`) + emit('completed', result) + return + } const message = legacyT(`批量操作完成:成功 ${result.success} 个,失败 ${result.failed} 个`) if (result.failed > 0) { warning(message) diff --git a/frontend/src/features/users/components/UserBatchResultSummary.vue b/frontend/src/features/users/components/UserBatchResultSummary.vue index 30ec4c956..a3a68f126 100644 --- a/frontend/src/features/users/components/UserBatchResultSummary.vue +++ b/frontend/src/features/users/components/UserBatchResultSummary.vue @@ -7,11 +7,26 @@ {{ failuresLabel }} + + {{ legacyT('查看批次中断详情') }} + + {{ legacyT('已完成用户 ID') }}:{{ userIds(result.completed_user_ids) }} + {{ legacyT('结果待核对用户 ID') }}:{{ userIds(result.uncertain_user_ids) }} + {{ legacyT('尚未执行用户 ID') }}:{{ userIds(result.unprocessed_user_ids) }} + + diff --git a/frontend/src/features/users/utils/__tests__/userBatchWalletIdempotency.spec.ts b/frontend/src/features/users/utils/__tests__/userBatchWalletIdempotency.spec.ts new file mode 100644 index 000000000..2452fe98e --- /dev/null +++ b/frontend/src/features/users/utils/__tests__/userBatchWalletIdempotency.spec.ts @@ -0,0 +1,312 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest' +import type { UserBatchActionResponse, UserBatchBalanceActionRequest } from '@/api/users' +import { + createUserBatchWalletRetryCoordinator, + UnresolvedWalletRequestMismatchError, + WalletIdempotencyPersistenceUnavailableError, + WalletIdempotencyUnavailableError, +} from '../userBatchWalletIdempotency' + +function createStorage() { + const values = new Map() + return { + getItem: (key: string) => values.get(key) ?? null, + setItem: (key: string, value: string) => values.set(key, value), + removeItem: (key: string) => values.delete(key), + } +} + +const walletRequest = { + selection: { user_ids: ['user-1', 'user-2'], group_ids: ['group-1'] }, + action: 'adjust_wallet_balance' as const, + payload: { operation: 'deduct' as const, amount: 17.25 }, +} +const defaultStorageKey = 'admin.users.batch.wallet-adjustment.pending.v1:default' +let testFallback: Map + +function response(interrupted = false): UserBatchActionResponse { + return { total: 2, success: 1, failed: 0, failures: [], interrupted } +} + +describe('user batch wallet idempotency', () => { + beforeEach(() => { + sessionStorage.clear() + localStorage.clear() + testFallback = new Map() + }) + + it('serializes the key with the exact top-level wallet request before sending', async () => { + const storage = createStorage() + const coordinator = createUserBatchWalletRetryCoordinator({ + storage, + fallback: testFallback, + createKey: () => 'wallet-key-1', + }) + let storedDuringSend: string | null = null + let sentRequest: UserBatchBalanceActionRequest | undefined + + await coordinator.execute(walletRequest, async (request) => { + sentRequest = request + storedDuringSend = storage.getItem(defaultStorageKey) + return response(true) + }) + + expect(sentRequest).toEqual({ ...walletRequest, idempotency_key: 'wallet-key-1' }) + expect(JSON.parse(storedDuringSend ?? 'null')).toEqual({ + idempotency_key: 'wallet-key-1', + request: sentRequest, + }) + }) + + it('retains a transport failure and reopens with the exact request for retry', async () => { + const storage = createStorage() + const first = createUserBatchWalletRetryCoordinator({ + storage, + fallback: testFallback, + createKey: () => 'wallet-key-2', + }) + const sendFailure = new Error('connection lost') + await expect(first.execute(walletRequest, async () => { throw sendFailure })).rejects.toBe(sendFailure) + + const reopened = createUserBatchWalletRetryCoordinator({ + storage, + fallback: testFallback, + createKey: () => 'must-not-be-used', + }) + const pending = reopened.getPending() + expect(pending).toEqual({ + idempotency_key: 'wallet-key-2', + request: { ...walletRequest, idempotency_key: 'wallet-key-2' }, + }) + + const send = vi.fn(async () => response()) + await expect(reopened.retry(send)).resolves.toEqual(response()) + expect(send).toHaveBeenCalledWith(pending?.request) + expect(storage.getItem(defaultStorageKey)).toBeNull() + }) + + it('keeps unresolved requests in persistent browser storage across coordinators', async () => { + const first = createUserBatchWalletRetryCoordinator({ + createKey: () => 'wallet-key-persistent', + scope: () => 'admin-1', + }) + await expect(first.execute(walletRequest, async () => { + throw new Error('connection lost') + })).rejects.toThrow('connection lost') + + const reopened = createUserBatchWalletRetryCoordinator({ scope: () => 'admin-1' }) + const pending = reopened.getPending() + expect(pending?.request).toEqual({ + ...walletRequest, + idempotency_key: 'wallet-key-persistent', + }) + + const send = vi.fn(async () => response()) + await reopened.retry(send) + expect(send).toHaveBeenCalledWith(pending?.request) + expect(localStorage.getItem( + 'admin.users.batch.wallet-adjustment.pending.v1:admin-1', + )).toBeNull() + }) + + it('keeps unresolved requests isolated by authenticated administrator', async () => { + const storage = createStorage() + const adminA = createUserBatchWalletRetryCoordinator({ + storage, + fallback: testFallback, + scope: () => 'admin-a', + createKey: () => 'wallet-key-admin-a', + }) + await expect(adminA.execute(walletRequest, async () => { throw new Error('connection lost') })) + .rejects.toThrow('connection lost') + + const adminB = createUserBatchWalletRetryCoordinator({ + storage, + fallback: testFallback, + scope: () => 'admin-b', + createKey: () => 'wallet-key-admin-b', + }) + expect(adminB.getPending()).toBeNull() + await adminB.execute(walletRequest, async () => response(true)) + + expect(adminA.getPending()?.idempotency_key).toBe('wallet-key-admin-a') + expect(adminB.getPending()?.idempotency_key).toBe('wallet-key-admin-b') + }) + + it('does not send if the authenticated administrator changes before dispatch', async () => { + const storage = createStorage() + let scopeReads = 0 + const send = vi.fn(async () => response()) + const coordinator = createUserBatchWalletRetryCoordinator({ + storage, + fallback: testFallback, + scope: () => (++scopeReads === 1 ? 'admin-a' : 'admin-b'), + createKey: () => 'wallet-key-scope-change', + }) + + await expect(coordinator.execute(walletRequest, send)).rejects.toThrow( + 'authenticated administrator changed', + ) + expect(send).not.toHaveBeenCalled() + expect(storage.getItem('admin.users.batch.wallet-adjustment.pending.v1:admin-a')).not.toBeNull() + }) + + it('retains interrupted requests and reuses their key until a terminal response', async () => { + const storage = createStorage() + const first = createUserBatchWalletRetryCoordinator({ + storage, + fallback: testFallback, + createKey: () => 'wallet-key-3', + }) + await first.execute(walletRequest, async () => response(true)) + const reopened = createUserBatchWalletRetryCoordinator({ storage, fallback: testFallback }) + const pending = reopened.getPending() + const send = vi.fn(async () => response(true)) + + await reopened.retry(send) + + expect(send).toHaveBeenCalledWith(pending?.request) + expect(reopened.getPending()).toEqual(pending) + await reopened.retry(async (request) => { + expect(request).toEqual(pending?.request) + return response() + }) + expect(reopened.getPending()).toBeNull() + }) + + it('matches the same serialized request when optional filter fields are omitted', async () => { + const storage = createStorage() + const requestWithUndefinedField = { + ...walletRequest, + selection: { filters: { search: 'active', is_active: undefined } }, + } + const first = createUserBatchWalletRetryCoordinator({ + storage, + fallback: testFallback, + createKey: () => 'wallet-key-filter', + }) + await first.execute(requestWithUndefinedField, async () => response(true)) + + const reopened = createUserBatchWalletRetryCoordinator({ storage, fallback: testFallback }) + const send = vi.fn(async () => response()) + await reopened.execute({ + ...walletRequest, + selection: { filters: { search: 'active' } }, + }, send) + + expect(send).toHaveBeenCalledWith({ + ...walletRequest, + selection: { filters: { search: 'active' } }, + idempotency_key: 'wallet-key-filter', + }) + }) + + it('blocks changed payloads while unresolved and gives a later adjustment a new key', async () => { + const storage = createStorage() + let nextKey = 0 + const coordinator = createUserBatchWalletRetryCoordinator({ + storage, + fallback: testFallback, + createKey: () => `wallet-key-${++nextKey}`, + }) + await expect(coordinator.execute(walletRequest, async () => { throw new Error('connection lost') })) + .rejects.toThrow('connection lost') + const changedRequest = { + ...walletRequest, + payload: { operation: 'add' as const, amount: 20 }, + } + const send = vi.fn(async () => response()) + + await expect(coordinator.execute(changedRequest, send)).rejects.toBeInstanceOf( + UnresolvedWalletRequestMismatchError, + ) + expect(send).not.toHaveBeenCalled() + expect(coordinator.getPending()?.idempotency_key).toBe('wallet-key-1') + + await coordinator.retry(async () => response()) + let newRequest + await coordinator.execute(changedRequest, async (request) => { + newRequest = request + return response() + }) + + expect(newRequest).toEqual({ ...changedRequest, idempotency_key: 'wallet-key-2' }) + }) + + it('fails closed when persistent storage cannot save the request', async () => { + const unavailableStorage = { + getItem: () => null, + setItem: () => { throw new Error('storage unavailable') }, + removeItem: () => undefined, + } + const send = vi.fn(async () => response(true)) + const coordinator = createUserBatchWalletRetryCoordinator({ + storage: unavailableStorage, + fallback: testFallback, + createKey: () => 'wallet-key-fallback', + }) + + await expect(coordinator.execute(walletRequest, send)).rejects.toBeInstanceOf( + WalletIdempotencyPersistenceUnavailableError, + ) + expect(send).not.toHaveBeenCalled() + expect(testFallback.size).toBe(0) + }) + + it('does not send when persistent storage readback does not match', async () => { + let wasWritten = false + const mismatchedStorage = { + getItem: () => wasWritten ? 'different request' : null, + setItem: () => { wasWritten = true }, + removeItem: () => undefined, + } + const send = vi.fn(async () => response()) + const coordinator = createUserBatchWalletRetryCoordinator({ + storage: mismatchedStorage, + fallback: testFallback, + createKey: () => 'wallet-key-readback', + }) + + await expect(coordinator.execute(walletRequest, send)).rejects.toBeInstanceOf( + WalletIdempotencyPersistenceUnavailableError, + ) + expect(send).not.toHaveBeenCalled() + expect(testFallback.size).toBe(0) + }) + + it('does not create a new request when persistent storage cannot be read', async () => { + const unavailableStorage = { + getItem: () => { throw new Error('storage unavailable') }, + setItem: () => undefined, + removeItem: () => undefined, + } + const send = vi.fn(async () => response()) + const coordinator = createUserBatchWalletRetryCoordinator({ + storage: unavailableStorage, + fallback: testFallback, + createKey: () => 'must-not-be-used', + }) + + await expect(coordinator.execute(walletRequest, send)).rejects.toBeInstanceOf( + WalletIdempotencyPersistenceUnavailableError, + ) + expect(send).not.toHaveBeenCalled() + expect(testFallback.size).toBe(0) + }) + + it('fails closed when secure UUID generation is unavailable', async () => { + const storage = createStorage() + const send = vi.fn(async () => response()) + const coordinator = createUserBatchWalletRetryCoordinator({ + storage, + fallback: testFallback, + createKey: () => { throw new WalletIdempotencyUnavailableError() }, + }) + + await expect(coordinator.execute(walletRequest, send)).rejects.toBeInstanceOf( + WalletIdempotencyUnavailableError, + ) + expect(send).not.toHaveBeenCalled() + expect(storage.getItem(defaultStorageKey)).toBeNull() + }) +}) diff --git a/frontend/src/features/users/utils/userBatchWalletIdempotency.ts b/frontend/src/features/users/utils/userBatchWalletIdempotency.ts new file mode 100644 index 000000000..d81fc412c --- /dev/null +++ b/frontend/src/features/users/utils/userBatchWalletIdempotency.ts @@ -0,0 +1,252 @@ +import type { + UserBatchActionResponse, + UserBatchBalanceActionRequest, +} from '@/api/users' + +export type UserBatchWalletAdjustmentRequest = Omit< + UserBatchBalanceActionRequest, + 'idempotency_key' +> + +export interface PendingUserBatchWalletRequest { + idempotency_key: string + request: UserBatchBalanceActionRequest +} + +interface StringStorage { + getItem(key: string): string | null + setItem(key: string, value: string): void + removeItem(key: string): void +} + +interface CoordinatorOptions { + storage?: StringStorage | null + createKey?: () => string + fallback?: Map + scope?: () => string | null +} + +const STORAGE_KEY = 'admin.users.batch.wallet-adjustment.pending.v1' +const inMemoryFallback = new Map() + +export class UnresolvedWalletRequestMismatchError extends Error { + constructor(readonly pending: PendingUserBatchWalletRequest) { + super('A different wallet batch request is still unresolved') + this.name = 'UnresolvedWalletRequestMismatchError' + } +} + +export class WalletIdempotencyUnavailableError extends Error { + constructor() { + super('crypto.randomUUID is unavailable') + this.name = 'WalletIdempotencyUnavailableError' + } +} + +export class WalletIdempotencyPersistenceUnavailableError extends Error { + constructor() { + super('Persistent browser storage is unavailable') + this.name = 'WalletIdempotencyPersistenceUnavailableError' + } +} + +export class WalletIdempotencyScopeUnavailableError extends Error { + constructor() { + super('The authenticated administrator identity is unavailable') + this.name = 'WalletIdempotencyScopeUnavailableError' + } +} + +export class WalletIdempotencyScopeChangedError extends Error { + constructor() { + super('The authenticated administrator changed before the request was sent') + this.name = 'WalletIdempotencyScopeChangedError' + } +} + +export class InvalidPendingWalletRequestError extends Error { + constructor() { + super('The stored wallet batch request is invalid') + this.name = 'InvalidPendingWalletRequestError' + } +} + +function browserPersistentStorage(): StringStorage | null { + try { + return globalThis.localStorage ?? null + } catch { + return null + } +} + +function secureRandomUUID(): string { + try { + const cryptoApi = globalThis.crypto + if (typeof cryptoApi?.randomUUID === 'function') { + return cryptoApi.randomUUID() + } + } catch { + // Treat unavailable secure randomness as a hard failure. + } + throw new WalletIdempotencyUnavailableError() +} + +function isPendingRequest(value: unknown): value is PendingUserBatchWalletRequest { + if (typeof value !== 'object' || value === null) return false + const record = value as Partial + const request = record.request + return typeof record.idempotency_key === 'string' + && record.idempotency_key.length > 0 + && typeof request === 'object' + && request !== null + && request.action === 'adjust_wallet_balance' + && request.idempotency_key === record.idempotency_key + && typeof request.selection === 'object' + && request.selection !== null + && typeof request.payload === 'object' + && request.payload !== null + && (request.payload.operation === 'add' || request.payload.operation === 'deduct') + && Number.isFinite(request.payload.amount) + && request.payload.amount > 0 +} + +function parsePendingRequest(serialized: string): PendingUserBatchWalletRequest { + try { + const value: unknown = JSON.parse(serialized) + if (isPendingRequest(value)) return value + } catch { + // Invalid persisted state must not allow a fresh adjustment to be sent. + } + throw new InvalidPendingWalletRequestError() +} + +function stableSerialize(value: unknown): string { + if (Array.isArray(value)) { + return `[${value.map((item) => item === undefined ? 'null' : stableSerialize(item)).join(',')}]` + } + if (typeof value === 'object' && value !== null) { + const fields = Object.entries(value as Record) + .filter(([, item]) => item !== undefined) + .sort(([left], [right]) => (left < right ? -1 : left > right ? 1 : 0)) + return `{${fields.map(([key, item]) => `${JSON.stringify(key)}:${stableSerialize(item)}`).join(',')}}` + } + return JSON.stringify(value) ?? 'null' +} + +export function matchesPendingWalletRequest( + pending: PendingUserBatchWalletRequest, + request: UserBatchWalletAdjustmentRequest, +): boolean { + const { idempotency_key: _key, ...pendingPayload } = pending.request + return stableSerialize(pendingPayload) === stableSerialize(request) +} + +export function createUserBatchWalletRetryCoordinator(options: CoordinatorOptions = {}) { + const storage = 'storage' in options ? options.storage ?? null : browserPersistentStorage() + const fallback = options.fallback ?? inMemoryFallback + const createKey = options.createKey ?? secureRandomUUID + const getScope = options.scope ?? (() => 'default') + + function getStorageKey(): string { + const scope = getScope() + if (!scope) throw new WalletIdempotencyScopeUnavailableError() + return `${STORAGE_KEY}:${encodeURIComponent(scope)}` + } + + function readPending(storageKey: string): PendingUserBatchWalletRequest | null { + let serialized: string | null = null + let storageReadFailed = false + try { + serialized = storage?.getItem(storageKey) ?? null + } catch { + storageReadFailed = true + serialized = null + } + serialized ??= fallback.get(storageKey) ?? null + if (serialized === null && (!storage || storageReadFailed)) { + throw new WalletIdempotencyPersistenceUnavailableError() + } + return serialized === null ? null : parsePendingRequest(serialized) + } + + function persist(storageKey: string, pending: PendingUserBatchWalletRequest): void { + const serialized = JSON.stringify(pending) + if (!storage) throw new WalletIdempotencyPersistenceUnavailableError() + try { + storage.setItem(storageKey, serialized) + if (storage.getItem(storageKey) !== serialized) { + throw new Error('Stored wallet batch request could not be verified') + } + } catch { + throw new WalletIdempotencyPersistenceUnavailableError() + } + fallback.set(storageKey, serialized) + } + + function clear(storageKey: string): void { + fallback.delete(storageKey) + try { + storage?.removeItem(storageKey) + } catch { + // A stale persisted request is safe to replay and will fail closed on mismatch. + } + } + + function getOrCreate( + storageKey: string, + request: UserBatchWalletAdjustmentRequest, + ): UserBatchBalanceActionRequest { + const pending = readPending(storageKey) + if (pending) { + if (!matchesPendingWalletRequest(pending, request)) { + throw new UnresolvedWalletRequestMismatchError(pending) + } + persist(storageKey, pending) + return pending.request + } + + const idempotencyKey = createKey() + if (!idempotencyKey) throw new WalletIdempotencyUnavailableError() + const keyedRequest: UserBatchBalanceActionRequest = { + ...request, + idempotency_key: idempotencyKey, + } + persist(storageKey, { idempotency_key: idempotencyKey, request: keyedRequest }) + return keyedRequest + } + + async function sendAndResolve( + request: UserBatchBalanceActionRequest, + send: (request: UserBatchBalanceActionRequest) => Promise, + storageKey: string, + ): Promise { + const response = await send(request) + if (!response.interrupted) clear(storageKey) + return response + } + + return { + getPending() { + return readPending(getStorageKey()) + }, + async execute( + request: UserBatchWalletAdjustmentRequest, + send: (request: UserBatchBalanceActionRequest) => Promise, + ) { + const storageKey = getStorageKey() + const keyedRequest = getOrCreate(storageKey, request) + if (getStorageKey() !== storageKey) throw new WalletIdempotencyScopeChangedError() + return sendAndResolve(keyedRequest, send, storageKey) + }, + async retry( + send: (request: UserBatchBalanceActionRequest) => Promise, + ): Promise { + const storageKey = getStorageKey() + const pending = readPending(storageKey) + if (!pending) return null + persist(storageKey, pending) + if (getStorageKey() !== storageKey) throw new WalletIdempotencyScopeChangedError() + return sendAndResolve(pending.request, send, storageKey) + }, + } +} From 390b73d4b6a95c65fcf986193c958be6ad4e4bf4 Mon Sep 17 00:00:00 2001 From: RWDai <27391645+RWDai@users.noreply.github.com> Date: Wed, 23 Sep 2026 20:34:25 +0800 Subject: [PATCH 5/9] test(data): include new migration in version expectation --- crates/aether-data/runtime/src/lifecycle/migrate/tests.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/crates/aether-data/runtime/src/lifecycle/migrate/tests.rs b/crates/aether-data/runtime/src/lifecycle/migrate/tests.rs index 02fa33e24..8d375d1f5 100644 --- a/crates/aether-data/runtime/src/lifecycle/migrate/tests.rs +++ b/crates/aether-data/runtime/src/lifecycle/migrate/tests.rs @@ -1575,6 +1575,7 @@ fn pending_migrations_from_applied_skips_versions_already_applied() { 20260901000000, 20260903000000, 20260908000000, + 20260923000000, ] ); } From 4c07d9fcfb6790babb094d81d81f6dbfcc1b7da2 Mon Sep 17 00:00:00 2001 From: RWDai <27391645+RWDai@users.noreply.github.com> Date: Mon, 28 Sep 2026 17:56:27 +0800 Subject: [PATCH 6/9] fix(admin): floor batch deductions at zero --- .../state/runtime/wallet/balance_mutations.rs | 6 ++++- .../src/tests/control/admin/users_batch.rs | 6 ++--- .../adapters/postgres/src/wallet.rs | 22 +++++++++++++------ 3 files changed, 23 insertions(+), 11 deletions(-) diff --git a/apps/aether-gateway/src/state/runtime/wallet/balance_mutations.rs b/apps/aether-gateway/src/state/runtime/wallet/balance_mutations.rs index 70a9ee1eb..c6b6b1b39 100644 --- a/apps/aether-gateway/src/state/runtime/wallet/balance_mutations.rs +++ b/apps/aether-gateway/src/state/runtime/wallet/balance_mutations.rs @@ -226,7 +226,11 @@ impl AppState { let before_gift = wallet.gift_balance; let before_total = before_recharge + before_gift; let amount_usd = if clamp_deduction_to_available_balance && amount_usd < 0.0 { - -(-amount_usd).min(before_total.max(0.0)) + if before_total < 0.0 { + -before_total + } else { + -(-amount_usd).min(before_total) + } } else { amount_usd }; diff --git a/apps/aether-gateway/src/tests/control/admin/users_batch.rs b/apps/aether-gateway/src/tests/control/admin/users_batch.rs index bb34d61c8..35efa942e 100644 --- a/apps/aether-gateway/src/tests/control/admin/users_batch.rs +++ b/apps/aether-gateway/src/tests/control/admin/users_batch.rs @@ -581,7 +581,7 @@ async fn gateway_reports_wallet_limit_lookup_failure_as_unprocessed() { } #[tokio::test] -async fn gateway_reports_missing_wallet_and_skips_zero_delta_for_non_positive_balance() { +async fn gateway_reports_missing_wallet_and_floors_negative_balance_on_deduction() { let state = AppState::new() .expect("gateway should build") .with_auth_users_for_tests([sample_user("user-negative"), sample_user("user-no-wallet")]) @@ -607,8 +607,8 @@ async fn gateway_reports_missing_wallet_and_skips_zero_delta_for_non_positive_ba assert_eq!(result["failures"][0]["reason"], "用户钱包不可用"); let wallet = wallet_detail(&client, &gateway_url, "user-negative").await; - assert_eq!(wallet["balance"], -1.0); - assert_eq!(wallet["total_adjusted"], 0.0); + assert_eq!(wallet["balance"], 0.0); + assert_eq!(wallet["total_adjusted"], 1.0); gateway_handle.abort(); } diff --git a/crates/aether-data/adapters/postgres/src/wallet.rs b/crates/aether-data/adapters/postgres/src/wallet.rs index 23f524d54..8e89f7a1d 100644 --- a/crates/aether-data/adapters/postgres/src/wallet.rs +++ b/crates/aether-data/adapters/postgres/src/wallet.rs @@ -812,7 +812,12 @@ impl SqlxWalletRepository { fn effective_wallet_adjustment_amount(input: &AdjustWalletBalanceInput, before_total: f64) -> f64 { if input.clamp_deduction_to_available_balance && input.amount_usd < 0.0 { - -(-input.amount_usd).min(before_total.max(0.0)) + if before_total < 0.0 { + // Clear legacy negative totals to zero and record the actual ledger delta. + -before_total + } else { + -(-input.amount_usd).min(before_total) + } } else { input.amount_usd } @@ -9598,7 +9603,7 @@ mod tests { }; assert_eq!(effective_wallet_adjustment_amount(&input, 13.0), -13.0); assert_eq!(effective_wallet_adjustment_amount(&input, 0.0), -0.0); - assert_eq!(effective_wallet_adjustment_amount(&input, -1.0), -0.0); + assert_eq!(effective_wallet_adjustment_amount(&input, -1.0), 1.0); let legacy_input = AdjustWalletBalanceInput { clamp_deduction_to_available_balance: false, @@ -9681,22 +9686,25 @@ mod tests { amount_usd: -1.0, balance_type: "recharge".to_string(), operator_id: Some("admin-user".to_string()), - description: Some("bulk deduction from negative balance".to_string()), + description: Some("bulk deduction floors a negative balance".to_string()), clamp_deduction_to_available_balance: true, batch_context: None, }) .await - .expect("legacy negative wallet should remain usable") + .expect("legacy negative wallet should be floored at zero") .expect("wallet should still exist"); - assert_eq!(wallet.balance + wallet.gift_balance, -1.0); - assert!(transaction.is_none()); + assert_eq!(wallet.balance + wallet.gift_balance, 0.0); + let transaction = transaction.expect("negative balance correction should be ledgered"); + assert_eq!(transaction.amount, 1.0); + assert_eq!(transaction.balance_before, -1.0); + assert_eq!(transaction.balance_after, 0.0); let transaction_count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM wallet_transactions WHERE wallet_id = $1") .bind(&wallet_id) .fetch_one(&pool) .await .expect("ledger row count should be readable"); - assert_eq!(transaction_count, 1); + assert_eq!(transaction_count, 2); pool.close().await; } From c96f8272f46b9433a1d9b2737a152b23f3083cf5 Mon Sep 17 00:00:00 2001 From: RWDai <27391645+RWDai@users.noreply.github.com> Date: Mon, 28 Sep 2026 18:12:21 +0800 Subject: [PATCH 7/9] test(ci): run batch wallet clamp against postgres --- .github/workflows/rust-ci.yml | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/.github/workflows/rust-ci.yml b/.github/workflows/rust-ci.yml index a01af779f..c03241452 100644 --- a/.github/workflows/rust-ci.yml +++ b/.github/workflows/rust-ci.yml @@ -556,6 +556,13 @@ jobs: AETHER_TEST_DATABASE_URL: postgres://aether:aether@127.0.0.1:5432/aether_test run: cargo test -p aether-data-postgres live_payment_callback --lib -- --ignored --nocapture + - name: Run Postgres batch wallet deduction regression + env: + RUSTC_WRAPPER: sccache + SCCACHE_GHA_ENABLED: "true" + AETHER_TEST_DATABASE_URL: postgres://aether:aether@127.0.0.1:5432/aether_test + run: cargo test -p aether-data-postgres live_bulk_wallet_adjustment_persists_actual_delta_and_skips_zero_ledger --lib -- --ignored --nocapture + - name: Run Postgres API key lifecycle tests env: RUSTC_WRAPPER: sccache From 74072e5007d08fe2a5c53d90ea94ee5f4f24dffa Mon Sep 17 00:00:00 2001 From: RWDai <27391645+RWDai@users.noreply.github.com> Date: Mon, 28 Sep 2026 18:23:20 +0800 Subject: [PATCH 8/9] test(postgres): decode ledger amount as float8 --- .../adapters/postgres/src/wallet.rs | 21 +++++++++++++------ 1 file changed, 15 insertions(+), 6 deletions(-) diff --git a/crates/aether-data/adapters/postgres/src/wallet.rs b/crates/aether-data/adapters/postgres/src/wallet.rs index 8e89f7a1d..ddf829a32 100644 --- a/crates/aether-data/adapters/postgres/src/wallet.rs +++ b/crates/aether-data/adapters/postgres/src/wallet.rs @@ -9641,12 +9641,13 @@ mod tests { assert_eq!(transaction.balance_before, 13.0); assert_eq!(transaction.balance_after, 0.0); assert_eq!(wallet.balance + wallet.gift_balance, 0.0); - let persisted_amount: f64 = - sqlx::query_scalar("SELECT amount FROM wallet_transactions WHERE id = $1") - .bind(&transaction.id) - .fetch_one(&pool) - .await - .expect("ledger should store the effective deduction"); + let persisted_amount: f64 = sqlx::query_scalar( + "SELECT amount::double precision FROM wallet_transactions WHERE id = $1", + ) + .bind(&transaction.id) + .fetch_one(&pool) + .await + .expect("ledger should store the effective deduction"); assert_eq!(persisted_amount, -13.0); let (wallet, transaction) = repository @@ -9698,6 +9699,14 @@ mod tests { assert_eq!(transaction.amount, 1.0); assert_eq!(transaction.balance_before, -1.0); assert_eq!(transaction.balance_after, 0.0); + let persisted_correction_amount: f64 = sqlx::query_scalar( + "SELECT amount::double precision FROM wallet_transactions WHERE id = $1", + ) + .bind(&transaction.id) + .fetch_one(&pool) + .await + .expect("ledger should store the negative-balance correction"); + assert_eq!(persisted_correction_amount, 1.0); let transaction_count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM wallet_transactions WHERE wallet_id = $1") .bind(&wallet_id) From bc0e9f94e2a5dc6df7c7624e6905e3f004470a86 Mon Sep 17 00:00:00 2001 From: RWDai <27391645+RWDai@users.noreply.github.com> Date: Mon, 28 Sep 2026 18:51:14 +0800 Subject: [PATCH 9/9] fix(admin): serialize wallet batches across tabs --- .../components/UserBatchActionDialog.vue | 10 ++ .../userBatchWalletIdempotency.spec.ts | 110 +++++++++++++++--- .../users/utils/userBatchWalletIdempotency.ts | 66 +++++++++-- frontend/src/i18n/legacy-admin-messages.ts | 2 + 4 files changed, 163 insertions(+), 25 deletions(-) diff --git a/frontend/src/features/users/components/UserBatchActionDialog.vue b/frontend/src/features/users/components/UserBatchActionDialog.vue index fff9293be..ac9175b37 100644 --- a/frontend/src/features/users/components/UserBatchActionDialog.vue +++ b/frontend/src/features/users/components/UserBatchActionDialog.vue @@ -172,7 +172,9 @@ import { buildUserBatchBalanceAdjustmentPayload } from '@/api/users' import { createUserBatchWalletRetryCoordinator, matchesPendingWalletRequest, + WalletIdempotencyCoordinationUnavailableError, WalletIdempotencyPersistenceUnavailableError, + WalletIdempotencyRequestInProgressError, WalletIdempotencyScopeChangedError, WalletIdempotencyScopeUnavailableError, WalletIdempotencyUnavailableError, @@ -438,6 +440,10 @@ async function executeBatchAction(): Promise { refreshPendingWalletBatch() if (err instanceof WalletIdempotencyPersistenceUnavailableError) { warning(legacyT('浏览器无法安全保存钱包批量请求,本次请求未发送。')) + } else if (err instanceof WalletIdempotencyRequestInProgressError) { + warning(legacyT('另一个标签页正在处理钱包批量调整,请稍后刷新状态再试。此次未发送新请求。')) + } else if (err instanceof WalletIdempotencyCoordinationUnavailableError) { + warning(legacyT('当前浏览器无法保护跨标签页的钱包批量请求,请使用支持此功能的浏览器。请求未发送。')) } else if ( err instanceof WalletIdempotencyUnavailableError || err instanceof WalletIdempotencyScopeUnavailableError @@ -468,6 +474,10 @@ async function retryPendingWalletBatch(): Promise { refreshPendingWalletBatch() if (err instanceof WalletIdempotencyPersistenceUnavailableError) { warning(legacyT('浏览器无法安全保存钱包批量请求,本次请求未发送。')) + } else if (err instanceof WalletIdempotencyRequestInProgressError) { + warning(legacyT('另一个标签页正在处理钱包批量调整,请稍后刷新状态再试。此次未发送新请求。')) + } else if (err instanceof WalletIdempotencyCoordinationUnavailableError) { + warning(legacyT('当前浏览器无法保护跨标签页的钱包批量请求,请使用支持此功能的浏览器。请求未发送。')) } else if ( err instanceof WalletIdempotencyScopeUnavailableError || err instanceof WalletIdempotencyScopeChangedError diff --git a/frontend/src/features/users/utils/__tests__/userBatchWalletIdempotency.spec.ts b/frontend/src/features/users/utils/__tests__/userBatchWalletIdempotency.spec.ts index 2452fe98e..ad4ae4810 100644 --- a/frontend/src/features/users/utils/__tests__/userBatchWalletIdempotency.spec.ts +++ b/frontend/src/features/users/utils/__tests__/userBatchWalletIdempotency.spec.ts @@ -2,6 +2,8 @@ import { beforeEach, describe, expect, it, vi } from 'vitest' import type { UserBatchActionResponse, UserBatchBalanceActionRequest } from '@/api/users' import { createUserBatchWalletRetryCoordinator, + WalletIdempotencyCoordinationUnavailableError, + WalletIdempotencyRequestInProgressError, UnresolvedWalletRequestMismatchError, WalletIdempotencyPersistenceUnavailableError, WalletIdempotencyUnavailableError, @@ -23,6 +25,31 @@ const walletRequest = { } const defaultStorageKey = 'admin.users.batch.wallet-adjustment.pending.v1:default' let testFallback: Map +let testLockManager: Pick + +type CoordinatorOptions = NonNullable[0]> + +function createTestLockManager(): Pick { + const heldNames = new Set() + const request = async ( + name: string, + _options: LockOptions, + callback: (lock: Lock | null) => Promise, + ): Promise => { + if (heldNames.has(name)) return callback(null) + heldNames.add(name) + try { + return await callback({ name, mode: 'exclusive' } as Lock) + } finally { + heldNames.delete(name) + } + } + return { request } as unknown as Pick +} + +function createCoordinator(options: Omit = {}) { + return createUserBatchWalletRetryCoordinator({ ...options, lockManager: testLockManager }) +} function response(interrupted = false): UserBatchActionResponse { return { total: 2, success: 1, failed: 0, failures: [], interrupted } @@ -33,11 +60,12 @@ describe('user batch wallet idempotency', () => { sessionStorage.clear() localStorage.clear() testFallback = new Map() + testLockManager = createTestLockManager() }) it('serializes the key with the exact top-level wallet request before sending', async () => { const storage = createStorage() - const coordinator = createUserBatchWalletRetryCoordinator({ + const coordinator = createCoordinator({ storage, fallback: testFallback, createKey: () => 'wallet-key-1', @@ -58,9 +86,57 @@ describe('user batch wallet idempotency', () => { }) }) + it('rejects a second tab while the same wallet batch request is in progress', async () => { + const storage = createStorage() + const first = createCoordinator({ + storage, + fallback: testFallback, + createKey: () => 'wallet-key-first-tab', + }) + const second = createCoordinator({ + storage, + fallback: testFallback, + createKey: () => 'wallet-key-second-tab', + }) + let resolveFirst!: (value: UserBatchActionResponse) => void + const firstSend = vi.fn(() => new Promise((resolve) => { + resolveFirst = resolve + })) + const secondSend = vi.fn(async () => response()) + const firstExecution = first.execute(walletRequest, firstSend) + + await vi.waitFor(() => expect(firstSend).toHaveBeenCalledOnce()) + await expect(second.execute(walletRequest, secondSend)).rejects.toBeInstanceOf( + WalletIdempotencyRequestInProgressError, + ) + expect(secondSend).not.toHaveBeenCalled() + + resolveFirst(response()) + await expect(firstExecution).resolves.toEqual(response()) + expect(firstSend).toHaveBeenCalledOnce() + expect(storage.getItem(defaultStorageKey)).toBeNull() + }) + + it('fails closed when cross-tab request coordination is unavailable', async () => { + const storage = createStorage() + const send = vi.fn(async () => response()) + const coordinator = createUserBatchWalletRetryCoordinator({ + storage, + fallback: testFallback, + lockManager: null, + createKey: () => 'wallet-key-without-locks', + }) + + await expect(coordinator.execute(walletRequest, send)).rejects.toBeInstanceOf( + WalletIdempotencyCoordinationUnavailableError, + ) + expect(send).not.toHaveBeenCalled() + expect(storage.getItem(defaultStorageKey)).toBeNull() + }) + it('retains a transport failure and reopens with the exact request for retry', async () => { const storage = createStorage() - const first = createUserBatchWalletRetryCoordinator({ + const first = createCoordinator({ storage, fallback: testFallback, createKey: () => 'wallet-key-2', @@ -68,7 +144,7 @@ describe('user batch wallet idempotency', () => { const sendFailure = new Error('connection lost') await expect(first.execute(walletRequest, async () => { throw sendFailure })).rejects.toBe(sendFailure) - const reopened = createUserBatchWalletRetryCoordinator({ + const reopened = createCoordinator({ storage, fallback: testFallback, createKey: () => 'must-not-be-used', @@ -86,7 +162,7 @@ describe('user batch wallet idempotency', () => { }) it('keeps unresolved requests in persistent browser storage across coordinators', async () => { - const first = createUserBatchWalletRetryCoordinator({ + const first = createCoordinator({ createKey: () => 'wallet-key-persistent', scope: () => 'admin-1', }) @@ -94,7 +170,7 @@ describe('user batch wallet idempotency', () => { throw new Error('connection lost') })).rejects.toThrow('connection lost') - const reopened = createUserBatchWalletRetryCoordinator({ scope: () => 'admin-1' }) + const reopened = createCoordinator({ scope: () => 'admin-1' }) const pending = reopened.getPending() expect(pending?.request).toEqual({ ...walletRequest, @@ -111,7 +187,7 @@ describe('user batch wallet idempotency', () => { it('keeps unresolved requests isolated by authenticated administrator', async () => { const storage = createStorage() - const adminA = createUserBatchWalletRetryCoordinator({ + const adminA = createCoordinator({ storage, fallback: testFallback, scope: () => 'admin-a', @@ -120,7 +196,7 @@ describe('user batch wallet idempotency', () => { await expect(adminA.execute(walletRequest, async () => { throw new Error('connection lost') })) .rejects.toThrow('connection lost') - const adminB = createUserBatchWalletRetryCoordinator({ + const adminB = createCoordinator({ storage, fallback: testFallback, scope: () => 'admin-b', @@ -137,7 +213,7 @@ describe('user batch wallet idempotency', () => { const storage = createStorage() let scopeReads = 0 const send = vi.fn(async () => response()) - const coordinator = createUserBatchWalletRetryCoordinator({ + const coordinator = createCoordinator({ storage, fallback: testFallback, scope: () => (++scopeReads === 1 ? 'admin-a' : 'admin-b'), @@ -153,13 +229,13 @@ describe('user batch wallet idempotency', () => { it('retains interrupted requests and reuses their key until a terminal response', async () => { const storage = createStorage() - const first = createUserBatchWalletRetryCoordinator({ + const first = createCoordinator({ storage, fallback: testFallback, createKey: () => 'wallet-key-3', }) await first.execute(walletRequest, async () => response(true)) - const reopened = createUserBatchWalletRetryCoordinator({ storage, fallback: testFallback }) + const reopened = createCoordinator({ storage, fallback: testFallback }) const pending = reopened.getPending() const send = vi.fn(async () => response(true)) @@ -180,14 +256,14 @@ describe('user batch wallet idempotency', () => { ...walletRequest, selection: { filters: { search: 'active', is_active: undefined } }, } - const first = createUserBatchWalletRetryCoordinator({ + const first = createCoordinator({ storage, fallback: testFallback, createKey: () => 'wallet-key-filter', }) await first.execute(requestWithUndefinedField, async () => response(true)) - const reopened = createUserBatchWalletRetryCoordinator({ storage, fallback: testFallback }) + const reopened = createCoordinator({ storage, fallback: testFallback }) const send = vi.fn(async () => response()) await reopened.execute({ ...walletRequest, @@ -204,7 +280,7 @@ describe('user batch wallet idempotency', () => { it('blocks changed payloads while unresolved and gives a later adjustment a new key', async () => { const storage = createStorage() let nextKey = 0 - const coordinator = createUserBatchWalletRetryCoordinator({ + const coordinator = createCoordinator({ storage, fallback: testFallback, createKey: () => `wallet-key-${++nextKey}`, @@ -240,7 +316,7 @@ describe('user batch wallet idempotency', () => { removeItem: () => undefined, } const send = vi.fn(async () => response(true)) - const coordinator = createUserBatchWalletRetryCoordinator({ + const coordinator = createCoordinator({ storage: unavailableStorage, fallback: testFallback, createKey: () => 'wallet-key-fallback', @@ -261,7 +337,7 @@ describe('user batch wallet idempotency', () => { removeItem: () => undefined, } const send = vi.fn(async () => response()) - const coordinator = createUserBatchWalletRetryCoordinator({ + const coordinator = createCoordinator({ storage: mismatchedStorage, fallback: testFallback, createKey: () => 'wallet-key-readback', @@ -281,7 +357,7 @@ describe('user batch wallet idempotency', () => { removeItem: () => undefined, } const send = vi.fn(async () => response()) - const coordinator = createUserBatchWalletRetryCoordinator({ + const coordinator = createCoordinator({ storage: unavailableStorage, fallback: testFallback, createKey: () => 'must-not-be-used', @@ -297,7 +373,7 @@ describe('user batch wallet idempotency', () => { it('fails closed when secure UUID generation is unavailable', async () => { const storage = createStorage() const send = vi.fn(async () => response()) - const coordinator = createUserBatchWalletRetryCoordinator({ + const coordinator = createCoordinator({ storage, fallback: testFallback, createKey: () => { throw new WalletIdempotencyUnavailableError() }, diff --git a/frontend/src/features/users/utils/userBatchWalletIdempotency.ts b/frontend/src/features/users/utils/userBatchWalletIdempotency.ts index d81fc412c..c0117525c 100644 --- a/frontend/src/features/users/utils/userBatchWalletIdempotency.ts +++ b/frontend/src/features/users/utils/userBatchWalletIdempotency.ts @@ -23,6 +23,7 @@ interface CoordinatorOptions { storage?: StringStorage | null createKey?: () => string fallback?: Map + lockManager?: Pick | null scope?: () => string | null } @@ -50,6 +51,20 @@ export class WalletIdempotencyPersistenceUnavailableError extends Error { } } +export class WalletIdempotencyCoordinationUnavailableError extends Error { + constructor() { + super('Cross-tab wallet request coordination is unavailable') + this.name = 'WalletIdempotencyCoordinationUnavailableError' + } +} + +export class WalletIdempotencyRequestInProgressError extends Error { + constructor() { + super('A wallet batch request is already in progress in another tab') + this.name = 'WalletIdempotencyRequestInProgressError' + } +} + export class WalletIdempotencyScopeUnavailableError extends Error { constructor() { super('The authenticated administrator identity is unavailable') @@ -79,6 +94,14 @@ function browserPersistentStorage(): StringStorage | null { } } +function browserLockManager(): Pick | null { + try { + return globalThis.navigator?.locks ?? null + } catch { + return null + } +} + function secureRandomUUID(): string { try { const cryptoApi = globalThis.crypto @@ -145,6 +168,9 @@ export function createUserBatchWalletRetryCoordinator(options: CoordinatorOption const storage = 'storage' in options ? options.storage ?? null : browserPersistentStorage() const fallback = options.fallback ?? inMemoryFallback const createKey = options.createKey ?? secureRandomUUID + const lockManager = 'lockManager' in options + ? options.lockManager ?? null + : browserLockManager() const getScope = options.scope ?? (() => 'default') function getStorageKey(): string { @@ -192,6 +218,26 @@ export function createUserBatchWalletRetryCoordinator(options: CoordinatorOption } } + async function withExclusiveLock(storageKey: string, task: () => Promise): Promise { + if (!lockManager) throw new WalletIdempotencyCoordinationUnavailableError() + + let taskStarted = false + try { + return await lockManager.request( + storageKey, + { mode: 'exclusive', ifAvailable: true }, + async (lock) => { + if (lock === null) throw new WalletIdempotencyRequestInProgressError() + taskStarted = true + return task() + }, + ) + } catch (error) { + if (taskStarted || error instanceof WalletIdempotencyRequestInProgressError) throw error + throw new WalletIdempotencyCoordinationUnavailableError() + } + } + function getOrCreate( storageKey: string, request: UserBatchWalletAdjustmentRequest, @@ -234,19 +280,23 @@ export function createUserBatchWalletRetryCoordinator(options: CoordinatorOption send: (request: UserBatchBalanceActionRequest) => Promise, ) { const storageKey = getStorageKey() - const keyedRequest = getOrCreate(storageKey, request) - if (getStorageKey() !== storageKey) throw new WalletIdempotencyScopeChangedError() - return sendAndResolve(keyedRequest, send, storageKey) + return withExclusiveLock(storageKey, async () => { + const keyedRequest = getOrCreate(storageKey, request) + if (getStorageKey() !== storageKey) throw new WalletIdempotencyScopeChangedError() + return sendAndResolve(keyedRequest, send, storageKey) + }) }, async retry( send: (request: UserBatchBalanceActionRequest) => Promise, ): Promise { const storageKey = getStorageKey() - const pending = readPending(storageKey) - if (!pending) return null - persist(storageKey, pending) - if (getStorageKey() !== storageKey) throw new WalletIdempotencyScopeChangedError() - return sendAndResolve(pending.request, send, storageKey) + return withExclusiveLock(storageKey, async () => { + const pending = readPending(storageKey) + if (!pending) return null + persist(storageKey, pending) + if (getStorageKey() !== storageKey) throw new WalletIdempotencyScopeChangedError() + return sendAndResolve(pending.request, send, storageKey) + }) }, } } diff --git a/frontend/src/i18n/legacy-admin-messages.ts b/frontend/src/i18n/legacy-admin-messages.ts index bad7721f1..c81f44423 100644 --- a/frontend/src/i18n/legacy-admin-messages.ts +++ b/frontend/src/i18n/legacy-admin-messages.ts @@ -613,6 +613,8 @@ export const legacyAdminEnglishMessages: Record = { '系统调账': 'System adjustment', '退款扣减': 'Refund debit', '退款回补': 'Refund recredit', + '另一个标签页正在处理钱包批量调整,请稍后刷新状态再试。此次未发送新请求。': 'Another tab is processing a wallet batch adjustment. Refresh the status and try again later. No new request was sent from this tab.', + '当前浏览器无法保护跨标签页的钱包批量请求,请使用支持此功能的浏览器。请求未发送。': 'This browser cannot protect wallet batch requests across tabs. Use a browser that supports this feature. The request was not sent.', '兑换码批次已创建': 'Redemption code batch created', 'CSV 已导出': 'CSV exported', '批次已停用': 'Batch disabled',
+ {{ legacyT('请输入大于 0 的有限金额') }} +
+ {{ legacyT('扣减超过单个用户可用余额时,该用户余额将归零。') }} +
{{ legacyT('已完成用户 ID') }}:{{ userIds(result.completed_user_ids) }}
{{ legacyT('结果待核对用户 ID') }}:{{ userIds(result.uncertain_user_ids) }}
{{ legacyT('尚未执行用户 ID') }}:{{ userIds(result.unprocessed_user_ids) }}