From bcd121d4471cd97c92a7b245ebd51177ff9c2df3 Mon Sep 17 00:00:00 2001 From: RWDai <27391645+RWDai@users.noreply.github.com> Date: Wed, 23 Sep 2026 19:43:18 +0800 Subject: [PATCH] fix(admin): make bulk wallet batches idempotent --- apps/aether-gateway/src/data/state/mod.rs | 12 +- apps/aether-gateway/src/data/state/runtime.rs | 113 +++- .../src/handlers/admin/request/state.rs | 58 ++ .../src/handlers/admin/users/batch.rs | 455 ++++++++++++++- apps/aether-gateway/src/state/app.rs | 15 + apps/aether-gateway/src/state/core.rs | 6 + .../state/runtime/wallet/balance_mutations.rs | 191 +++++++ apps/aether-gateway/src/state/testing.rs | 8 + .../src/tests/control/admin/users_batch.rs | 183 ++++++- ...000_add_admin_wallet_batch_idempotency.sql | 12 + .../adapters/postgres/src/wallet.rs | 518 +++++++++++++++++- .../contracts/src/repository/wallet/types.rs | 92 ++++ .../postgres/baseline/005_wallet_billing.sql | 14 + .../schema/logical/005_wallet_billing.toml | 44 ++ .../runtime/src/repository/wallet/mod.rs | 43 +- frontend/src/api/users.ts | 1 + .../components/UserBatchActionDialog.vue | 162 +++++- .../userBatchWalletIdempotency.spec.ts | 312 +++++++++++ .../users/utils/userBatchWalletIdempotency.ts | 252 +++++++++ 19 files changed, 2396 insertions(+), 95 deletions(-) create mode 100644 crates/aether-data/adapters/postgres/migrations/20260923000000_add_admin_wallet_batch_idempotency.sql create mode 100644 frontend/src/features/users/utils/__tests__/userBatchWalletIdempotency.spec.ts create mode 100644 frontend/src/features/users/utils/userBatchWalletIdempotency.ts diff --git a/apps/aether-gateway/src/data/state/mod.rs b/apps/aether-gateway/src/data/state/mod.rs index 0a455684b..6837220c7 100644 --- a/apps/aether-gateway/src/data/state/mod.rs +++ b/apps/aether-gateway/src/data/state/mod.rs @@ -62,8 +62,9 @@ pub(crate) use aether_data::repository::users::{ StoredUserPreferenceRecord, StoredUserSessionRecord, }; use aether_data::repository::wallet::{ - AdjustWalletBalanceInput, AdminPaymentOrderListQuery, AdminRedeemCodeBatchListQuery, - AdminRedeemCodeListQuery, AdminWalletLedgerQuery, AdminWalletListQuery, + AdjustWalletBalanceInBatchInput, AdjustWalletBalanceInput, AdminPaymentOrderListQuery, + AdminRedeemCodeBatchListQuery, AdminRedeemCodeListQuery, + AdminUserWalletBalanceBatchUserOutcome, AdminWalletLedgerQuery, AdminWalletListQuery, AdminWalletRefundRequestListQuery, CompareAndSwapPaymentOrderStripeClientSecretInput, CompleteAdminWalletRefundInput, CreateAdminRedeemCodeBatchInput, CreateAdminRedeemCodeBatchResult, CreateManualWalletRechargeInput, @@ -72,13 +73,14 @@ use aether_data::repository::wallet::{ CreateWalletRefundRequestOutcome, CreditAdminPaymentOrderInput, DeleteAdminRedeemCodeBatchInput, DisableAdminRedeemCodeBatchInput, DisableAdminRedeemCodeInput, FailAdminWalletRefundInput, FailWalletRechargeCheckoutInput, InitializeAuthWalletOutcome, + PrepareAdminUserWalletBalanceBatchInput, PrepareAdminUserWalletBalanceBatchOutcome, ProcessAdminWalletRefundInput, ProcessPaymentCallbackInput, ProcessPaymentCallbackOutcome, ReclaimWalletRechargeCheckoutInput, RedeemWalletCodeInput, RedeemWalletCodeOutcome, StoredAdminPaymentCallback, StoredAdminPaymentCallbackPage, StoredAdminPaymentOrder, StoredAdminPaymentOrderPage, StoredAdminRedeemCode, StoredAdminRedeemCodeBatch, - StoredAdminRedeemCodeBatchPage, StoredAdminRedeemCodePage, StoredAdminWalletLedgerPage, - StoredAdminWalletListPage, StoredAdminWalletRefund, StoredAdminWalletRefundPage, - StoredAdminWalletRefundRequestPage, StoredAdminWalletTransaction, + StoredAdminRedeemCodeBatchPage, StoredAdminRedeemCodePage, StoredAdminUserWalletBalanceBatch, + StoredAdminWalletLedgerPage, StoredAdminWalletListPage, StoredAdminWalletRefund, + StoredAdminWalletRefundPage, StoredAdminWalletRefundRequestPage, StoredAdminWalletTransaction, StoredAdminWalletTransactionPage, StoredWalletDailyUsageLedger, StoredWalletDailyUsageLedgerPage, StoredWalletSnapshot, UpdateAdminWalletRefundGatewayInput, UpdateWalletRechargeCheckoutInput, WalletLookupKey, WalletMutationOutcome, diff --git a/apps/aether-gateway/src/data/state/runtime.rs b/apps/aether-gateway/src/data/state/runtime.rs index b06449a80..5cac1999b 100644 --- a/apps/aether-gateway/src/data/state/runtime.rs +++ b/apps/aether-gateway/src/data/state/runtime.rs @@ -1,9 +1,10 @@ use super::{ read_decision_trace, read_provider_transport_snapshot, read_request_candidate_trace, - AdjustWalletBalanceInput, AdminBillingCollectorRecord, AdminBillingCollectorWriteInput, - AdminBillingMutationOutcome, AdminBillingPresetApplyResult, AdminBillingRuleRecord, - AdminBillingRuleWriteInput, AdminPaymentOrderListQuery, AdminRedeemCodeBatchListQuery, - AdminRedeemCodeListQuery, AdminWalletLedgerQuery, AdminWalletListQuery, + AdjustWalletBalanceInBatchInput, AdjustWalletBalanceInput, AdminBillingCollectorRecord, + AdminBillingCollectorWriteInput, AdminBillingMutationOutcome, AdminBillingPresetApplyResult, + AdminBillingRuleRecord, AdminBillingRuleWriteInput, AdminPaymentOrderListQuery, + AdminRedeemCodeBatchListQuery, AdminRedeemCodeListQuery, + AdminUserWalletBalanceBatchUserOutcome, AdminWalletLedgerQuery, AdminWalletListQuery, AdminWalletRefundRequestListQuery, AnnouncementListQuery, AuditLogListQuery, BackgroundTaskListQuery, BackgroundTaskSummary, BillingModelContextCacheKey, BillingModelContextCacheState, BillingModelContextInflightState, BillingPlanRecord, @@ -17,24 +18,25 @@ use super::{ DisableAdminRedeemCodeBatchInput, DisableAdminRedeemCodeInput, FailAdminWalletRefundInput, FailWalletRechargeCheckoutInput, GatewayDataState, GatewayProviderTransportSnapshot, LocalVideoTaskReadResponse, PaymentGatewayConfigCasWriteInput, PaymentGatewayConfigRecord, - PaymentGatewayConfigWriteInput, PaymentGatewaySecretCasUpdate, ProcessAdminWalletRefundInput, - ProcessPaymentCallbackInput, ProcessPaymentCallbackOutcome, ReclaimWalletRechargeCheckoutInput, - ReconcileUsagePolicyCostInput, RedeemWalletCodeInput, RedeemWalletCodeOutcome, - ReleaseUsagePolicyRequestAdmissionInput, RequestAuditBundle, RequestCandidateTrace, - ReserveUsagePolicyCostInput, ReserveUsagePolicyCostOutcome, ReserveUsagePolicyRequestInput, - ReserveUsagePolicyRequestOutcome, StoredAdminAuditLogPage, StoredAdminPaymentCallbackPage, - StoredAdminPaymentOrder, StoredAdminPaymentOrderPage, StoredAdminRedeemCodeBatch, - StoredAdminRedeemCodeBatchPage, StoredAdminRedeemCodePage, StoredAdminWalletLedgerPage, - StoredAdminWalletListPage, StoredAdminWalletRefund, StoredAdminWalletRefundPage, - StoredAdminWalletRefundRequestPage, StoredAdminWalletTransaction, - StoredAdminWalletTransactionPage, StoredAnnouncement, StoredAnnouncementPage, - StoredBackgroundTaskEvent, StoredBackgroundTaskRun, StoredBackgroundTaskRunPage, - StoredBillingModelContext, StoredProviderQuotaSnapshot, StoredProviderUsageSummary, - StoredRequestUsageAudit, StoredSuspiciousActivity, StoredUsagePolicyCostReservation, - StoredUsagePolicyRequestAdmission, StoredUsageSettlement, StoredUserAuditLogPage, - StoredUserAuthRecord, StoredUserExportRow, StoredUserSummary, StoredVideoTask, - StoredWalletDailyUsageLedger, StoredWalletDailyUsageLedgerPage, StoredWalletSnapshot, - UpdateAdminWalletRefundGatewayInput, UpdateAnnouncementRecord, + PaymentGatewayConfigWriteInput, PaymentGatewaySecretCasUpdate, + PrepareAdminUserWalletBalanceBatchInput, PrepareAdminUserWalletBalanceBatchOutcome, + ProcessAdminWalletRefundInput, ProcessPaymentCallbackInput, ProcessPaymentCallbackOutcome, + ReclaimWalletRechargeCheckoutInput, ReconcileUsagePolicyCostInput, RedeemWalletCodeInput, + RedeemWalletCodeOutcome, ReleaseUsagePolicyRequestAdmissionInput, RequestAuditBundle, + RequestCandidateTrace, ReserveUsagePolicyCostInput, ReserveUsagePolicyCostOutcome, + ReserveUsagePolicyRequestInput, ReserveUsagePolicyRequestOutcome, StoredAdminAuditLogPage, + StoredAdminPaymentCallbackPage, StoredAdminPaymentOrder, StoredAdminPaymentOrderPage, + StoredAdminRedeemCodeBatch, StoredAdminRedeemCodeBatchPage, StoredAdminRedeemCodePage, + StoredAdminUserWalletBalanceBatch, StoredAdminWalletLedgerPage, StoredAdminWalletListPage, + StoredAdminWalletRefund, StoredAdminWalletRefundPage, StoredAdminWalletRefundRequestPage, + StoredAdminWalletTransaction, StoredAdminWalletTransactionPage, StoredAnnouncement, + StoredAnnouncementPage, StoredBackgroundTaskEvent, StoredBackgroundTaskRun, + StoredBackgroundTaskRunPage, StoredBillingModelContext, StoredProviderQuotaSnapshot, + StoredProviderUsageSummary, StoredRequestUsageAudit, StoredSuspiciousActivity, + StoredUsagePolicyCostReservation, StoredUsagePolicyRequestAdmission, StoredUsageSettlement, + StoredUserAuditLogPage, StoredUserAuthRecord, StoredUserExportRow, StoredUserSummary, + StoredVideoTask, StoredWalletDailyUsageLedger, StoredWalletDailyUsageLedgerPage, + StoredWalletSnapshot, UpdateAdminWalletRefundGatewayInput, UpdateAnnouncementRecord, UpdateWalletRechargeCheckoutInput, UpsertBackgroundTaskEvent, UpsertBackgroundTaskRun, UpsertUsageRecord, UpsertVideoTask, UsageSettlementInput, UserDailyQuotaAvailabilityRecord, UserPlanEntitlementRecord, VideoTaskLookupKey, VideoTaskModelCount, VideoTaskQueryFilter, @@ -1074,6 +1076,73 @@ impl GatewayDataState { } } + pub(crate) async fn prepare_admin_user_wallet_balance_batch( + &self, + input: PrepareAdminUserWalletBalanceBatchInput, + ) -> Result, DataLayerError> { + match &self.wallet_writer { + Some(repository) => repository + .prepare_admin_user_wallet_balance_batch(input) + .await + .map(Some), + None => Ok(None), + } + } + + pub(crate) async fn get_admin_user_wallet_balance_batch( + &self, + admin_user_id: &str, + idempotency_key: &str, + request_fingerprint: &str, + ) -> Result, DataLayerError> { + match &self.wallet_writer { + Some(repository) => { + repository + .get_admin_user_wallet_balance_batch( + admin_user_id, + idempotency_key, + request_fingerprint, + ) + .await + } + None => Ok(None), + } + } + + pub(crate) async fn adjust_admin_user_wallet_balance_batch_user( + &self, + input: AdjustWalletBalanceInBatchInput, + ) -> Result, DataLayerError> { + match &self.wallet_writer { + Some(repository) => repository + .adjust_admin_user_wallet_balance_batch_user(input) + .await + .map(Some), + None => Ok(None), + } + } + + pub(crate) async fn record_admin_user_wallet_balance_batch_failure( + &self, + admin_user_id: &str, + idempotency_key: &str, + user_id: &str, + reason: &str, + ) -> Result, DataLayerError> { + match &self.wallet_writer { + Some(repository) => repository + .record_admin_user_wallet_balance_batch_failure( + admin_user_id, + idempotency_key, + user_id, + reason, + ) + .await + .map(Some), + None => Ok(None), + } + } + pub(crate) async fn create_manual_wallet_recharge( &self, input: CreateManualWalletRechargeInput, diff --git a/apps/aether-gateway/src/handlers/admin/request/state.rs b/apps/aether-gateway/src/handlers/admin/request/state.rs index 7de1273e3..9564e03af 100644 --- a/apps/aether-gateway/src/handlers/admin/request/state.rs +++ b/apps/aether-gateway/src/handlers/admin/request/state.rs @@ -17,6 +17,64 @@ impl<'a> AdminAppState<'a> { pub(crate) fn cloned_app(&self) -> AppState { self.app.clone() } + + pub(crate) async fn get_admin_user_wallet_balance_batch( + &self, + admin_user_id: &str, + idempotency_key: &str, + request_fingerprint: &str, + ) -> Result< + Option, + GatewayError, + > { + self.app + .get_admin_user_wallet_balance_batch( + admin_user_id, + idempotency_key, + request_fingerprint, + ) + .await + } + + pub(crate) async fn prepare_admin_user_wallet_balance_batch( + &self, + input: aether_data::repository::wallet::PrepareAdminUserWalletBalanceBatchInput, + ) -> Result< + aether_data::repository::wallet::PrepareAdminUserWalletBalanceBatchOutcome, + GatewayError, + > { + self.app + .prepare_admin_user_wallet_balance_batch(input) + .await + } + + pub(crate) async fn record_admin_user_wallet_balance_batch_failure( + &self, + admin_user_id: &str, + idempotency_key: &str, + user_id: &str, + reason: &str, + ) -> Result + { + self.app + .record_admin_user_wallet_balance_batch_failure( + admin_user_id, + idempotency_key, + user_id, + reason, + ) + .await + } + + pub(crate) async fn adjust_admin_user_wallet_balance_batch_user( + &self, + input: aether_data::repository::wallet::AdjustWalletBalanceInBatchInput, + ) -> Result + { + self.app + .adjust_admin_user_wallet_balance_batch_user(input) + .await + } } impl<'a> AsRef for AdminAppState<'a> { diff --git a/apps/aether-gateway/src/handlers/admin/users/batch.rs b/apps/aether-gateway/src/handlers/admin/users/batch.rs index 22584fea8..d56b3f08d 100644 --- a/apps/aether-gateway/src/handlers/admin/users/batch.rs +++ b/apps/aether-gateway/src/handlers/admin/users/batch.rs @@ -8,6 +8,9 @@ use super::{ use crate::handlers::admin::request::{AdminAppState, AdminRequestContext}; use crate::handlers::admin::shared::attach_admin_audit_response; use crate::GatewayError; +use aether_data::repository::wallet::{ + AdminUserWalletBalanceBatchUserOutcome, PrepareAdminUserWalletBalanceBatchOutcome, +}; use axum::{ body::{Body, Bytes}, http, @@ -15,9 +18,10 @@ use axum::{ Json, }; use serde_json::{json, Value}; +use sha2::{Digest as _, Sha256}; use std::collections::{BTreeMap, BTreeSet}; -#[derive(Debug, Clone, Default, serde::Deserialize)] +#[derive(Debug, Clone, Default, serde::Deserialize, serde::Serialize)] struct AdminUserSelectionFilters { #[serde(default)] search: Option, @@ -29,7 +33,7 @@ struct AdminUserSelectionFilters { group_id: Option, } -#[derive(Debug, Clone, Default)] +#[derive(Debug, Clone, Default, serde::Serialize)] struct AdminUserSelectionRequest { user_ids: Vec, group_ids: Vec, @@ -42,6 +46,7 @@ struct AdminUserBatchActionRequest { selection: AdminUserSelectionRequest, action: String, payload: Option, + idempotency_key: Option, } #[derive(Debug, serde::Deserialize)] @@ -50,6 +55,8 @@ struct RawAdminUserBatchActionRequest { action: String, #[serde(default)] payload: Option, + #[serde(default)] + idempotency_key: Option, } #[derive(Debug, Clone, Default)] @@ -70,7 +77,7 @@ struct AdminUserSelectionItem { matched_by: Vec, } -#[derive(Debug, Clone, serde::Serialize)] +#[derive(Debug, Clone, serde::Deserialize, serde::Serialize)] struct AdminUserSelectionWarning { #[serde(rename = "type")] warning_type: String, @@ -117,6 +124,11 @@ enum AdminBatchWalletBalanceAdjustmentError { BalanceAdjustment, } +enum AdminBatchWalletLimitModeError { + WalletLookup, + Mutation, +} + pub(in super::super) async fn build_admin_resolve_user_selection_response( state: &AdminAppState<'_>, _request_context: &AdminRequestContext<'_>, @@ -148,10 +160,29 @@ pub(in super::super) async fn build_admin_user_batch_action_response( Ok(value) => value, Err(detail) => return Ok(build_admin_user_batch_bad_request_response(detail)), }; - let mutation = match parse_batch_mutation(&request.action, request.payload) { + let mutation = match parse_batch_mutation(&request.action, request.payload.clone()) { Ok(value) => value, Err(detail) => return Ok(build_admin_user_batch_bad_request_response(detail)), }; + if mutation.wallet_balance_adjustment.is_some() { + if !management_token_may_adjust_admin_wallet_balance(request_context) { + return Ok(build_admin_users_wallet_permission_denied_response( + request_context, + )); + } + if !state.has_auth_wallet_write_capability() { + return Ok(build_admin_users_read_only_response( + "当前为只读模式,无法批量调整用户钱包余额", + )); + } + return build_admin_user_wallet_balance_batch_response( + state, + request_context, + request, + mutation, + ) + .await; + } let resolved = match resolve_admin_user_selection(state, request.selection).await { Ok(value) => value, Err(detail) => return Ok(build_admin_user_batch_bad_request_response(detail)), @@ -180,18 +211,6 @@ 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_wallet_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( - "当前为只读模式,无法批量调整用户钱包余额", - )); - } 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? @@ -263,7 +282,20 @@ pub(in super::super) async fn build_admin_user_batch_action_response( })); continue; } - Err(_) => { + Err(AdminBatchWalletLimitModeError::WalletLookup) => { + record_batch_action_interruption( + &resolved.items, + item_index, + false, + "读取用户钱包失败,批次已中止,该用户未执行", + &mut failures, + &mut uncertain_user_ids, + &mut unprocessed_user_ids, + ); + interrupted = true; + break; + } + Err(AdminBatchWalletLimitModeError::Mutation) => { record_batch_action_interruption( &resolved.items, item_index, @@ -427,6 +459,379 @@ fn record_batch_action_interruption( } } +async fn build_admin_user_wallet_balance_batch_response( + state: &AdminAppState<'_>, + request_context: &AdminRequestContext<'_>, + request: AdminUserBatchActionRequest, + mutation: AdminUserBatchMutation, +) -> Result, GatewayError> { + let Some(idempotency_key) = request + .idempotency_key + .as_deref() + .map(str::trim) + .filter(|value| { + !value.is_empty() + && value.len() <= 128 + && value.bytes().all(|byte| (0x21..=0x7e).contains(&byte)) + }) + else { + return Ok(build_admin_user_batch_bad_request_response( + "余额批量操作必须提供有效的 idempotency_key".to_string(), + )); + }; + let Some(admin_user_id) = request_context + .decision() + .and_then(|decision| decision.admin_principal.as_ref()) + .map(|principal| principal.user_id.clone()) + else { + return Ok(build_admin_users_permission_denied_response( + request_context, + )); + }; + let action = request.action.trim().to_ascii_lowercase(); + let fingerprint_payload = json!({ + "selection": &request.selection, + "action": &action, + "payload": &request.payload, + }); + let encoded = serde_json::to_vec(&fingerprint_payload) + .map_err(|error| GatewayError::Internal(error.to_string()))?; + let request_fingerprint = format!("{:x}", Sha256::digest(encoded)); + + let existing = state + .get_admin_user_wallet_balance_batch(&admin_user_id, idempotency_key, &request_fingerprint) + .await?; + let batch = match existing { + Some(PrepareAdminUserWalletBalanceBatchOutcome::Conflict) => { + return Ok(build_admin_user_batch_idempotency_conflict_response()); + } + Some(PrepareAdminUserWalletBalanceBatchOutcome::Ready(batch)) => batch, + None => { + let resolved = + match resolve_admin_user_selection(state, request.selection.clone()).await { + Ok(value) => value, + Err(detail) => return Ok(build_admin_user_batch_bad_request_response(detail)), + }; + let warnings = serde_json::to_value(&resolved.warnings) + .ok() + .and_then(|value| value.as_array().cloned()) + .unwrap_or_default(); + let prepared = state + .prepare_admin_user_wallet_balance_batch( + aether_data::repository::wallet::PrepareAdminUserWalletBalanceBatchInput { + admin_user_id: admin_user_id.clone(), + idempotency_key: idempotency_key.to_string(), + request_fingerprint: request_fingerprint.clone(), + target_user_ids: resolved + .items + .iter() + .map(|item| item.user_id.clone()) + .collect(), + missing_user_ids: resolved.missing_user_ids, + warnings, + }, + ) + .await?; + match prepared { + PrepareAdminUserWalletBalanceBatchOutcome::Conflict => { + return Ok(build_admin_user_batch_idempotency_conflict_response()); + } + PrepareAdminUserWalletBalanceBatchOutcome::Ready(batch) => batch, + } + } + }; + let warnings: Vec = + serde_json::from_value(Value::Array(batch.warnings.clone())) + .map_err(|error| GatewayError::Internal(error.to_string()))?; + let resolved = ResolvedAdminUserSelection { + items: batch + .target_user_ids + .iter() + .map(|user_id| AdminUserSelectionItem { + user_id: user_id.clone(), + username: String::new(), + email: None, + role: "user".to_string(), + is_active: true, + matched_by: Vec::new(), + }) + .collect(), + missing_user_ids: batch.missing_user_ids.clone(), + warnings, + }; + let adjustment = mutation + .wallet_balance_adjustment + .expect("wallet balance action should have an adjustment"); + let signed_amount = match adjustment.operation { + AdminUserWalletBalanceOperation::Add => adjustment.amount, + AdminUserWalletBalanceOperation::Deduct => -adjustment.amount, + }; + + let mut outcomes = batch.user_outcomes; + let mut completed_user_ids = outcomes + .iter() + .filter_map(|(user_id, outcome)| { + matches!(outcome, AdminUserWalletBalanceBatchUserOutcome::Succeeded) + .then_some(user_id.clone()) + }) + .collect::>(); + let mut success = completed_user_ids.len(); + let mut failures = resolved + .missing_user_ids + .iter() + .map(|user_id| json!({ "user_id": user_id, "reason": "用户不存在或已删除" })) + .collect::>(); + for (user_id, outcome) in &outcomes { + if let AdminUserWalletBalanceBatchUserOutcome::Failed(reason) = outcome { + failures.push(json!({ "user_id": user_id, "reason": reason })); + } + } + let mut uncertain_user_ids = Vec::new(); + let mut unprocessed_user_ids = Vec::new(); + let mut interrupted = false; + + for (item_index, item) in resolved.items.iter().enumerate() { + if outcomes.contains_key(&item.user_id) { + continue; + } + match state.find_user_auth_by_id(&item.user_id).await { + Err(_) => { + failures.push(json!({ + "user_id": item.user_id, + "reason": "读取用户状态失败,批次已中止,该用户未执行", + })); + unprocessed_user_ids.push(item.user_id.clone()); + interrupted = true; + } + Ok(None) => { + let outcome = match state + .record_admin_user_wallet_balance_batch_failure( + &admin_user_id, + idempotency_key, + &item.user_id, + "用户不存在或已删除", + ) + .await + { + Ok(outcome) => outcome, + Err(_) => { + failures.push(json!({ + "user_id": item.user_id, + "reason": "记录用户状态失败,批次已中止,该用户未执行", + })); + unprocessed_user_ids.push(item.user_id.clone()); + interrupted = true; + append_wallet_batch_unprocessed_suffix( + &resolved.items, + item_index + 1, + &outcomes, + &mut failures, + &mut unprocessed_user_ids, + ); + break; + } + }; + outcomes.insert(item.user_id.clone(), outcome.clone()); + match outcome { + AdminUserWalletBalanceBatchUserOutcome::Succeeded => { + completed_user_ids.push(item.user_id.clone()); + success += 1; + } + AdminUserWalletBalanceBatchUserOutcome::Failed(reason) => { + failures.push(json!({ "user_id": item.user_id, "reason": reason })); + } + } + } + Ok(Some(_)) => { + let wallet = match state + .find_wallet(aether_data::repository::wallet::WalletLookupKey::UserId( + &item.user_id, + )) + .await + { + Ok(Some(wallet)) => wallet, + Ok(None) => { + let outcome = match state + .record_admin_user_wallet_balance_batch_failure( + &admin_user_id, + idempotency_key, + &item.user_id, + "用户钱包不可用", + ) + .await + { + Ok(outcome) => outcome, + Err(_) => { + failures.push(json!({ + "user_id": item.user_id, + "reason": "记录用户钱包状态失败,批次已中止,该用户未执行", + })); + unprocessed_user_ids.push(item.user_id.clone()); + interrupted = true; + append_wallet_batch_unprocessed_suffix( + &resolved.items, + item_index + 1, + &outcomes, + &mut failures, + &mut unprocessed_user_ids, + ); + break; + } + }; + outcomes.insert(item.user_id.clone(), outcome.clone()); + match outcome { + AdminUserWalletBalanceBatchUserOutcome::Succeeded => { + completed_user_ids.push(item.user_id.clone()); + success += 1; + } + AdminUserWalletBalanceBatchUserOutcome::Failed(reason) => { + failures.push(json!({ "user_id": item.user_id, "reason": reason })); + } + } + continue; + } + Err(_) => { + failures.push(json!({ + "user_id": item.user_id, + "reason": "读取用户钱包失败,批次已中止,该用户未执行", + })); + unprocessed_user_ids.push(item.user_id.clone()); + interrupted = true; + append_wallet_batch_unprocessed_suffix( + &resolved.items, + item_index + 1, + &outcomes, + &mut failures, + &mut unprocessed_user_ids, + ); + break; + } + }; + let result = state + .adjust_admin_user_wallet_balance_batch_user( + aether_data::repository::wallet::AdjustWalletBalanceInBatchInput { + admin_user_id: admin_user_id.clone(), + idempotency_key: idempotency_key.to_string(), + user_id: item.user_id.clone(), + adjustment: aether_data::repository::wallet::AdjustWalletBalanceInput { + wallet_id: wallet.id, + amount_usd: signed_amount, + balance_type: "recharge".to_string(), + operator_id: Some(admin_user_id.clone()), + description: Some("管理员批量调整用户余额".to_string()), + clamp_deduction_to_available_balance: true, + batch_context: None, + }, + }, + ) + .await; + match result { + Ok(AdminUserWalletBalanceBatchUserOutcome::Succeeded) => { + outcomes.insert( + item.user_id.clone(), + AdminUserWalletBalanceBatchUserOutcome::Succeeded, + ); + completed_user_ids.push(item.user_id.clone()); + success += 1; + } + Ok(AdminUserWalletBalanceBatchUserOutcome::Failed(reason)) => { + outcomes.insert( + item.user_id.clone(), + AdminUserWalletBalanceBatchUserOutcome::Failed(reason.clone()), + ); + failures.push(json!({ "user_id": item.user_id, "reason": reason })); + } + Err(_) => { + failures.push(json!({ + "user_id": item.user_id, + "reason": "余额调整结果未确认,批次已中止,请使用同一批次重试以核对结果", + })); + uncertain_user_ids.push(item.user_id.clone()); + interrupted = true; + append_wallet_batch_unprocessed_suffix( + &resolved.items, + item_index + 1, + &outcomes, + &mut failures, + &mut unprocessed_user_ids, + ); + break; + } + } + } + } + if interrupted { + for pending in resolved.items.iter().skip(item_index + 1) { + if outcomes.contains_key(&pending.user_id) { + continue; + } + failures.push(json!({ + "user_id": pending.user_id, + "reason": "因前序错误未执行", + })); + unprocessed_user_ids.push(pending.user_id.clone()); + } + break; + } + } + + let failed = failures.len(); + let total = success + failed; + let mut response_payload = json!({ + "total": total, + "success": success, + "failed": failed, + "failures": failures, + "warnings": resolved.warnings, + "action": action, + "modified_fields": mutation.modified_fields, + "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, + "admin_users_batch_action_executed", + "batch_update_users", + "user_batch", + "users", + )) +} + +fn append_wallet_batch_unprocessed_suffix( + items: &[AdminUserSelectionItem], + start_index: usize, + outcomes: &BTreeMap, + failures: &mut Vec, + unprocessed_user_ids: &mut Vec, +) { + for pending in items.iter().skip(start_index) { + if outcomes.contains_key(&pending.user_id) { + continue; + } + failures.push(json!({ + "user_id": pending.user_id, + "reason": "因前序错误未执行", + })); + unprocessed_user_ids.push(pending.user_id.clone()); + } +} + +fn build_admin_user_batch_idempotency_conflict_response() -> Response { + ( + http::StatusCode::CONFLICT, + Json(json!({ + "detail": "idempotency_key was already used with a different request", + "error_code": "idempotency_key_conflict", + })), + ) + .into_response() +} + fn parse_resolve_selection_request( request_body: Option<&Bytes>, ) -> Result { @@ -452,6 +857,7 @@ fn parse_batch_action_request( selection: parse_selection_request_value(raw.selection)?, action: raw.action, payload: raw.payload, + idempotency_key: raw.idempotency_key, }) } _ => Err("Invalid JSON request body".to_string()), @@ -893,26 +1299,29 @@ async fn apply_batch_user_wallet_limit_mode( state: &AdminAppState<'_>, user_id: &str, unlimited: bool, -) -> Result { +) -> Result { let desired_limit_mode = if unlimited { "unlimited" } else { "finite" }; - match state + let wallet = state .find_wallet(aether_data::repository::wallet::WalletLookupKey::UserId( user_id, )) - .await? - { + .await + .map_err(|_| AdminBatchWalletLimitModeError::WalletLookup)?; + match wallet { Some(wallet) => { if wallet.limit_mode.eq_ignore_ascii_case(desired_limit_mode) { return Ok(true); } Ok(state .update_auth_user_wallet_limit_mode(user_id, desired_limit_mode) - .await? + .await + .map_err(|_| AdminBatchWalletLimitModeError::Mutation)? .is_some()) } None => Ok(state .initialize_auth_user_wallet(user_id, 0.0, unlimited) - .await? + .await + .map_err(|_| AdminBatchWalletLimitModeError::Mutation)? .is_some()), } } diff --git a/apps/aether-gateway/src/state/app.rs b/apps/aether-gateway/src/state/app.rs index a3996a1ab..962c003fc 100644 --- a/apps/aether-gateway/src/state/app.rs +++ b/apps/aether-gateway/src/state/app.rs @@ -496,6 +496,21 @@ pub struct AppState { #[cfg(test)] pub(crate) auth_wallet_lookup_error_for_tests: Option, #[cfg(test)] + pub(crate) auth_wallet_batch_store_for_tests: Option< + Arc< + StdMutex< + HashMap< + (String, String), + aether_data::repository::wallet::StoredAdminUserWalletBalanceBatch, + >, + >, + >, + >, + #[cfg(test)] + pub(crate) auth_wallet_batch_operation_lock_for_tests: Arc>, + #[cfg(test)] + pub(crate) auth_wallet_batch_failure_record_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 8acfaf40c..f8a46c47e 100644 --- a/apps/aether-gateway/src/state/core.rs +++ b/apps/aether-gateway/src/state/core.rs @@ -470,6 +470,12 @@ impl AppState { #[cfg(test)] auth_wallet_lookup_error_for_tests: None, #[cfg(test)] + auth_wallet_batch_store_for_tests: Some(Arc::new(StdMutex::new(HashMap::new()))), + #[cfg(test)] + auth_wallet_batch_operation_lock_for_tests: Arc::new(TokioMutex::new(())), + #[cfg(test)] + auth_wallet_batch_failure_record_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 212233803..70a9ee1eb 100644 --- a/apps/aether-gateway/src/state/runtime/wallet/balance_mutations.rs +++ b/apps/aether-gateway/src/state/runtime/wallet/balance_mutations.rs @@ -1,8 +1,198 @@ use crate::{AdminWalletPaymentOrderRecord, AdminWalletTransactionRecord, AppState, GatewayError}; +use aether_data::repository::wallet::{ + AdjustWalletBalanceInBatchInput, AdminUserWalletBalanceBatchUserOutcome, + PrepareAdminUserWalletBalanceBatchInput, PrepareAdminUserWalletBalanceBatchOutcome, + StoredAdminUserWalletBalanceBatch, +}; +use std::collections::BTreeMap; use super::admin_wallet_build_order_no; impl AppState { + pub(crate) async fn prepare_admin_user_wallet_balance_batch( + &self, + input: PrepareAdminUserWalletBalanceBatchInput, + ) -> Result { + #[cfg(test)] + if let Some(store) = self.auth_wallet_batch_store_for_tests.as_ref() { + let mut batches = store.lock().expect("auth wallet batch store should lock"); + let key = (input.admin_user_id.clone(), input.idempotency_key.clone()); + if let Some(existing) = batches.get(&key) { + if existing.request_fingerprint != input.request_fingerprint { + return Ok(PrepareAdminUserWalletBalanceBatchOutcome::Conflict); + } + return Ok(PrepareAdminUserWalletBalanceBatchOutcome::Ready( + existing.clone(), + )); + } + let batch = StoredAdminUserWalletBalanceBatch { + admin_user_id: input.admin_user_id, + idempotency_key: input.idempotency_key, + request_fingerprint: input.request_fingerprint, + target_user_ids: input.target_user_ids, + missing_user_ids: input.missing_user_ids, + warnings: input.warnings, + user_outcomes: BTreeMap::new(), + }; + batches.insert(key, batch.clone()); + return Ok(PrepareAdminUserWalletBalanceBatchOutcome::Ready(batch)); + } + + self.data + .prepare_admin_user_wallet_balance_batch(input) + .await + .map_err(|error| GatewayError::Internal(error.to_string()))? + .ok_or_else(|| { + GatewayError::Internal("admin wallet batch storage is unavailable".to_string()) + }) + } + + pub(crate) async fn get_admin_user_wallet_balance_batch( + &self, + admin_user_id: &str, + idempotency_key: &str, + request_fingerprint: &str, + ) -> Result, GatewayError> { + #[cfg(test)] + if let Some(store) = self.auth_wallet_batch_store_for_tests.as_ref() { + let batches = store.lock().expect("auth wallet batch store should lock"); + return Ok(batches + .get(&(admin_user_id.to_string(), idempotency_key.to_string())) + .map(|existing| { + if existing.request_fingerprint != request_fingerprint { + PrepareAdminUserWalletBalanceBatchOutcome::Conflict + } else { + PrepareAdminUserWalletBalanceBatchOutcome::Ready(existing.clone()) + } + })); + } + + self.data + .get_admin_user_wallet_balance_batch( + admin_user_id, + idempotency_key, + request_fingerprint, + ) + .await + .map_err(|error| GatewayError::Internal(error.to_string())) + } + + pub(crate) async fn adjust_admin_user_wallet_balance_batch_user( + &self, + input: AdjustWalletBalanceInBatchInput, + ) -> Result { + #[cfg(test)] + if let Some(store) = self.auth_wallet_batch_store_for_tests.as_ref() { + let _operation = self.auth_wallet_batch_operation_lock_for_tests.lock().await; + let key = (input.admin_user_id.clone(), input.idempotency_key.clone()); + if let Some(existing) = store + .lock() + .expect("auth wallet batch store should lock") + .get(&key) + .and_then(|batch| batch.user_outcomes.get(&input.user_id)) + .cloned() + { + return Ok(existing); + } + let target_exists = store + .lock() + .expect("auth wallet batch store should lock") + .get(&key) + .is_some_and(|batch| batch.target_user_ids.contains(&input.user_id)); + if !target_exists { + return Err(GatewayError::Internal( + "user is outside the prepared admin wallet batch".to_string(), + )); + } + let adjustment = input.adjustment; + let result = self + .admin_adjust_wallet_balance( + &adjustment.wallet_id, + adjustment.amount_usd, + &adjustment.balance_type, + adjustment.operator_id.as_deref(), + adjustment.description.as_deref(), + adjustment.clamp_deduction_to_available_balance, + ) + .await?; + let outcome = if result.is_some() { + AdminUserWalletBalanceBatchUserOutcome::Succeeded + } else { + AdminUserWalletBalanceBatchUserOutcome::Failed("用户钱包不可用".to_string()) + }; + if let Some(batch) = store + .lock() + .expect("auth wallet batch store should lock") + .get_mut(&key) + { + batch.user_outcomes.insert(input.user_id, outcome.clone()); + } + return Ok(outcome); + } + + self.data + .adjust_admin_user_wallet_balance_batch_user(input) + .await + .map_err(|error| GatewayError::Internal(error.to_string()))? + .ok_or_else(|| { + GatewayError::Internal("admin wallet batch storage is unavailable".to_string()) + }) + } + + pub(crate) async fn record_admin_user_wallet_balance_batch_failure( + &self, + admin_user_id: &str, + idempotency_key: &str, + user_id: &str, + reason: &str, + ) -> Result { + #[cfg(test)] + if self + .auth_wallet_batch_failure_record_error_for_tests + .as_deref() + == Some(user_id) + { + return Err(GatewayError::Internal( + "injected wallet batch failure-record error".to_string(), + )); + } + + #[cfg(test)] + if let Some(store) = self.auth_wallet_batch_store_for_tests.as_ref() { + let _operation = self.auth_wallet_batch_operation_lock_for_tests.lock().await; + let key = (admin_user_id.to_string(), idempotency_key.to_string()); + let mut batches = store.lock().expect("auth wallet batch store should lock"); + let batch = batches.get_mut(&key).ok_or_else(|| { + GatewayError::Internal("admin wallet batch was not prepared".to_string()) + })?; + if !batch.target_user_ids.iter().any(|target| target == user_id) { + return Err(GatewayError::Internal( + "user is outside the prepared admin wallet batch".to_string(), + )); + } + return Ok(batch + .user_outcomes + .entry(user_id.to_string()) + .or_insert_with(|| { + AdminUserWalletBalanceBatchUserOutcome::Failed(reason.to_string()) + }) + .clone()); + } + + self.data + .record_admin_user_wallet_balance_batch_failure( + admin_user_id, + idempotency_key, + user_id, + reason, + ) + .await + .map_err(|error| GatewayError::Internal(error.to_string()))? + .ok_or_else(|| { + GatewayError::Internal("admin wallet batch storage is unavailable".to_string()) + }) + } + pub(crate) async fn admin_adjust_wallet_balance( &self, wallet_id: &str, @@ -117,6 +307,7 @@ impl AppState { operator_id: operator_id.map(ToOwned::to_owned), description: description.map(ToOwned::to_owned), clamp_deduction_to_available_balance, + batch_context: None, }) .await? .map(|(wallet, transaction)| { diff --git a/apps/aether-gateway/src/state/testing.rs b/apps/aether-gateway/src/state/testing.rs index 30cb1c079..ebc27b469 100644 --- a/apps/aether-gateway/src/state/testing.rs +++ b/apps/aether-gateway/src/state/testing.rs @@ -485,6 +485,14 @@ impl AppState { self } + pub(crate) fn fail_auth_wallet_batch_failure_record_for_tests( + mut self, + user_id: impl Into, + ) -> Self { + self.auth_wallet_batch_failure_record_error_for_tests = Some(user_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 bb6fa0308..bb34d61c8 100644 --- a/apps/aether-gateway/src/tests/control/admin/users_batch.rs +++ b/apps/aether-gateway/src/tests/control/admin/users_batch.rs @@ -1,3 +1,4 @@ +use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::Arc; use aether_data::repository::management_tokens::InMemoryManagementTokenRepository; @@ -68,6 +69,14 @@ fn sample_wallet(user_id: &str, balance: f64, gift_balance: f64) -> StoredWallet } async fn post_batch_action(client: &Client, gateway_url: &str, payload: Value) -> Response { + static NEXT_TEST_IDEMPOTENCY_KEY: AtomicU64 = AtomicU64::new(1); + let mut payload = payload; + if payload.get("action").and_then(Value::as_str) == Some("adjust_wallet_balance") + && payload.get("idempotency_key").is_none() + { + let sequence = NEXT_TEST_IDEMPOTENCY_KEY.fetch_add(1, Ordering::Relaxed); + payload["idempotency_key"] = json!(format!("test-wallet-batch-{sequence}")); + } admin_headers(client.post(format!("{gateway_url}/api/admin/users/batch-action"))) .json(&payload) .send() @@ -172,6 +181,132 @@ async fn gateway_batches_wallet_addition_deduction_and_clamped_deduction_per_use gateway_handle.abort(); } +#[tokio::test] +async fn gateway_replays_wallet_batch_idempotently_and_rejects_key_reuse() { + 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 request = json!({ + "selection": { "user_ids": ["user-1"] }, + "action": "adjust_wallet_balance", + "payload": { "operation": "add", "amount": 5.0 }, + "idempotency_key": "same-wallet-batch" + }); + + let first = post_batch_action(&client, &gateway_url, request.clone()).await; + assert_eq!(first.status(), StatusCode::OK); + let first_result: Value = first.json().await.expect("response should parse"); + assert_eq!(first_result["success"], 1); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 15.0 + ); + + let replay = post_batch_action(&client, &gateway_url, request.clone()).await; + assert_eq!(replay.status(), StatusCode::OK); + let replay_result: Value = replay.json().await.expect("response should parse"); + assert_eq!(replay_result["success"], 1); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 15.0 + ); + + let changed_request = json!({ + "selection": { "user_ids": ["user-1"] }, + "action": "adjust_wallet_balance", + "payload": { "operation": "add", "amount": 50.0 }, + "idempotency_key": "same-wallet-batch" + }); + let conflict = post_batch_action(&client, &gateway_url, changed_request).await; + assert_eq!(conflict.status(), StatusCode::CONFLICT); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 15.0 + ); + + gateway_handle.abort(); +} + +#[tokio::test] +async fn gateway_returns_partial_results_when_failure_recording_fails() { + 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, 0.0)]) + .fail_auth_wallet_batch_failure_record_for_tests("user-2"); + let (gateway_url, gateway_handle) = start_server(build_router_with_state(state)).await; + let client = Client::new(); + let request = json!({ + "selection": { "user_ids": ["user-1", "user-2"] }, + "action": "adjust_wallet_balance", + "payload": { "operation": "add", "amount": 5.0 }, + "idempotency_key": "failure-record-wallet-batch" + }); + + let response = post_batch_action(&client, &gateway_url, request).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["interrupted"], true); + assert_eq!(result["completed_user_ids"], json!(["user-1"])); + assert_eq!(result["uncertain_user_ids"], json!([])); + assert_eq!(result["unprocessed_user_ids"], json!(["user-2"])); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 15.0 + ); + + gateway_handle.abort(); +} + +#[tokio::test] +async fn gateway_records_zero_delta_batch_for_later_replay() { + 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", 0.0, 0.0)]); + let (gateway_url, gateway_handle) = start_server(build_router_with_state(state)).await; + let client = Client::new(); + let deduction = json!({ + "selection": { "user_ids": ["user-1"] }, + "action": "adjust_wallet_balance", + "payload": { "operation": "deduct", "amount": 5.0 }, + "idempotency_key": "zero-delta-wallet-batch" + }); + let first = post_batch_action(&client, &gateway_url, deduction.clone()).await; + assert_eq!(first.status(), StatusCode::OK); + + let top_up = post_batch_action( + &client, + &gateway_url, + json!({ + "selection": { "user_ids": ["user-1"] }, + "action": "adjust_wallet_balance", + "payload": { "operation": "add", "amount": 10.0 }, + "idempotency_key": "wallet-top-up-batch" + }), + ) + .await; + assert_eq!(top_up.status(), StatusCode::OK); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 10.0 + ); + + let replay = post_batch_action(&client, &gateway_url, deduction).await; + assert_eq!(replay.status(), StatusCode::OK); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 10.0 + ); + + 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"; @@ -236,7 +371,8 @@ async fn gateway_requires_wallet_write_permission_for_batch_balance_adjustments( let payload = json!({ "selection": { "user_ids": ["user-1"] }, "action": "adjust_wallet_balance", - "payload": { "operation": "add", "amount": 5.0 } + "payload": { "operation": "add", "amount": 5.0 }, + "idempotency_key": "wallet-write-batch" }); let denied = client @@ -279,7 +415,12 @@ async fn gateway_requires_wallet_write_permission_for_batch_balance_adjustments( .post(format!("{gateway_url}/api/admin/users/batch-action")) .header(crate::constants::GATEWAY_HEADER, "rust-phase3b") .bearer_auth(wallet_admin_token) - .json(&payload) + .json(&json!({ + "selection": payload["selection"], + "action": payload["action"], + "payload": payload["payload"], + "idempotency_key": "wallet-admin-batch" + })) .send() .await .expect("wallet-admin management token request should complete"); @@ -401,6 +542,44 @@ async fn gateway_reports_wallet_lookup_failure_as_unprocessed() { gateway_handle.abort(); } +#[tokio::test] +async fn gateway_reports_wallet_limit_lookup_failure_as_unprocessed() { + 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_lookup_for_tests("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": "update_access_control", + "payload": { "unlimited": true } + }), + ) + .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["completed_user_ids"], json!(["user-1"])); + assert_eq!(result["uncertain_user_ids"], json!([])); + assert_eq!(result["unprocessed_user_ids"], json!(["user-2", "user-3"])); + + 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/crates/aether-data/adapters/postgres/migrations/20260923000000_add_admin_wallet_batch_idempotency.sql b/crates/aether-data/adapters/postgres/migrations/20260923000000_add_admin_wallet_batch_idempotency.sql new file mode 100644 index 000000000..edd5cfdc8 --- /dev/null +++ b/crates/aether-data/adapters/postgres/migrations/20260923000000_add_admin_wallet_batch_idempotency.sql @@ -0,0 +1,12 @@ +CREATE TABLE IF NOT EXISTS public.admin_user_wallet_balance_batches ( + admin_user_id character varying(64) NOT NULL, + idempotency_key character varying(128) NOT NULL, + request_fingerprint character varying(64) NOT NULL, + target_user_ids jsonb NOT NULL, + missing_user_ids jsonb NOT NULL, + warnings jsonb NOT NULL, + user_outcomes jsonb NOT NULL, + created_at_unix_secs bigint NOT NULL, + updated_at_unix_secs bigint NOT NULL, + CONSTRAINT admin_user_wallet_balance_batches_pkey PRIMARY KEY (admin_user_id, idempotency_key) +); diff --git a/crates/aether-data/adapters/postgres/src/wallet.rs b/crates/aether-data/adapters/postgres/src/wallet.rs index 4d26e8b35..23f524d54 100644 --- a/crates/aether-data/adapters/postgres/src/wallet.rs +++ b/crates/aether-data/adapters/postgres/src/wallet.rs @@ -24,8 +24,9 @@ use aether_data_contracts::repository::wallet::{ wallet_recharge_checkout_failed_response, wallet_recharge_checkout_uncertain_response, wallet_recharge_order_is_reclaimable_placeholder, wallet_recharge_replay_matches, wallet_recharge_response_is_checkout_placeholder, wallet_refund_proof_is_success, - AdjustWalletBalanceInput, AdminPaymentOrderListQuery, AdminRedeemCodeBatchListQuery, - AdminRedeemCodeListQuery, AdminWalletLedgerQuery, AdminWalletListQuery, + AdjustWalletBalanceInBatchInput, AdjustWalletBalanceInput, AdminPaymentOrderListQuery, + AdminRedeemCodeBatchListQuery, AdminRedeemCodeListQuery, AdminUserWalletBalanceBatchContext, + AdminUserWalletBalanceBatchUserOutcome, AdminWalletLedgerQuery, AdminWalletListQuery, AdminWalletRefundRequestListQuery, CompareAndSwapPaymentOrderStripeClientSecretInput, CompleteAdminWalletRefundInput, CreateAdminRedeemCodeBatchInput, CreateAdminRedeemCodeBatchResult, CreateManualWalletRechargeInput, @@ -34,21 +35,23 @@ use aether_data_contracts::repository::wallet::{ CreateWalletRefundRequestOutcome, CreatedAdminRedeemCodePlaintext, CreditAdminPaymentOrderInput, DeleteAdminRedeemCodeBatchInput, DisableAdminRedeemCodeBatchInput, DisableAdminRedeemCodeInput, FailAdminWalletRefundInput, - FailWalletRechargeCheckoutInput, InitializeAuthWalletOutcome, ProcessAdminWalletRefundInput, - ProcessPaymentCallbackInput, ProcessPaymentCallbackOutcome, ReclaimWalletRechargeCheckoutInput, - RedeemWalletCodeInput, RedeemWalletCodeOutcome, StoredAdminPaymentCallback, - StoredAdminPaymentCallbackPage, StoredAdminPaymentOrder, StoredAdminPaymentOrderPage, - StoredAdminRedeemCode, StoredAdminRedeemCodeBatch, StoredAdminRedeemCodeBatchPage, - StoredAdminRedeemCodePage, StoredAdminWalletLedgerItem, StoredAdminWalletLedgerPage, - StoredAdminWalletListItem, StoredAdminWalletListPage, StoredAdminWalletRefund, - StoredAdminWalletRefundPage, StoredAdminWalletRefundRequestItem, - StoredAdminWalletRefundRequestPage, StoredAdminWalletTransaction, - StoredAdminWalletTransactionPage, StoredWalletDailyUsageLedger, + FailWalletRechargeCheckoutInput, InitializeAuthWalletOutcome, + PrepareAdminUserWalletBalanceBatchInput, PrepareAdminUserWalletBalanceBatchOutcome, + ProcessAdminWalletRefundInput, ProcessPaymentCallbackInput, ProcessPaymentCallbackOutcome, + ReclaimWalletRechargeCheckoutInput, RedeemWalletCodeInput, RedeemWalletCodeOutcome, + StoredAdminPaymentCallback, StoredAdminPaymentCallbackPage, StoredAdminPaymentOrder, + StoredAdminPaymentOrderPage, StoredAdminRedeemCode, StoredAdminRedeemCodeBatch, + StoredAdminRedeemCodeBatchPage, StoredAdminRedeemCodePage, StoredAdminUserWalletBalanceBatch, + StoredAdminWalletLedgerItem, StoredAdminWalletLedgerPage, StoredAdminWalletListItem, + StoredAdminWalletListPage, StoredAdminWalletRefund, StoredAdminWalletRefundPage, + StoredAdminWalletRefundRequestItem, StoredAdminWalletRefundRequestPage, + StoredAdminWalletTransaction, StoredAdminWalletTransactionPage, StoredWalletDailyUsageLedger, StoredWalletDailyUsageLedgerPage, StoredWalletSnapshot, UpdateAdminWalletRefundGatewayInput, UpdateWalletRechargeCheckoutInput, WalletLookupKey, WalletMutationOutcome, WalletReadRepository, WalletWriteRepository, }; use aether_data_contracts::DataLayerError; +use std::collections::BTreeMap; use crate::{ error::{postgres_error, SqlxResultExt}, @@ -815,6 +818,31 @@ fn effective_wallet_adjustment_amount(input: &AdjustWalletBalanceInput, before_t } } +async fn persist_admin_wallet_batch_outcomes( + connection: &mut sqlx::PgConnection, + context: &AdminUserWalletBalanceBatchContext, + outcomes: &BTreeMap, +) -> Result<(), DataLayerError> { + let outcomes = serde_json::to_value(outcomes).map_err(|error| { + DataLayerError::UnexpectedValue(format!("admin wallet batch outcomes are invalid: {error}")) + })?; + sqlx::query( + r#" +UPDATE admin_user_wallet_balance_batches +SET user_outcomes = $3, updated_at_unix_secs = $4 +WHERE admin_user_id = $1 AND idempotency_key = $2 + "#, + ) + .bind(&context.admin_user_id) + .bind(&context.idempotency_key) + .bind(outcomes) + .bind(Utc::now().timestamp().max(0)) + .execute(connection) + .await + .map_postgres_err()?; + Ok(()) +} + #[async_trait] impl WalletReadRepository for SqlxWalletRepository { async fn find( @@ -4280,6 +4308,58 @@ RETURNING self.tx_runner .run_read_write(|tx| { Box::pin(async move { + let batch_context = input.batch_context.clone(); + let mut batch_user_outcomes = if let Some(context) = &batch_context { + let batch_row = sqlx::query( + r#" +SELECT target_user_ids, user_outcomes +FROM admin_user_wallet_balance_batches +WHERE admin_user_id = $1 AND idempotency_key = $2 +FOR UPDATE + "#, + ) + .bind(&context.admin_user_id) + .bind(&context.idempotency_key) + .fetch_optional(&mut **tx) + .await + .map_postgres_err()? + .ok_or_else(|| { + DataLayerError::InvalidInput( + "admin wallet batch was not prepared".to_string(), + ) + })?; + let target_user_ids: Vec = + serde_json::from_value(row_get(&batch_row, "target_user_ids")?) + .map_err(|error| { + DataLayerError::UnexpectedValue(format!( + "admin wallet batch target list is invalid: {error}" + )) + })?; + if !target_user_ids.iter().any(|id| id == &context.user_id) { + return Err(DataLayerError::InvalidInput( + "user is outside the prepared admin wallet batch".to_string(), + )); + } + let outcomes: BTreeMap = + serde_json::from_value(row_get(&batch_row, "user_outcomes")?).map_err( + |error| { + DataLayerError::UnexpectedValue(format!( + "admin wallet batch outcomes are invalid: {error}" + )) + }, + )?; + if let Some(outcome) = outcomes.get(&context.user_id) { + match outcome { + AdminUserWalletBalanceBatchUserOutcome::Succeeded => {} + AdminUserWalletBalanceBatchUserOutcome::Failed(_) => { + return Ok(None); + } + } + } + Some(outcomes) + } else { + None + }; let Some(row) = sqlx::query( r#" SELECT @@ -4298,10 +4378,16 @@ SELECT CAST(EXTRACT(EPOCH FROM updated_at) AS BIGINT) AS updated_at_unix_secs FROM wallets WHERE id = $1 + AND ($2::character varying IS NULL OR user_id = $2::character varying) FOR UPDATE "#, ) .bind(&input.wallet_id) + .bind( + batch_context + .as_ref() + .map(|context| context.user_id.as_str()), + ) .fetch_optional(&mut **tx) .await .map_postgres_err()? @@ -4322,9 +4408,28 @@ FOR UPDATE "wallet balance is invalid".to_string(), )); } + let wallet = map_wallet_row(&row)?; + let already_applied = batch_context.as_ref().is_some_and(|context| { + batch_user_outcomes.as_ref().is_some_and(|outcomes| { + outcomes.get(&context.user_id) + == Some(&AdminUserWalletBalanceBatchUserOutcome::Succeeded) + }) + }); + if already_applied { + return Ok(Some((wallet, None))); + } let amount_usd = effective_wallet_adjustment_amount(&input, before_total); if amount_usd == 0.0 { - return Ok(Some((map_wallet_row(&row)?, None))); + if let (Some(context), Some(outcomes)) = + (batch_context.as_ref(), batch_user_outcomes.as_mut()) + { + outcomes.insert( + context.user_id.clone(), + AdminUserWalletBalanceBatchUserOutcome::Succeeded, + ); + persist_admin_wallet_batch_outcomes(tx, context, outcomes).await?; + } + return Ok(Some((wallet, None))); } let mut after_recharge = before_recharge; let mut after_gift = before_gift; @@ -4467,6 +4572,16 @@ VALUES ( .await .map_postgres_err()?; + if let (Some(context), Some(outcomes)) = + (batch_context.as_ref(), batch_user_outcomes.as_mut()) + { + outcomes.insert( + context.user_id.clone(), + AdminUserWalletBalanceBatchUserOutcome::Succeeded, + ); + persist_admin_wallet_batch_outcomes(tx, context, outcomes).await?; + } + Ok(Some(( wallet, Some(StoredAdminWalletTransaction { @@ -4495,6 +4610,237 @@ VALUES ( .await } + async fn prepare_admin_user_wallet_balance_batch( + &self, + input: PrepareAdminUserWalletBalanceBatchInput, + ) -> Result { + let target_user_ids = serde_json::to_value(&input.target_user_ids).map_err(|error| { + DataLayerError::InvalidInput(format!("invalid batch target users: {error}")) + })?; + let missing_user_ids = serde_json::to_value(&input.missing_user_ids).map_err(|error| { + DataLayerError::InvalidInput(format!("invalid missing batch users: {error}")) + })?; + let warnings = serde_json::to_value(&input.warnings).map_err(|error| { + DataLayerError::InvalidInput(format!("invalid batch warnings: {error}")) + })?; + let now = Utc::now().timestamp().max(0); + self.tx_runner + .run_read_write(|tx| { + Box::pin(async move { + sqlx::query( + r#" +INSERT INTO admin_user_wallet_balance_batches ( + admin_user_id, idempotency_key, request_fingerprint, target_user_ids, + missing_user_ids, warnings, user_outcomes, created_at_unix_secs, updated_at_unix_secs +) +VALUES ($1, $2, $3, $4, $5, $6, '{}'::jsonb, $7, $7) +ON CONFLICT (admin_user_id, idempotency_key) DO NOTHING + "#, + ) + .bind(&input.admin_user_id) + .bind(&input.idempotency_key) + .bind(&input.request_fingerprint) + .bind(target_user_ids) + .bind(missing_user_ids) + .bind(warnings) + .bind(now) + .execute(&mut **tx) + .await + .map_postgres_err()?; + + let row = sqlx::query( + r#" +SELECT admin_user_id, idempotency_key, request_fingerprint, target_user_ids, + missing_user_ids, warnings, user_outcomes +FROM admin_user_wallet_balance_batches +WHERE admin_user_id = $1 AND idempotency_key = $2 +FOR UPDATE + "#, + ) + .bind(&input.admin_user_id) + .bind(&input.idempotency_key) + .fetch_one(&mut **tx) + .await + .map_postgres_err()?; + let request_fingerprint: String = row_get(&row, "request_fingerprint")?; + if request_fingerprint != input.request_fingerprint { + return Ok(PrepareAdminUserWalletBalanceBatchOutcome::Conflict); + } + + let read_json = |column: &str| -> Result { + row_get(&row, column) + }; + let batch = StoredAdminUserWalletBalanceBatch { + admin_user_id: row_get(&row, "admin_user_id")?, + idempotency_key: row_get(&row, "idempotency_key")?, + request_fingerprint, + target_user_ids: serde_json::from_value(read_json("target_user_ids")?) + .map_err(|error| DataLayerError::UnexpectedValue(error.to_string()))?, + missing_user_ids: serde_json::from_value(read_json("missing_user_ids")?) + .map_err(|error| DataLayerError::UnexpectedValue(error.to_string()))?, + warnings: serde_json::from_value(read_json("warnings")?) + .map_err(|error| DataLayerError::UnexpectedValue(error.to_string()))?, + user_outcomes: serde_json::from_value(read_json("user_outcomes")?) + .map_err(|error| DataLayerError::UnexpectedValue(error.to_string()))?, + }; + Ok(PrepareAdminUserWalletBalanceBatchOutcome::Ready(batch)) + }) + }) + .await + } + + async fn get_admin_user_wallet_balance_batch( + &self, + admin_user_id: &str, + idempotency_key: &str, + expected_fingerprint: &str, + ) -> Result, DataLayerError> { + let row = sqlx::query( + r#" +SELECT admin_user_id, idempotency_key, request_fingerprint, target_user_ids, + missing_user_ids, warnings, user_outcomes +FROM admin_user_wallet_balance_batches +WHERE admin_user_id = $1 AND idempotency_key = $2 + "#, + ) + .bind(admin_user_id) + .bind(idempotency_key) + .fetch_optional(&self.pool) + .await + .map_postgres_err()?; + let Some(row) = row else { + return Ok(None); + }; + let request_fingerprint: String = row_get(&row, "request_fingerprint")?; + if request_fingerprint != expected_fingerprint { + return Ok(Some(PrepareAdminUserWalletBalanceBatchOutcome::Conflict)); + } + let read_json = + |column: &str| -> Result { row_get(&row, column) }; + let batch = StoredAdminUserWalletBalanceBatch { + admin_user_id: row_get(&row, "admin_user_id")?, + idempotency_key: row_get(&row, "idempotency_key")?, + request_fingerprint, + target_user_ids: serde_json::from_value(read_json("target_user_ids")?) + .map_err(|error| DataLayerError::UnexpectedValue(error.to_string()))?, + missing_user_ids: serde_json::from_value(read_json("missing_user_ids")?) + .map_err(|error| DataLayerError::UnexpectedValue(error.to_string()))?, + warnings: serde_json::from_value(read_json("warnings")?) + .map_err(|error| DataLayerError::UnexpectedValue(error.to_string()))?, + user_outcomes: serde_json::from_value(read_json("user_outcomes")?) + .map_err(|error| DataLayerError::UnexpectedValue(error.to_string()))?, + }; + Ok(Some(PrepareAdminUserWalletBalanceBatchOutcome::Ready( + batch, + ))) + } + + async fn adjust_admin_user_wallet_balance_batch_user( + &self, + input: AdjustWalletBalanceInBatchInput, + ) -> Result { + let context = AdminUserWalletBalanceBatchContext { + admin_user_id: input.admin_user_id.clone(), + idempotency_key: input.idempotency_key.clone(), + user_id: input.user_id.clone(), + }; + let mut adjustment = input.adjustment; + adjustment.batch_context = Some(context); + if self.adjust_wallet_balance(adjustment).await?.is_some() { + return Ok(AdminUserWalletBalanceBatchUserOutcome::Succeeded); + } + self.record_admin_user_wallet_balance_batch_failure( + &input.admin_user_id, + &input.idempotency_key, + &input.user_id, + "用户钱包不可用", + ) + .await + } + + async fn record_admin_user_wallet_balance_batch_failure( + &self, + admin_user_id: &str, + idempotency_key: &str, + user_id: &str, + reason: &str, + ) -> Result { + let admin_user_id = admin_user_id.to_string(); + let idempotency_key = idempotency_key.to_string(); + let user_id = user_id.to_string(); + let reason = reason.to_string(); + self.tx_runner + .run_read_write(|tx| { + Box::pin(async move { + let row = sqlx::query( + r#" +SELECT target_user_ids, user_outcomes +FROM admin_user_wallet_balance_batches +WHERE admin_user_id = $1 AND idempotency_key = $2 +FOR UPDATE + "#, + ) + .bind(&admin_user_id) + .bind(&idempotency_key) + .fetch_optional(&mut **tx) + .await + .map_postgres_err()? + .ok_or_else(|| { + DataLayerError::InvalidInput( + "admin wallet batch was not prepared".to_string(), + ) + })?; + let target_user_ids: Vec = + serde_json::from_value(row_get(&row, "target_user_ids")?).map_err( + |error| { + DataLayerError::UnexpectedValue(format!( + "admin wallet batch target list is invalid: {error}" + )) + }, + )?; + if !target_user_ids.iter().any(|id| id == &user_id) { + return Err(DataLayerError::InvalidInput( + "user is outside the prepared admin wallet batch".to_string(), + )); + } + let mut outcomes: BTreeMap = + serde_json::from_value(row_get(&row, "user_outcomes")?).map_err( + |error| { + DataLayerError::UnexpectedValue(format!( + "admin wallet batch outcomes are invalid: {error}" + )) + }, + )?; + if let Some(outcome) = outcomes.get(&user_id) { + return Ok(outcome.clone()); + } + let outcome = AdminUserWalletBalanceBatchUserOutcome::Failed(reason); + outcomes.insert(user_id, outcome.clone()); + let value = serde_json::to_value(&outcomes).map_err(|error| { + DataLayerError::UnexpectedValue(format!( + "admin wallet batch outcomes are invalid: {error}" + )) + })?; + sqlx::query( + r#" +UPDATE admin_user_wallet_balance_batches +SET user_outcomes = $3, updated_at_unix_secs = $4 +WHERE admin_user_id = $1 AND idempotency_key = $2 + "#, + ) + .bind(&admin_user_id) + .bind(&idempotency_key) + .bind(value) + .bind(Utc::now().timestamp().max(0)) + .execute(&mut **tx) + .await + .map_postgres_err()?; + Ok(outcome) + }) + }) + .await + } + async fn create_manual_wallet_recharge( &self, mut input: CreateManualWalletRechargeInput, @@ -8503,10 +8849,12 @@ 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::{ - AdjustWalletBalanceInput, CreateManualWalletRechargeInput, CreditAdminPaymentOrderInput, - ProcessPaymentCallbackInput, ProcessPaymentCallbackOutcome, RedeemWalletCodeInput, - RedeemWalletCodeOutcome, WalletLookupKey, WalletMutationOutcome, WalletReadRepository, - WalletWriteRepository, + AdjustWalletBalanceInBatchInput, AdjustWalletBalanceInput, + AdminUserWalletBalanceBatchUserOutcome, CreateManualWalletRechargeInput, + CreditAdminPaymentOrderInput, PrepareAdminUserWalletBalanceBatchInput, + PrepareAdminUserWalletBalanceBatchOutcome, ProcessPaymentCallbackInput, + ProcessPaymentCallbackOutcome, RedeemWalletCodeInput, RedeemWalletCodeOutcome, + WalletLookupKey, WalletMutationOutcome, WalletReadRepository, WalletWriteRepository, }; use sqlx::Row; @@ -8590,6 +8938,7 @@ mod tests { "payment_orders", "payment_callbacks", "wallet_transactions", + "admin_user_wallet_balance_batches", "user_plan_entitlements", "redeem_code_batches", "redeem_codes", @@ -9245,6 +9594,7 @@ mod tests { operator_id: None, description: None, clamp_deduction_to_available_balance: true, + batch_context: None, }; assert_eq!(effective_wallet_adjustment_amount(&input, 13.0), -13.0); assert_eq!(effective_wallet_adjustment_amount(&input, 0.0), -0.0); @@ -9275,6 +9625,7 @@ mod tests { operator_id: Some("admin-user".to_string()), description: Some("bulk deduction".to_string()), clamp_deduction_to_available_balance: true, + batch_context: None, }) .await .expect("bulk adjustment should succeed") @@ -9301,6 +9652,7 @@ mod tests { operator_id: Some("admin-user".to_string()), description: Some("bulk deduction at zero".to_string()), clamp_deduction_to_available_balance: true, + batch_context: None, }) .await .expect("zero-balance adjustment should succeed") @@ -9331,6 +9683,7 @@ mod tests { operator_id: Some("admin-user".to_string()), description: Some("bulk deduction from negative balance".to_string()), clamp_deduction_to_available_balance: true, + batch_context: None, }) .await .expect("legacy negative wallet should remain usable") @@ -9347,6 +9700,137 @@ mod tests { pool.close().await; } + #[tokio::test] + #[ignore = "requires AETHER_TEST_DATABASE_URL and PostgreSQL bootstrap schema"] + async fn live_admin_wallet_balance_batch_replays_committed_and_zero_delta_results_once() { + let pool = isolated_wallet_test_pool().await; + let (wallet_id, user_id) = seed_wallet(&pool).await; + let repository = SqlxWalletRepository::new(pool.clone()); + let admin_user_id = "admin-user".to_string(); + let make_adjustment = + |idempotency_key: &str, amount_usd: f64| AdjustWalletBalanceInBatchInput { + admin_user_id: admin_user_id.clone(), + idempotency_key: idempotency_key.to_string(), + user_id: user_id.clone(), + adjustment: AdjustWalletBalanceInput { + wallet_id: wallet_id.clone(), + amount_usd, + balance_type: "recharge".to_string(), + operator_id: Some(admin_user_id.clone()), + description: Some("idempotency integration test".to_string()), + clamp_deduction_to_available_balance: true, + batch_context: None, + }, + }; + let prepare_batch = |idempotency_key: &str, request_fingerprint: &str| { + PrepareAdminUserWalletBalanceBatchInput { + admin_user_id: admin_user_id.clone(), + idempotency_key: idempotency_key.to_string(), + request_fingerprint: request_fingerprint.to_string(), + target_user_ids: vec![user_id.clone()], + missing_user_ids: Vec::new(), + warnings: Vec::new(), + } + }; + + assert!(matches!( + repository + .prepare_admin_user_wallet_balance_batch(prepare_batch( + "deduct-key", + "deduct-fingerprint" + )) + .await + .unwrap(), + PrepareAdminUserWalletBalanceBatchOutcome::Ready(_) + )); + let deduct = make_adjustment("deduct-key", -50.0); + assert_eq!( + repository + .adjust_admin_user_wallet_balance_batch_user(deduct.clone()) + .await + .unwrap(), + AdminUserWalletBalanceBatchUserOutcome::Succeeded + ); + assert_eq!( + repository + .adjust_admin_user_wallet_balance_batch_user(deduct.clone()) + .await + .unwrap(), + AdminUserWalletBalanceBatchUserOutcome::Succeeded + ); + let wallet = repository + .find(WalletLookupKey::WalletId(&wallet_id)) + .await + .unwrap() + .unwrap(); + assert_eq!(wallet.balance + wallet.gift_balance, 0.0); + + assert!(matches!( + repository + .prepare_admin_user_wallet_balance_batch(prepare_batch( + "zero-key", + "zero-fingerprint" + )) + .await + .unwrap(), + PrepareAdminUserWalletBalanceBatchOutcome::Ready(_) + )); + let zero_delta = make_adjustment("zero-key", -5.0); + assert_eq!( + repository + .adjust_admin_user_wallet_balance_batch_user(zero_delta.clone()) + .await + .unwrap(), + AdminUserWalletBalanceBatchUserOutcome::Succeeded + ); + repository + .adjust_wallet_balance(AdjustWalletBalanceInput { + wallet_id: wallet_id.clone(), + amount_usd: 8.0, + balance_type: "recharge".to_string(), + operator_id: Some(admin_user_id.clone()), + description: Some("recharge after zero-delta batch".to_string()), + clamp_deduction_to_available_balance: true, + batch_context: None, + }) + .await + .unwrap() + .unwrap(); + for replay in [zero_delta, deduct] { + assert_eq!( + repository + .adjust_admin_user_wallet_balance_batch_user(replay) + .await + .unwrap(), + AdminUserWalletBalanceBatchUserOutcome::Succeeded + ); + } + assert_eq!( + repository + .prepare_admin_user_wallet_balance_batch(prepare_batch( + "zero-key", + "different-fingerprint" + )) + .await + .unwrap(), + PrepareAdminUserWalletBalanceBatchOutcome::Conflict + ); + let wallet = repository + .find(WalletLookupKey::WalletId(&wallet_id)) + .await + .unwrap() + .unwrap(); + assert_eq!(wallet.balance + wallet.gift_balance, 8.0); + let transaction_count: i64 = + sqlx::query_scalar("SELECT COUNT(*) FROM wallet_transactions WHERE wallet_id = $1") + .bind(&wallet_id) + .fetch_one(&pool) + .await + .unwrap(); + assert_eq!(transaction_count, 2); + 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 f391d49d2..01262885f 100644 --- a/crates/aether-data/contracts/src/repository/wallet/types.rs +++ b/crates/aether-data/contracts/src/repository/wallet/types.rs @@ -2464,6 +2464,57 @@ pub struct AdjustWalletBalanceInput { pub description: Option, #[serde(default)] pub clamp_deduction_to_available_balance: bool, + #[serde(default)] + pub batch_context: Option, +} + +#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] +pub struct AdminUserWalletBalanceBatchContext { + pub admin_user_id: String, + pub idempotency_key: String, + pub user_id: String, +} + +#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] +pub struct PrepareAdminUserWalletBalanceBatchInput { + pub admin_user_id: String, + pub idempotency_key: String, + pub request_fingerprint: String, + pub target_user_ids: Vec, + pub missing_user_ids: Vec, + pub warnings: Vec, +} + +#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] +pub struct StoredAdminUserWalletBalanceBatch { + pub admin_user_id: String, + pub idempotency_key: String, + pub request_fingerprint: String, + pub target_user_ids: Vec, + pub missing_user_ids: Vec, + pub warnings: Vec, + pub user_outcomes: std::collections::BTreeMap, +} + +#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] +pub enum PrepareAdminUserWalletBalanceBatchOutcome { + Ready(StoredAdminUserWalletBalanceBatch), + Conflict, +} + +#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] +#[serde(tag = "status", content = "reason", rename_all = "snake_case")] +pub enum AdminUserWalletBalanceBatchUserOutcome { + Succeeded, + Failed(String), +} + +#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] +pub struct AdjustWalletBalanceInBatchInput { + pub admin_user_id: String, + pub idempotency_key: String, + pub user_id: String, + pub adjustment: AdjustWalletBalanceInput, } #[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] @@ -2900,6 +2951,47 @@ pub trait WalletWriteRepository: Send + Sync { crate::DataLayerError, >; + async fn prepare_admin_user_wallet_balance_batch( + &self, + _input: PrepareAdminUserWalletBalanceBatchInput, + ) -> Result { + Err(crate::DataLayerError::InvalidInput( + "idempotent admin wallet batches are not available".to_string(), + )) + } + + async fn get_admin_user_wallet_balance_batch( + &self, + _admin_user_id: &str, + _idempotency_key: &str, + _request_fingerprint: &str, + ) -> Result, crate::DataLayerError> { + Err(crate::DataLayerError::InvalidInput( + "idempotent admin wallet batches are not available".to_string(), + )) + } + + async fn adjust_admin_user_wallet_balance_batch_user( + &self, + _input: AdjustWalletBalanceInBatchInput, + ) -> Result { + Err(crate::DataLayerError::InvalidInput( + "idempotent admin wallet batches are not available".to_string(), + )) + } + + async fn record_admin_user_wallet_balance_batch_failure( + &self, + _admin_user_id: &str, + _idempotency_key: &str, + _user_id: &str, + _reason: &str, + ) -> Result { + Err(crate::DataLayerError::InvalidInput( + "idempotent admin wallet batches are not available".to_string(), + )) + } + async fn create_manual_wallet_recharge( &self, input: CreateManualWalletRechargeInput, diff --git a/crates/aether-data/runtime/schema/generated/postgres/baseline/005_wallet_billing.sql b/crates/aether-data/runtime/schema/generated/postgres/baseline/005_wallet_billing.sql index da83615e5..cd7a51f60 100644 --- a/crates/aether-data/runtime/schema/generated/postgres/baseline/005_wallet_billing.sql +++ b/crates/aether-data/runtime/schema/generated/postgres/baseline/005_wallet_billing.sql @@ -343,3 +343,17 @@ CREATE INDEX IF NOT EXISTS idx_redeem_codes_status ON public.redeem_codes USING CREATE INDEX IF NOT EXISTS idx_redeem_codes_redeemed_user ON public.redeem_codes USING btree (redeemed_by_user_id, redeemed_at); CREATE INDEX IF NOT EXISTS idx_redeem_codes_redeemed_order ON public.redeem_codes USING btree (redeemed_payment_order_id); +CREATE TABLE IF NOT EXISTS public.admin_user_wallet_balance_batches ( + admin_user_id character varying(64) NOT NULL, + idempotency_key character varying(128) NOT NULL, + request_fingerprint character varying(64) NOT NULL, + target_user_ids jsonb NOT NULL, + missing_user_ids jsonb NOT NULL, + warnings jsonb NOT NULL, + user_outcomes jsonb NOT NULL, + created_at_unix_secs bigint NOT NULL, + updated_at_unix_secs bigint NOT NULL +); + +ALTER TABLE ONLY public.admin_user_wallet_balance_batches ADD CONSTRAINT admin_user_wallet_balance_batches_pkey PRIMARY KEY (admin_user_id, idempotency_key); + diff --git a/crates/aether-data/runtime/schema/logical/005_wallet_billing.toml b/crates/aether-data/runtime/schema/logical/005_wallet_billing.toml index 99bda2f71..7d1079c6c 100644 --- a/crates/aether-data/runtime/schema/logical/005_wallet_billing.toml +++ b/crates/aether-data/runtime/schema/logical/005_wallet_billing.toml @@ -1366,3 +1366,47 @@ columns = ["redeemed_by_user_id", "redeemed_at"] [[table.redeem_codes.indexes]] name = "idx_redeem_codes_redeemed_order" columns = ["redeemed_payment_order_id"] + +[table.admin_user_wallet_balance_batches] +domain = "wallet_billing" +order = 110 +primary_key = ["admin_user_id", "idempotency_key"] + +[[table.admin_user_wallet_balance_batches.columns]] +name = "admin_user_id" +type = "text_id" +length = 64 + +[[table.admin_user_wallet_balance_batches.columns]] +name = "idempotency_key" +type = "text" +length = 128 + +[[table.admin_user_wallet_balance_batches.columns]] +name = "request_fingerprint" +type = "text" +length = 64 + +[[table.admin_user_wallet_balance_batches.columns]] +name = "target_user_ids" +type = "json" + +[[table.admin_user_wallet_balance_batches.columns]] +name = "missing_user_ids" +type = "json" + +[[table.admin_user_wallet_balance_batches.columns]] +name = "warnings" +type = "json" + +[[table.admin_user_wallet_balance_batches.columns]] +name = "user_outcomes" +type = "json" + +[[table.admin_user_wallet_balance_batches.columns]] +name = "created_at_unix_secs" +type = "unix_seconds" + +[[table.admin_user_wallet_balance_batches.columns]] +name = "updated_at_unix_secs" +type = "unix_seconds" diff --git a/crates/aether-data/runtime/src/repository/wallet/mod.rs b/crates/aether-data/runtime/src/repository/wallet/mod.rs index 8029a4e77..387a8346f 100644 --- a/crates/aether-data/runtime/src/repository/wallet/mod.rs +++ b/crates/aether-data/runtime/src/repository/wallet/mod.rs @@ -16,27 +16,30 @@ pub use aether_data_contracts::repository::wallet::{ wallet_recharge_order_is_checkout_placeholder, wallet_recharge_order_is_reclaimable_placeholder, wallet_recharge_replay_matches, wallet_recharge_response_is_checkout_placeholder, wallet_refund_proof_is_success, - AdjustWalletBalanceInput, AdminPaymentCallbackRecord, AdminPaymentOrderListQuery, - AdminRedeemCodeBatchListQuery, AdminRedeemCodeListQuery, AdminWalletLedgerQuery, - AdminWalletListQuery, AdminWalletPaymentOrderRecord, AdminWalletRefundRecord, - AdminWalletRefundRequestListQuery, AdminWalletTransactionRecord, CanonicalWalletRefundFields, - CompareAndSwapPaymentOrderStripeClientSecretInput, CompleteAdminWalletRefundInput, - CreateAdminRedeemCodeBatchInput, CreateAdminRedeemCodeBatchResult, - CreateManualWalletRechargeInput, CreatePlanPurchaseOrderInput, CreatePlanPurchaseOrderOutcome, - CreateWalletRechargeOrderInput, CreateWalletRechargeOrderOutcome, - CreateWalletRefundRequestInput, CreateWalletRefundRequestOutcome, - CreatedAdminRedeemCodePlaintext, CreditAdminPaymentOrderInput, DeleteAdminRedeemCodeBatchInput, + AdjustWalletBalanceInBatchInput, AdjustWalletBalanceInput, AdminPaymentCallbackRecord, + AdminPaymentOrderListQuery, AdminRedeemCodeBatchListQuery, AdminRedeemCodeListQuery, + AdminUserWalletBalanceBatchContext, AdminUserWalletBalanceBatchUserOutcome, + AdminWalletLedgerQuery, AdminWalletListQuery, AdminWalletPaymentOrderRecord, + AdminWalletRefundRecord, AdminWalletRefundRequestListQuery, AdminWalletTransactionRecord, + CanonicalWalletRefundFields, CompareAndSwapPaymentOrderStripeClientSecretInput, + CompleteAdminWalletRefundInput, CreateAdminRedeemCodeBatchInput, + CreateAdminRedeemCodeBatchResult, CreateManualWalletRechargeInput, + CreatePlanPurchaseOrderInput, CreatePlanPurchaseOrderOutcome, CreateWalletRechargeOrderInput, + CreateWalletRechargeOrderOutcome, CreateWalletRefundRequestInput, + CreateWalletRefundRequestOutcome, CreatedAdminRedeemCodePlaintext, + CreditAdminPaymentOrderInput, DeleteAdminRedeemCodeBatchInput, DisableAdminRedeemCodeBatchInput, DisableAdminRedeemCodeInput, FailAdminWalletRefundInput, - FailWalletRechargeCheckoutInput, InitializeAuthWalletOutcome, ProcessAdminWalletRefundInput, - ProcessPaymentCallbackInput, ProcessPaymentCallbackOutcome, ReclaimWalletRechargeCheckoutInput, - RedeemWalletCodeInput, RedeemWalletCodeOutcome, StoredAdminPaymentCallback, - StoredAdminPaymentCallbackPage, StoredAdminPaymentOrder, StoredAdminPaymentOrderPage, - StoredAdminRedeemCode, StoredAdminRedeemCodeBatch, StoredAdminRedeemCodeBatchPage, - StoredAdminRedeemCodePage, StoredAdminWalletLedgerItem, StoredAdminWalletLedgerPage, - StoredAdminWalletListItem, StoredAdminWalletListPage, StoredAdminWalletRefund, - StoredAdminWalletRefundPage, StoredAdminWalletRefundRequestItem, - StoredAdminWalletRefundRequestPage, StoredAdminWalletTransaction, - StoredAdminWalletTransactionPage, StoredWalletDailyUsageLedger, + FailWalletRechargeCheckoutInput, InitializeAuthWalletOutcome, + PrepareAdminUserWalletBalanceBatchInput, PrepareAdminUserWalletBalanceBatchOutcome, + ProcessAdminWalletRefundInput, ProcessPaymentCallbackInput, ProcessPaymentCallbackOutcome, + ReclaimWalletRechargeCheckoutInput, RedeemWalletCodeInput, RedeemWalletCodeOutcome, + StoredAdminPaymentCallback, StoredAdminPaymentCallbackPage, StoredAdminPaymentOrder, + StoredAdminPaymentOrderPage, StoredAdminRedeemCode, StoredAdminRedeemCodeBatch, + StoredAdminRedeemCodeBatchPage, StoredAdminRedeemCodePage, StoredAdminUserWalletBalanceBatch, + StoredAdminWalletLedgerItem, StoredAdminWalletLedgerPage, StoredAdminWalletListItem, + StoredAdminWalletListPage, StoredAdminWalletRefund, StoredAdminWalletRefundPage, + StoredAdminWalletRefundRequestItem, StoredAdminWalletRefundRequestPage, + StoredAdminWalletTransaction, StoredAdminWalletTransactionPage, StoredWalletDailyUsageLedger, StoredWalletDailyUsageLedgerPage, StoredWalletSnapshot, UpdateAdminWalletRefundGatewayInput, UpdateWalletRechargeCheckoutInput, WalletLookupKey, WalletMutationOutcome, WalletReadRepository, WalletReadSeed, WalletReadSnapshot, WalletRepository, diff --git a/frontend/src/api/users.ts b/frontend/src/api/users.ts index 1f216e31f..0902b7775 100644 --- a/frontend/src/api/users.ts +++ b/frontend/src/api/users.ts @@ -165,6 +165,7 @@ export interface UserBatchBalanceActionRequest { selection: UserBatchSelection action: 'adjust_wallet_balance' payload: UserBatchBalanceAdjustmentPayload + idempotency_key: string } export type UserBatchActionRequest = diff --git a/frontend/src/features/users/components/UserBatchActionDialog.vue b/frontend/src/features/users/components/UserBatchActionDialog.vue index 4712016a3..fff9293be 100644 --- a/frontend/src/features/users/components/UserBatchActionDialog.vue +++ b/frontend/src/features/users/components/UserBatchActionDialog.vue @@ -88,6 +88,39 @@

+ + + () const usersStore = useUsersStore() +const authStore = useAuthStore() +const walletRetryCoordinator = createUserBatchWalletRetryCoordinator({ + scope: () => authStore.user?.id ?? null, +}) const { success, warning, error } = useToast() const { legacyT, locale } = useI18n() @@ -178,6 +228,8 @@ const previewItems = ref([]) const resolvedTotal = ref(null) const executing = ref(false) const lastResult = ref(null) +const pendingWalletBatch = ref(null) +const pendingWalletReadError = ref(false) const hasAnyTarget = computed(() => props.selectedCount > 0 || selectedGroupIds.value.length > 0) const impactCount = computed(() => resolvedTotal.value ?? props.selectedCount) @@ -185,11 +237,23 @@ const balancePayload = computed(() => buildUserBatchBalanceAdjustmentPayload( balanceOperation.value, balanceAmount.value, )) +const walletAdjustmentRequest = computed(() => ( + balancePayload.value === null + ? null + : { selection: buildSelection(), action: 'adjust_wallet_balance', payload: balancePayload.value } +)) +const walletRequestMismatch = computed(() => ( + pendingWalletBatch.value !== null + && selectedAction.value === 'adjust_wallet_balance' + && (walletAdjustmentRequest.value === null + || !matchesPendingWalletRequest(pendingWalletBatch.value, walletAdjustmentRequest.value)) +)) const canExecute = computed(() => ( hasAnyTarget.value && !previewLoading.value && !executing.value - && (selectedAction.value !== 'adjust_wallet_balance' || balancePayload.value !== null) + && (selectedAction.value !== 'adjust_wallet_balance' + || (balancePayload.value !== null && !pendingWalletReadError.value && !walletRequestMismatch.value)) )) const selectedActionLabel = computed(() => ( USER_BATCH_ACTION_OPTIONS.find((action) => action.value === selectedAction.value)?.label ?? '批量操作' @@ -208,6 +272,11 @@ const targetRoleWarning = computed(() => { return legacyT('提示:设置为普通用户会移除目标用户的管理员权限。') }) const executeButtonLabel = computed(() => legacyT(`确认${selectedActionLabel.value}(${impactCount.value})`)) +const pendingWalletOperationLabel = computed(() => ( + pendingWalletBatch.value?.request.payload.operation === 'deduct' + ? legacyT('扣减') + : legacyT('增加') +)) const lastResultLabel = computed(() => { if (!lastResult.value) return '' if (lastResult.value.interrupted) { @@ -230,8 +299,15 @@ watch( (open) => { if (!open) return resetLocalState() + refreshPendingWalletBatch() void resolvePreview() }, + { immediate: true }, +) + +watch( + () => authStore.user?.id, + () => refreshPendingWalletBatch(), ) watch( @@ -255,6 +331,16 @@ function resetLocalState(): void { lastResult.value = null } +function refreshPendingWalletBatch(): void { + try { + pendingWalletBatch.value = walletRetryCoordinator.getPending() + pendingWalletReadError.value = false + } catch { + pendingWalletBatch.value = null + pendingWalletReadError.value = true + } +} + function buildSelection(): UserBatchSelection { const group_ids = selectedGroupIds.value.length > 0 ? [...selectedGroupIds.value] : undefined if (props.selectAllFiltered) { @@ -301,7 +387,8 @@ function buildRolePayload(): UserBatchRolePayload { async function executeBatchAction(): Promise { if (!canExecute.value) return const selection = buildSelection() - let request: UserBatchActionRequest + let request: Exclude | null = null + let walletRequest: UserBatchWalletAdjustmentRequest | null = null if (selectedAction.value === 'update_access_control') { const payload = buildAccessControlPayload() if (payload === null) { @@ -310,11 +397,11 @@ async function executeBatchAction(): Promise { } request = { selection, action: 'update_access_control', payload } } else if (selectedAction.value === 'adjust_wallet_balance') { - if (balancePayload.value === null) { + if (walletAdjustmentRequest.value === null) { warning(legacyT('请输入大于 0 的有限金额')) return } - request = { selection, action: 'adjust_wallet_balance', payload: balancePayload.value } + walletRequest = walletAdjustmentRequest.value } else if (selectedAction.value === 'update_role') { request = { selection, action: 'update_role', payload: buildRolePayload() } } else { @@ -323,6 +410,15 @@ async function executeBatchAction(): Promise { executing.value = true try { + if (walletRequest) { + const result = await walletRetryCoordinator.execute( + walletRequest, + (keyedRequest) => usersStore.batchAction(keyedRequest), + ) + handleWalletBatchResult(result) + return + } + if (request === null) return const result = await usersStore.batchAction(request) lastResult.value = result if (result.interrupted) { @@ -338,9 +434,63 @@ async function executeBatchAction(): Promise { } emit('completed', result) } catch (err) { - error(legacyT(parseApiError(err, '批量操作失败')), legacyT('批量操作失败')) + if (walletRequest) { + refreshPendingWalletBatch() + if (err instanceof WalletIdempotencyPersistenceUnavailableError) { + warning(legacyT('浏览器无法安全保存钱包批量请求,本次请求未发送。')) + } else if ( + err instanceof WalletIdempotencyUnavailableError + || err instanceof WalletIdempotencyScopeUnavailableError + || err instanceof WalletIdempotencyScopeChangedError + ) { + warning(legacyT('无法确认管理员身份或安全生成钱包批量请求标识,请求未发送。')) + } else { + warning(legacyT('钱包批量调整结果未知或可能部分完成。请重试原请求,不要开始新的余额调整。')) + } + } else { + error(legacyT(parseApiError(err, '批量操作失败')), legacyT('批量操作失败')) + } } finally { executing.value = false } } + +async function retryPendingWalletBatch(): Promise { + if (executing.value) return + executing.value = true + try { + const result = await walletRetryCoordinator.retry( + (request) => usersStore.batchAction(request), + ) + if (result) handleWalletBatchResult(result) + else refreshPendingWalletBatch() + } catch (err) { + refreshPendingWalletBatch() + if (err instanceof WalletIdempotencyPersistenceUnavailableError) { + warning(legacyT('浏览器无法安全保存钱包批量请求,本次请求未发送。')) + } else if ( + err instanceof WalletIdempotencyScopeUnavailableError + || err instanceof WalletIdempotencyScopeChangedError + ) { + warning(legacyT('无法确认管理员身份,请求未发送。')) + } else { + warning(legacyT('钱包批量调整结果未知或可能部分完成。请重试原请求,不要开始新的余额调整。')) + } + } finally { + executing.value = false + } +} + +function handleWalletBatchResult(result: UserBatchActionResponse): void { + lastResult.value = result + refreshPendingWalletBatch() + if (result.interrupted) { + warning(`${lastResultLabel.value};${legacyT('结果可能部分完成,请仅重试原请求,不要开始新的余额调整。')}`) + } else { + const message = legacyT(`批量操作完成:成功 ${result.success} 个,失败 ${result.failed} 个`) + if (result.failed > 0) warning(message) + else success(message) + } + emit('completed', result) +} 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) + }, + } +}