diff --git a/.github/workflows/rust-ci.yml b/.github/workflows/rust-ci.yml index c8d6dfdc2..3eaa76a92 100644 --- a/.github/workflows/rust-ci.yml +++ b/.github/workflows/rust-ci.yml @@ -603,6 +603,13 @@ jobs: AETHER_TEST_DATABASE_URL: postgres://aether:aether@127.0.0.1:5432/aether_test run: cargo test -p aether-data-postgres live_payment_callback --lib -- --ignored --nocapture + - name: Run Postgres batch wallet deduction regression + env: + RUSTC_WRAPPER: sccache + SCCACHE_GHA_ENABLED: "true" + AETHER_TEST_DATABASE_URL: postgres://aether:aether@127.0.0.1:5432/aether_test + run: cargo test -p aether-data-postgres live_bulk_wallet_adjustment_persists_actual_delta_and_skips_zero_ledger --lib -- --ignored --nocapture + - name: Run Postgres API key lifecycle tests env: RUSTC_WRAPPER: sccache 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 901563379..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, @@ -1066,13 +1068,81 @@ impl GatewayDataState { pub(crate) async fn adjust_wallet_balance( &self, input: AdjustWalletBalanceInput, - ) -> Result, DataLayerError> { + ) -> Result)>, DataLayerError> + { match &self.wallet_writer { Some(repository) => repository.adjust_wallet_balance(input).await, None => Ok(None), } } + 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/billing/wallets/mutations/adjust.rs b/apps/aether-gateway/src/handlers/admin/billing/wallets/mutations/adjust.rs index 8efabaa0b..f3d2809af 100644 --- a/apps/aether-gateway/src/handlers/admin/billing/wallets/mutations/adjust.rs +++ b/apps/aether-gateway/src/handlers/admin/billing/wallets/mutations/adjust.rs @@ -63,13 +63,14 @@ pub(in super::super) async fn build_admin_wallet_adjust_response( } let operator_id = admin_wallet_operator_id(request_context); let has_wallet_writer = state.has_wallet_data_writer(); - let Some((wallet, transaction)) = state + let Some((wallet, Some(transaction))) = state .admin_adjust_wallet_balance( &wallet_id, amount_usd, &balance_type, operator_id.as_deref(), description.as_deref(), + false, ) .await? else { diff --git a/apps/aether-gateway/src/handlers/admin/request/billing.rs b/apps/aether-gateway/src/handlers/admin/request/billing.rs index 3ecde6481..38a7b9ddf 100644 --- a/apps/aether-gateway/src/handlers/admin/request/billing.rs +++ b/apps/aether-gateway/src/handlers/admin/request/billing.rs @@ -387,10 +387,11 @@ impl<'a> AdminAppState<'a> { balance_type: &str, operator_id: Option<&str>, description: Option<&str>, + clamp_deduction_to_available_balance: bool, ) -> Result< Option<( aether_data::repository::wallet::StoredWalletSnapshot, - crate::AdminWalletTransactionRecord, + Option, )>, GatewayError, > { @@ -401,6 +402,7 @@ impl<'a> AdminAppState<'a> { balance_type, operator_id, description, + clamp_deduction_to_available_balance, ) .await } diff --git a/apps/aether-gateway/src/handlers/admin/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 75d945001..d56b3f08d 100644 --- a/apps/aether-gateway/src/handlers/admin/users/batch.rs +++ b/apps/aether-gateway/src/handlers/admin/users/batch.rs @@ -1,11 +1,16 @@ use super::{ build_admin_users_bad_request_response, build_admin_users_permission_denied_response, - build_admin_users_read_only_response, disabled_user_policy_detail, disabled_user_policy_field, + build_admin_users_read_only_response, build_admin_users_wallet_permission_denied_response, + disabled_user_policy_detail, disabled_user_policy_field, + management_token_may_adjust_admin_wallet_balance, management_token_may_administer_user_accounts, normalize_admin_user_role, }; use crate::handlers::admin::request::{AdminAppState, AdminRequestContext}; 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, @@ -13,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, @@ -27,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, @@ -40,6 +46,7 @@ struct AdminUserBatchActionRequest { selection: AdminUserSelectionRequest, action: String, payload: Option, + idempotency_key: Option, } #[derive(Debug, serde::Deserialize)] @@ -48,6 +55,8 @@ struct RawAdminUserBatchActionRequest { action: String, #[serde(default)] payload: Option, + #[serde(default)] + idempotency_key: Option, } #[derive(Debug, Clone, Default)] @@ -68,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, @@ -88,6 +97,7 @@ struct AdminUserBatchMutation { role: Option, is_active: Option, unlimited: Option, + wallet_balance_adjustment: Option, modified_fields: Vec<&'static str>, } @@ -97,6 +107,28 @@ impl AdminUserBatchMutation { } } +#[derive(Debug, Clone, Copy)] +struct AdminUserWalletBalanceAdjustment { + operation: AdminUserWalletBalanceOperation, + amount: f64, +} + +#[derive(Debug, Clone, Copy)] +enum AdminUserWalletBalanceOperation { + Add, + Deduct, +} + +enum AdminBatchWalletBalanceAdjustmentError { + WalletLookup, + BalanceAdjustment, +} + +enum AdminBatchWalletLimitModeError { + WalletLookup, + Mutation, +} + pub(in super::super) async fn build_admin_resolve_user_selection_response( state: &AdminAppState<'_>, _request_context: &AdminRequestContext<'_>, @@ -128,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)), @@ -177,9 +228,29 @@ pub(in super::super) async fn build_admin_user_batch_action_response( .iter() .map(|user_id| json!({ "user_id": user_id, "reason": "用户不存在或已删除" })) .collect::>(); + let mut completed_user_ids = Vec::new(); + let mut uncertain_user_ids = Vec::new(); + let mut unprocessed_user_ids = Vec::new(); + let mut interrupted = false; - for item in &resolved.items { - if state.find_user_auth_by_id(&item.user_id).await?.is_none() { + for (item_index, item) in resolved.items.iter().enumerate() { + let user = match state.find_user_auth_by_id(&item.user_id).await { + Ok(user) => user, + Err(_) => { + record_batch_action_interruption( + &resolved.items, + item_index, + false, + "读取用户状态失败,批次已中止,该用户未执行", + &mut failures, + &mut uncertain_user_ids, + &mut unprocessed_user_ids, + ); + interrupted = true; + break; + } + }; + if user.is_none() { failures.push(json!({ "user_id": item.user_id, "reason": "用户不存在或已删除", @@ -202,17 +273,92 @@ pub(in super::super) async fn build_admin_user_batch_action_response( } if let Some(unlimited) = mutation.unlimited { - if !apply_batch_user_wallet_limit_mode(state, &item.user_id, unlimited).await? { - failures.push(json!({ - "user_id": item.user_id, - "reason": "用户钱包不可用", - })); - continue; + match apply_batch_user_wallet_limit_mode(state, &item.user_id, unlimited).await { + Ok(true) => {} + Ok(false) => { + failures.push(json!({ + "user_id": item.user_id, + "reason": "用户钱包不可用", + })); + continue; + } + Err(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, + true, + "用户钱包更新结果未确认,批次已中止,请核对钱包后再重试", + &mut failures, + &mut uncertain_user_ids, + &mut unprocessed_user_ids, + ); + interrupted = true; + break; + } } } - if mutation.has_auth_user_fields() - && state + if let Some(adjustment) = mutation.wallet_balance_adjustment { + match apply_batch_user_wallet_balance_adjustment( + state, + &item.user_id, + adjustment, + current_admin_user_id, + ) + .await + { + Ok(true) => {} + Ok(false) => { + failures.push(json!({ + "user_id": item.user_id, + "reason": "用户钱包不可用", + })); + continue; + } + Err(AdminBatchWalletBalanceAdjustmentError::WalletLookup) => { + record_batch_action_interruption( + &resolved.items, + item_index, + false, + "读取用户钱包失败,批次已中止,该用户未执行", + &mut failures, + &mut uncertain_user_ids, + &mut unprocessed_user_ids, + ); + interrupted = true; + break; + } + Err(AdminBatchWalletBalanceAdjustmentError::BalanceAdjustment) => { + record_batch_action_interruption( + &resolved.items, + item_index, + true, + "余额调整结果未确认,批次已中止,请核对钱包后再重试", + &mut failures, + &mut uncertain_user_ids, + &mut unprocessed_user_ids, + ); + interrupted = true; + break; + } + } + } + + if mutation.has_auth_user_fields() { + let updated_user = match state .update_local_auth_user_admin_fields( &item.user_id, mutation.role.clone(), @@ -226,22 +372,39 @@ pub(in super::super) async fn build_admin_user_batch_action_response( None, mutation.is_active, ) - .await? - .is_none() - { - failures.push(json!({ - "user_id": item.user_id, - "reason": "用户不存在或已删除", - })); - continue; + .await + { + Ok(user) => user, + Err(_) => { + record_batch_action_interruption( + &resolved.items, + item_index, + true, + "用户更新结果未确认,批次已中止,请核对后再重试", + &mut failures, + &mut uncertain_user_ids, + &mut unprocessed_user_ids, + ); + interrupted = true; + break; + } + }; + if updated_user.is_none() { + failures.push(json!({ + "user_id": item.user_id, + "reason": "用户不存在或已删除", + })); + continue; + } } success += 1; + completed_user_ids.push(item.user_id.clone()); } let failed = failures.len(); let total = success + failed; - let response = Json(json!({ + let mut response_payload = json!({ "total": total, "success": success, "failed": failed, @@ -249,8 +412,14 @@ pub(in super::super) async fn build_admin_user_batch_action_response( "warnings": resolved.warnings, "action": request.action.trim().to_ascii_lowercase(), "modified_fields": mutation.modified_fields, - })) - .into_response(); + "interrupted": interrupted, + }); + if interrupted { + response_payload["completed_user_ids"] = json!(completed_user_ids); + response_payload["uncertain_user_ids"] = json!(uncertain_user_ids); + response_payload["unprocessed_user_ids"] = json!(unprocessed_user_ids); + } + let response = Json(response_payload).into_response(); Ok(attach_admin_audit_response( response, @@ -261,6 +430,408 @@ pub(in super::super) async fn build_admin_user_batch_action_response( )) } +fn record_batch_action_interruption( + items: &[AdminUserSelectionItem], + item_index: usize, + current_result_uncertain: bool, + reason: &str, + failures: &mut Vec, + uncertain_user_ids: &mut Vec, + unprocessed_user_ids: &mut Vec, +) { + let current_item = &items[item_index]; + failures.push(json!({ + "user_id": current_item.user_id, + "reason": reason, + })); + if current_result_uncertain { + uncertain_user_ids.push(current_item.user_id.clone()); + } else { + unprocessed_user_ids.push(current_item.user_id.clone()); + } + + for item in items.iter().skip(item_index + 1) { + failures.push(json!({ + "user_id": item.user_id, + "reason": "因前序错误未执行", + })); + unprocessed_user_ids.push(item.user_id.clone()); + } +} + +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 { @@ -286,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()), @@ -604,10 +1176,37 @@ fn parse_batch_mutation( }), "update_access_control" => parse_access_control_mutation(payload), "update_role" => parse_role_mutation(payload), + "adjust_wallet_balance" => parse_wallet_balance_adjustment_mutation(payload), _ => Err("不支持的批量操作".to_string()), } } +fn parse_wallet_balance_adjustment_mutation( + payload: Option, +) -> Result { + let Some(Value::Object(payload)) = payload else { + return Err("payload 必须是对象".to_string()); + }; + let operation = match payload.get("operation").and_then(Value::as_str) { + Some("add") => AdminUserWalletBalanceOperation::Add, + Some("deduct") => AdminUserWalletBalanceOperation::Deduct, + _ => return Err("operation 必须为 add 或 deduct".to_string()), + }; + let amount = payload + .get("amount") + .and_then(Value::as_f64) + .ok_or_else(|| "amount 必须为大于 0 的有限数字".to_string())?; + if !amount.is_finite() || amount <= 0.0 { + return Err("amount 必须为大于 0 的有限数字".to_string()); + } + + Ok(AdminUserBatchMutation { + wallet_balance_adjustment: Some(AdminUserWalletBalanceAdjustment { operation, amount }), + modified_fields: vec!["wallet_balance"], + ..AdminUserBatchMutation::default() + }) +} + fn parse_role_mutation(payload: Option) -> Result { let Some(Value::Object(payload)) = payload else { return Err("payload 必须是对象".to_string()); @@ -700,30 +1299,68 @@ 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()), } } +async fn apply_batch_user_wallet_balance_adjustment( + state: &AdminAppState<'_>, + user_id: &str, + adjustment: AdminUserWalletBalanceAdjustment, + operator_id: Option<&str>, +) -> Result { + // Resolve only the wallet ID; the repository clamps the deduction under its row lock. + let Some(wallet) = state + .find_wallet(aether_data::repository::wallet::WalletLookupKey::UserId( + user_id, + )) + .await + .map_err(|_| AdminBatchWalletBalanceAdjustmentError::WalletLookup)? + else { + return Ok(false); + }; + let amount = match adjustment.operation { + AdminUserWalletBalanceOperation::Add => adjustment.amount, + AdminUserWalletBalanceOperation::Deduct => -adjustment.amount, + }; + + state + .admin_adjust_wallet_balance( + &wallet.id, + amount, + "recharge", + operator_id, + Some("管理员批量调整用户余额"), + true, + ) + .await + .map_err(|_| AdminBatchWalletBalanceAdjustmentError::BalanceAdjustment) + .map(|result| result.is_some()) +} + fn build_admin_user_batch_bad_request_response(detail: String) -> Response { if detail.as_str() == "缺少 user_id" { return build_admin_users_bad_request_response("缺少 user_id"); diff --git a/apps/aether-gateway/src/handlers/admin/users/mod.rs b/apps/aether-gateway/src/handlers/admin/users/mod.rs index e5c854748..5a5e9d16a 100644 --- a/apps/aether-gateway/src/handlers/admin/users/mod.rs +++ b/apps/aether-gateway/src/handlers/admin/users/mod.rs @@ -49,12 +49,14 @@ use self::shared::AdminUpdateUserPatch; use self::shared::{ admin_default_user_initial_gift, build_admin_users_bad_request_response, build_admin_users_data_unavailable_response, build_admin_users_permission_denied_response, - build_admin_users_read_only_response, disabled_user_policy_detail, disabled_user_policy_field, - format_optional_datetime_iso8601, legacy_admin_list_policy_mode, - legacy_admin_rate_limit_policy_mode, management_token_may_administer_user_accounts, - normalize_admin_optional_user_email, normalize_admin_user_group_ids, normalize_admin_user_role, - normalize_admin_username, validate_admin_user_password, AdminCreateUserApiKeyRequest, - AdminCreateUserRequest, AdminToggleUserApiKeyLockRequest, AdminUpdateUserApiKeyRequest, + build_admin_users_read_only_response, build_admin_users_wallet_permission_denied_response, + disabled_user_policy_detail, disabled_user_policy_field, format_optional_datetime_iso8601, + legacy_admin_list_policy_mode, legacy_admin_rate_limit_policy_mode, + management_token_may_adjust_admin_wallet_balance, + management_token_may_administer_user_accounts, normalize_admin_optional_user_email, + normalize_admin_user_group_ids, normalize_admin_user_role, normalize_admin_username, + validate_admin_user_password, AdminCreateUserApiKeyRequest, AdminCreateUserRequest, + AdminToggleUserApiKeyLockRequest, AdminUpdateUserApiKeyRequest, }; pub(crate) use self::shared::{ normalize_admin_list_policy_mode, normalize_admin_rate_limit_policy_mode, diff --git a/apps/aether-gateway/src/handlers/admin/users/shared.rs b/apps/aether-gateway/src/handlers/admin/users/shared.rs index b304932c9..b74bd0c9a 100644 --- a/apps/aether-gateway/src/handlers/admin/users/shared.rs +++ b/apps/aether-gateway/src/handlers/admin/users/shared.rs @@ -165,6 +165,18 @@ pub(super) fn management_token_may_administer_user_accounts( }) } +pub(super) fn management_token_may_adjust_admin_wallet_balance( + request_context: &crate::handlers::admin::request::AdminRequestContext<'_>, +) -> bool { + request_context.decision().is_some_and(|decision| { + crate::control::management_token_principal_has_permission(decision, "admin:wallets:write") + || crate::control::management_token_principal_has_permission( + decision, + "admin:wallets:admin", + ) + }) +} + pub(super) fn build_admin_users_permission_denied_response( request_context: &crate::handlers::admin::request::AdminRequestContext<'_>, ) -> Response { @@ -192,6 +204,34 @@ pub(super) fn build_admin_users_permission_denied_response( ) } +pub(super) fn build_admin_users_wallet_permission_denied_response( + request_context: &crate::handlers::admin::request::AdminRequestContext<'_>, +) -> Response { + let actor_id = request_context + .decision() + .and_then(|decision| decision.admin_principal.as_ref()) + .and_then(|principal| principal.management_token_id.as_deref()) + .unwrap_or("unknown"); + crate::handlers::admin::shared::attach_admin_audit_response( + ( + http::StatusCode::FORBIDDEN, + Json(json!({ + "detail": "management token permission denied", + "required_permissions": ["admin:wallets:write", "admin:wallets:admin"], + "permission_mode": "any_of", + "route_family": request_context.route_family(), + "route_kind": request_context.route_kind(), + "request_path": request_context.path(), + })), + ) + .into_response(), + "admin_user_wallet_balance_permission_denied", + "permission_denied", + "admin_user_wallet_balance", + actor_id, + ) +} + pub(super) fn normalize_admin_optional_user_email( value: Option<&str>, ) -> Result, String> { @@ -397,9 +437,52 @@ pub(super) fn format_optional_datetime_iso8601( #[cfg(test)] mod tests { - use super::{normalize_admin_user_api_formats, AdminUpdateUserApiKeyRequest}; + use super::{ + build_admin_users_wallet_permission_denied_response, normalize_admin_user_api_formats, + AdminUpdateUserApiKeyRequest, + }; + use crate::control::{GatewayControlDecision, GatewayPublicRequestContext}; + use crate::handlers::admin::request::AdminRequestContext; + use axum::http::{HeaderMap, Method, Uri}; use serde_json::json; + #[test] + fn wallet_permission_denial_uses_wallet_audit_category() { + let uri: Uri = "/api/admin/users/batch-action" + .parse() + .expect("uri should parse"); + let method = Method::POST; + let headers = HeaderMap::new(); + let decision = GatewayControlDecision::synthetic( + uri.path(), + Some("admin_proxy".to_string()), + Some("users_manage".to_string()), + Some("batch_user_action".to_string()), + Some("admin:users".to_string()), + ); + let context = GatewayPublicRequestContext::from_request_parts( + "trace-wallet-permission-denied", + &method, + &uri, + &headers, + Some(decision), + ); + let request_context = AdminRequestContext::new(&context); + + let response = build_admin_users_wallet_permission_denied_response(&request_context); + let event = response + .extensions() + .get::() + .expect("wallet denial should attach an audit event"); + + assert_eq!( + event.event_name, + "admin_user_wallet_balance_permission_denied" + ); + assert_eq!(event.action, "permission_denied"); + assert_eq!(event.target_type, "admin_user_wallet_balance"); + } + #[test] fn admin_user_api_formats_accept_current_canonical_signatures() { assert_eq!( diff --git a/apps/aether-gateway/src/state/app.rs b/apps/aether-gateway/src/state/app.rs index 1851b7c0c..962c003fc 100644 --- a/apps/aether-gateway/src/state/app.rs +++ b/apps/aether-gateway/src/state/app.rs @@ -492,6 +492,25 @@ pub struct AppState { Arc>>, >, #[cfg(test)] + pub(crate) auth_wallet_adjustment_error_for_tests: Option, + #[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 e7cde0f8a..52420dfaf 100644 --- a/apps/aether-gateway/src/state/core.rs +++ b/apps/aether-gateway/src/state/core.rs @@ -471,6 +471,16 @@ impl AppState { #[cfg(test)] auth_wallet_store: Some(Arc::new(StdMutex::new(HashMap::new()))), #[cfg(test)] + auth_wallet_adjustment_error_for_tests: None, + #[cfg(test)] + 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 8501b155a..c6b6b1b39 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, @@ -10,13 +200,21 @@ impl AppState { balance_type: &str, operator_id: Option<&str>, description: Option<&str>, + clamp_deduction_to_available_balance: bool, ) -> Result< Option<( aether_data::repository::wallet::StoredWalletSnapshot, - AdminWalletTransactionRecord, + Option, )>, GatewayError, > { + #[cfg(test)] + if self.auth_wallet_adjustment_error_for_tests.as_deref() == Some(wallet_id) { + return Err(GatewayError::Internal( + "injected test wallet adjustment failure".to_string(), + )); + } + #[cfg(test)] if let Some(store) = self.auth_wallet_store.as_ref() { let mut guard = store.lock().expect("auth wallet store should lock"); @@ -27,6 +225,18 @@ impl AppState { let before_recharge = wallet.balance; let before_gift = wallet.gift_balance; let before_total = before_recharge + before_gift; + let amount_usd = if clamp_deduction_to_available_balance && amount_usd < 0.0 { + if before_total < 0.0 { + -before_total + } else { + -(-amount_usd).min(before_total) + } + } else { + amount_usd + }; + if amount_usd == 0.0 { + return Ok(Some((wallet.clone(), None))); + } let mut after_recharge = before_recharge; let mut after_gift = before_gift; @@ -90,7 +300,7 @@ impl AppState { let updated_wallet = wallet.clone(); drop(guard); self.invalidate_auth_context_cache(); - return Ok(Some((updated_wallet, transaction))); + return Ok(Some((updated_wallet, Some(transaction)))); } Ok(self @@ -100,10 +310,15 @@ impl AppState { balance_type: balance_type.to_string(), operator_id: operator_id.map(ToOwned::to_owned), description: description.map(ToOwned::to_owned), + clamp_deduction_to_available_balance, + batch_context: None, }) .await? .map(|(wallet, transaction)| { - (wallet, stored_wallet_transaction_to_gateway(transaction)) + ( + wallet, + transaction.map(stored_wallet_transaction_to_gateway), + ) })) } diff --git a/apps/aether-gateway/src/state/runtime/wallet/mutations.rs b/apps/aether-gateway/src/state/runtime/wallet/mutations.rs index 995ba626f..40b4960a8 100644 --- a/apps/aether-gateway/src/state/runtime/wallet/mutations.rs +++ b/apps/aether-gateway/src/state/runtime/wallet/mutations.rs @@ -112,7 +112,7 @@ impl AppState { ) -> Result< Option<( aether_data::repository::wallet::StoredWalletSnapshot, - aether_data::repository::wallet::StoredAdminWalletTransaction, + Option, )>, GatewayError, > { diff --git a/apps/aether-gateway/src/state/runtime/wallet/reads.rs b/apps/aether-gateway/src/state/runtime/wallet/reads.rs index 90896a419..a0bf2d187 100644 --- a/apps/aether-gateway/src/state/runtime/wallet/reads.rs +++ b/apps/aether-gateway/src/state/runtime/wallet/reads.rs @@ -5,6 +5,19 @@ impl AppState { &self, lookup: aether_data::repository::wallet::WalletLookupKey<'_>, ) -> Result, GatewayError> { + #[cfg(test)] + if let Some(failed_user_id) = self.auth_wallet_lookup_error_for_tests.as_deref() { + let lookup_user_id = match &lookup { + aether_data::repository::wallet::WalletLookupKey::UserId(user_id) => Some(*user_id), + _ => None, + }; + if lookup_user_id == Some(failed_user_id) { + return Err(GatewayError::Internal( + "injected test wallet lookup failure".to_string(), + )); + } + } + #[cfg(test)] if let Some(store) = self.auth_wallet_store.as_ref() { let wallet = { diff --git a/apps/aether-gateway/src/state/testing.rs b/apps/aether-gateway/src/state/testing.rs index bdbf662e9..ebc27b469 100644 --- a/apps/aether-gateway/src/state/testing.rs +++ b/apps/aether-gateway/src/state/testing.rs @@ -472,6 +472,27 @@ impl AppState { self } + pub(crate) fn fail_auth_wallet_adjustment_for_tests( + mut self, + wallet_id: impl Into, + ) -> Self { + self.auth_wallet_adjustment_error_for_tests = Some(wallet_id.into()); + self + } + + pub(crate) fn fail_auth_wallet_lookup_for_tests(mut self, user_id: impl Into) -> Self { + self.auth_wallet_lookup_error_for_tests = Some(user_id.into()); + 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.rs b/apps/aether-gateway/src/tests/control/admin.rs index 442d422e7..ab54ff946 100644 --- a/apps/aether-gateway/src/tests/control/admin.rs +++ b/apps/aether-gateway/src/tests/control/admin.rs @@ -21,5 +21,6 @@ mod system; mod system_import; mod usage; mod users; +mod users_batch; mod video_tasks; mod wallets; diff --git a/apps/aether-gateway/src/tests/control/admin/users_batch.rs b/apps/aether-gateway/src/tests/control/admin/users_batch.rs new file mode 100644 index 000000000..35efa942e --- /dev/null +++ b/apps/aether-gateway/src/tests/control/admin/users_batch.rs @@ -0,0 +1,650 @@ +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::Arc; + +use aether_data::repository::management_tokens::InMemoryManagementTokenRepository; +use aether_data::repository::users::{InMemoryUserReadRepository, StoredUserAuthRecord}; +use aether_data::repository::wallet::StoredWalletSnapshot; +use axum::http::StatusCode; +use chrono::Utc; +use reqwest::{Client, RequestBuilder, Response}; +use serde_json::{json, Value}; + +use super::super::{ + build_router_with_state, hash_management_token, sample_management_token, start_server, AppState, +}; +use crate::data::GatewayDataState; + +fn admin_headers(request: RequestBuilder) -> RequestBuilder { + request + .header(crate::constants::GATEWAY_HEADER, "rust-phase3b") + .header(crate::constants::TRUSTED_ADMIN_USER_ID_HEADER, "admin-user") + .header(crate::constants::TRUSTED_ADMIN_USER_ROLE_HEADER, "admin") + .header( + crate::constants::TRUSTED_ADMIN_SESSION_ID_HEADER, + "session-admin", + ) +} + +fn sample_user(user_id: &str) -> StoredUserAuthRecord { + sample_user_with_role(user_id, "user") +} + +fn sample_user_with_role(user_id: &str, role: &str) -> StoredUserAuthRecord { + StoredUserAuthRecord::new( + user_id.to_string(), + Some(format!("{user_id}@example.com")), + true, + user_id.to_string(), + Some("hash".to_string()), + role.to_string(), + "local".to_string(), + Some(json!(["openai"])), + Some(json!(["openai:chat"])), + Some(json!(["gpt-4.1"])), + true, + false, + Some(Utc::now()), + Some(Utc::now()), + ) + .expect("test user should build") +} + +fn sample_wallet(user_id: &str, balance: f64, gift_balance: f64) -> StoredWalletSnapshot { + StoredWalletSnapshot::new( + format!("wallet-{user_id}"), + Some(user_id.to_string()), + None, + balance, + gift_balance, + "finite".to_string(), + "USD".to_string(), + "active".to_string(), + balance.max(0.0), + 0.0, + 0.0, + 0.0, + 1_710_000_000, + ) + .expect("test wallet should build") +} + +async fn post_batch_action(client: &Client, gateway_url: &str, payload: Value) -> Response { + 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() + .await + .expect("batch request should complete") +} + +async fn wallet_detail(client: &Client, gateway_url: &str, user_id: &str) -> Value { + admin_headers(client.get(format!("{gateway_url}/api/admin/wallets/wallet-{user_id}"))) + .send() + .await + .expect("wallet lookup should complete") + .json() + .await + .expect("wallet response should parse") +} + +#[tokio::test] +async fn gateway_batches_wallet_addition_deduction_and_clamped_deduction_per_user() { + let state = AppState::new() + .expect("gateway should build") + .with_auth_users_for_tests([sample_user("user-1"), sample_user("user-2")]) + .with_auth_wallets_for_tests([ + sample_wallet("user-1", 10.0, 3.0), + sample_wallet("user-2", 2.0, 1.0), + ]); + let (gateway_url, gateway_handle) = start_server(build_router_with_state(state)).await; + let client = Client::new(); + let selection = json!({ "user_ids": ["user-1", "user-2"] }); + + let add_response = post_batch_action( + &client, + &gateway_url, + json!({ + "selection": selection.clone(), + "action": "adjust_wallet_balance", + "payload": { "operation": "add", "amount": 5.0 } + }), + ) + .await; + assert_eq!(add_response.status(), StatusCode::OK); + let add_result: Value = add_response.json().await.expect("response should parse"); + assert_eq!(add_result["success"], 2); + assert_eq!(add_result["failed"], 0); + assert_eq!(add_result["modified_fields"], json!(["wallet_balance"])); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 18.0 + ); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-2").await["balance"], + 8.0 + ); + + let deduct_response = post_batch_action( + &client, + &gateway_url, + json!({ + "selection": selection.clone(), + "action": "adjust_wallet_balance", + "payload": { "operation": "deduct", "amount": 4.0 } + }), + ) + .await; + assert_eq!(deduct_response.status(), StatusCode::OK); + let deduct_result: Value = deduct_response.json().await.expect("response should parse"); + assert_eq!(deduct_result["success"], 2); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 14.0 + ); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-2").await["balance"], + 4.0 + ); + + let over_deduct_response = post_batch_action( + &client, + &gateway_url, + json!({ + "selection": selection, + "action": "adjust_wallet_balance", + "payload": { "operation": "deduct", "amount": 100.0 } + }), + ) + .await; + assert_eq!(over_deduct_response.status(), StatusCode::OK); + let over_deduct_result: Value = over_deduct_response + .json() + .await + .expect("response should parse"); + assert_eq!(over_deduct_result["success"], 2); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 0.0 + ); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-2").await["balance"], + 0.0 + ); + + gateway_handle.abort(); +} + +#[tokio::test] +async fn gateway_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"; + let wallet_write_token = "ae-batch-users-wallet-write"; + let wallet_admin_token = "ae-batch-users-wallet-admin"; + let token_owner = sample_user_with_role("token-owner", "admin"); + let target_user = sample_user("user-1"); + let mut users_only = sample_management_token( + "token-users-write-only", + &token_owner.id, + &token_owner.username, + true, + ); + users_only.token.allowed_ips = None; + users_only.token.permissions = Some(json!(["admin:users:write"])); + let mut users_and_wallets = sample_management_token( + "token-users-and-wallet-write", + &token_owner.id, + &token_owner.username, + true, + ); + users_and_wallets.token.allowed_ips = None; + users_and_wallets.token.permissions = Some(json!(["admin:users:write", "admin:wallets:write"])); + let mut wallets_admin = sample_management_token( + "token-users-wallet-admin", + &token_owner.id, + &token_owner.username, + true, + ); + wallets_admin.token.allowed_ips = None; + wallets_admin.token.permissions = Some(json!(["admin:users:write", "admin:wallets:admin"])); + let token_repository = Arc::new(InMemoryManagementTokenRepository::seed_with_hashes( + vec![users_only, users_and_wallets, wallets_admin], + vec![ + ( + hash_management_token(users_write_token), + "token-users-write-only".to_string(), + ), + ( + hash_management_token(wallet_write_token), + "token-users-and-wallet-write".to_string(), + ), + ( + hash_management_token(wallet_admin_token), + "token-users-wallet-admin".to_string(), + ), + ], + )); + let user_repository = Arc::new(InMemoryUserReadRepository::seed_auth_users(vec![ + token_owner.clone(), + target_user.clone(), + ])); + let data = GatewayDataState::with_management_token_repository_for_tests(token_repository) + .with_user_reader(user_repository); + let state = AppState::new() + .expect("gateway should build") + .with_data_state_for_tests(data) + .with_auth_users_for_tests([token_owner, target_user]) + .with_auth_wallets_for_tests([sample_wallet("user-1", 10.0, 0.0)]); + let (gateway_url, gateway_handle) = start_server(build_router_with_state(state)).await; + let client = Client::new(); + let payload = json!({ + "selection": { "user_ids": ["user-1"] }, + "action": "adjust_wallet_balance", + "payload": { "operation": "add", "amount": 5.0 }, + "idempotency_key": "wallet-write-batch" + }); + + let denied = client + .post(format!("{gateway_url}/api/admin/users/batch-action")) + .header(crate::constants::GATEWAY_HEADER, "rust-phase3b") + .bearer_auth(users_write_token) + .json(&payload) + .send() + .await + .expect("users-only management token request should complete"); + assert_eq!(denied.status(), StatusCode::FORBIDDEN); + let denied_payload: Value = denied.json().await.expect("response should parse"); + assert_eq!( + denied_payload["required_permissions"], + json!(["admin:wallets:write", "admin:wallets:admin"]) + ); + assert_eq!(denied_payload["permission_mode"], "any_of"); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 10.0 + ); + + let allowed = client + .post(format!("{gateway_url}/api/admin/users/batch-action")) + .header(crate::constants::GATEWAY_HEADER, "rust-phase3b") + .bearer_auth(wallet_write_token) + .json(&payload) + .send() + .await + .expect("wallet-write management token request should complete"); + assert_eq!(allowed.status(), StatusCode::OK); + let result: Value = allowed.json().await.expect("response should parse"); + assert_eq!(result["success"], 1); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 15.0 + ); + + let allowed_with_wallet_admin = client + .post(format!("{gateway_url}/api/admin/users/batch-action")) + .header(crate::constants::GATEWAY_HEADER, "rust-phase3b") + .bearer_auth(wallet_admin_token) + .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"); + assert_eq!(allowed_with_wallet_admin.status(), StatusCode::OK); + let result: Value = allowed_with_wallet_admin + .json() + .await + .expect("response should parse"); + assert_eq!(result["success"], 1); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 20.0 + ); + + gateway_handle.abort(); +} + +#[tokio::test] +async fn gateway_reports_completed_uncertain_and_unprocessed_users_after_adjustment_error() { + let state = AppState::new() + .expect("gateway should build") + .with_auth_users_for_tests([ + sample_user("user-1"), + sample_user("user-2"), + sample_user("user-3"), + ]) + .with_auth_wallets_for_tests([ + sample_wallet("user-1", 10.0, 0.0), + sample_wallet("user-2", 20.0, 0.0), + sample_wallet("user-3", 30.0, 0.0), + ]) + .fail_auth_wallet_adjustment_for_tests("wallet-user-2"); + let (gateway_url, gateway_handle) = start_server(build_router_with_state(state)).await; + let client = Client::new(); + + let response = post_batch_action( + &client, + &gateway_url, + json!({ + "selection": { "user_ids": ["user-1", "user-2", "user-3"] }, + "action": "adjust_wallet_balance", + "payload": { "operation": "add", "amount": 5.0 } + }), + ) + .await; + assert_eq!(response.status(), StatusCode::OK); + let result: Value = response.json().await.expect("response should parse"); + assert_eq!(result["interrupted"], true); + assert_eq!(result["success"], 1); + assert_eq!(result["failed"], 2); + assert_eq!(result["completed_user_ids"], json!(["user-1"])); + assert_eq!(result["uncertain_user_ids"], json!(["user-2"])); + assert_eq!(result["unprocessed_user_ids"], json!(["user-3"])); + assert_eq!(result["failures"][0]["user_id"], "user-2"); + assert_eq!(result["failures"][1]["user_id"], "user-3"); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 15.0 + ); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-2").await["balance"], + 20.0 + ); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-3").await["balance"], + 30.0 + ); + + gateway_handle.abort(); +} + +#[tokio::test] +async fn gateway_reports_wallet_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": "adjust_wallet_balance", + "payload": { "operation": "add", "amount": 5.0 } + }), + ) + .await; + assert_eq!(response.status(), StatusCode::OK); + let result: Value = response.json().await.expect("response should parse"); + assert_eq!(result["interrupted"], true); + assert_eq!(result["completed_user_ids"], json!(["user-1"])); + assert_eq!(result["uncertain_user_ids"], json!([])); + assert_eq!(result["unprocessed_user_ids"], json!(["user-2", "user-3"])); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-1").await["balance"], + 15.0 + ); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-2").await["balance"], + 20.0 + ); + assert_eq!( + wallet_detail(&client, &gateway_url, "user-3").await["balance"], + 30.0 + ); + + gateway_handle.abort(); +} + +#[tokio::test] +async fn gateway_reports_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_floors_negative_balance_on_deduction() { + let state = AppState::new() + .expect("gateway should build") + .with_auth_users_for_tests([sample_user("user-negative"), sample_user("user-no-wallet")]) + .with_auth_wallets_for_tests([sample_wallet("user-negative", -2.0, 1.0)]); + let (gateway_url, gateway_handle) = start_server(build_router_with_state(state)).await; + let client = Client::new(); + + let response = post_batch_action( + &client, + &gateway_url, + json!({ + "selection": { "user_ids": ["user-negative", "user-no-wallet"] }, + "action": "adjust_wallet_balance", + "payload": { "operation": "deduct", "amount": 10.0 } + }), + ) + .await; + assert_eq!(response.status(), StatusCode::OK); + let result: Value = response.json().await.expect("response should parse"); + assert_eq!(result["success"], 1); + assert_eq!(result["failed"], 1); + assert_eq!(result["failures"][0]["user_id"], "user-no-wallet"); + assert_eq!(result["failures"][0]["reason"], "用户钱包不可用"); + + let wallet = wallet_detail(&client, &gateway_url, "user-negative").await; + assert_eq!(wallet["balance"], 0.0); + assert_eq!(wallet["total_adjusted"], 1.0); + + gateway_handle.abort(); +} + +#[tokio::test] +async fn gateway_rejects_zero_and_non_finite_batch_wallet_adjustments() { + let state = AppState::new() + .expect("gateway should build") + .with_auth_users_for_tests([sample_user("user-1")]) + .with_auth_wallets_for_tests([sample_wallet("user-1", 10.0, 0.0)]); + let (gateway_url, gateway_handle) = start_server(build_router_with_state(state)).await; + let client = Client::new(); + + let zero_response = post_batch_action( + &client, + &gateway_url, + json!({ + "selection": { "user_ids": ["user-1"] }, + "action": "adjust_wallet_balance", + "payload": { "operation": "add", "amount": 0.0 } + }), + ) + .await; + assert_eq!(zero_response.status(), StatusCode::BAD_REQUEST); + + let non_finite_response = admin_headers( + client.post(format!("{gateway_url}/api/admin/users/batch-action")), + ) + .header(reqwest::header::CONTENT_TYPE, "application/json") + .body( + r#"{"selection":{"user_ids":["user-1"]},"action":"adjust_wallet_balance","payload":{"operation":"add","amount":1e999}}"#, + ) + .send() + .await + .expect("non-finite amount request should complete"); + assert_eq!(non_finite_response.status(), StatusCode::BAD_REQUEST); + + gateway_handle.abort(); +} diff --git a/crates/aether-data/adapters/postgres/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 eb9311f8f..ddf829a32 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}, @@ -807,6 +810,44 @@ impl SqlxWalletRepository { } } +fn effective_wallet_adjustment_amount(input: &AdjustWalletBalanceInput, before_total: f64) -> f64 { + if input.clamp_deduction_to_available_balance && input.amount_usd < 0.0 { + if before_total < 0.0 { + // Clear legacy negative totals to zero and record the actual ledger delta. + -before_total + } else { + -(-input.amount_usd).min(before_total) + } + } else { + input.amount_usd + } +} + +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( @@ -4262,7 +4303,8 @@ RETURNING async fn adjust_wallet_balance( &self, input: AdjustWalletBalanceInput, - ) -> Result, DataLayerError> { + ) -> Result)>, DataLayerError> + { if !input.amount_usd.is_finite() || input.amount_usd == 0.0 { return Err(DataLayerError::InvalidInput( "adjustment amount must be finite and non-zero".to_string(), @@ -4271,6 +4313,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 @@ -4285,13 +4379,20 @@ SELECT CAST(total_recharged AS DOUBLE PRECISION) AS total_recharged, CAST(total_consumed AS DOUBLE PRECISION) AS total_consumed, CAST(total_refunded AS DOUBLE PRECISION) AS total_refunded, - CAST(total_adjusted AS DOUBLE PRECISION) AS total_adjusted + CAST(total_adjusted AS DOUBLE PRECISION) AS total_adjusted, + CAST(EXTRACT(EPOCH FROM updated_at) AS BIGINT) AS updated_at_unix_secs FROM wallets WHERE id = $1 + 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()? @@ -4312,17 +4413,40 @@ 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 { + 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; - if input.amount_usd > 0.0 { + if amount_usd > 0.0 { if input.balance_type.eq_ignore_ascii_case("gift") { - after_gift += input.amount_usd; + after_gift += amount_usd; } else { - after_recharge += input.amount_usd; + after_recharge += amount_usd; } } else { - let mut remaining = -input.amount_usd; + let mut remaining = -amount_usd; let consume_positive_bucket = |balance: &mut f64, to_consume: &mut f64| { if *to_consume <= 0.0 { return; @@ -4344,7 +4468,7 @@ FOR UPDATE } } let after_total = after_recharge + after_gift; - let after_total_adjusted = before_total_adjusted + input.amount_usd; + let after_total_adjusted = before_total_adjusted + amount_usd; if !after_recharge.is_finite() || !after_gift.is_finite() || !after_total.is_finite() @@ -4383,7 +4507,7 @@ RETURNING .bind(&input.wallet_id) .bind(after_recharge) .bind(after_gift) - .bind(input.amount_usd) + .bind(amount_usd) .fetch_one(&mut **tx) .await .map_postgres_err()?; @@ -4439,7 +4563,7 @@ VALUES ( ) .bind(&transaction_id) .bind(&input.wallet_id) - .bind(input.amount_usd) + .bind(amount_usd) .bind(before_total) .bind(after_total) .bind(before_recharge) @@ -4453,14 +4577,24 @@ 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, - StoredAdminWalletTransaction { + Some(StoredAdminWalletTransaction { id: transaction_id, wallet_id: input.wallet_id, category: "adjust".to_string(), reason_code: "adjust_admin".to_string(), - amount: input.amount_usd, + amount: amount_usd, balance_before: before_total, balance_after: after_total, recharge_balance_before: before_recharge, @@ -4474,13 +4608,244 @@ VALUES ( operator_email: None, description: Some(description), created_at_unix_ms: Some(created_at), - }, + }), ))) }) }) .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, @@ -8489,13 +8854,16 @@ VALUES ($1, $2, 'gift', 'gift_initial', $3, 0, $3, 0, 0, 0, $3, 'system_task', $ #[cfg(test)] mod tests { use aether_data_contracts::repository::wallet::{ - CreateManualWalletRechargeInput, CreditAdminPaymentOrderInput, ProcessPaymentCallbackInput, + AdjustWalletBalanceInBatchInput, AdjustWalletBalanceInput, + AdminUserWalletBalanceBatchUserOutcome, CreateManualWalletRechargeInput, + CreditAdminPaymentOrderInput, PrepareAdminUserWalletBalanceBatchInput, + PrepareAdminUserWalletBalanceBatchOutcome, ProcessPaymentCallbackInput, ProcessPaymentCallbackOutcome, RedeemWalletCodeInput, RedeemWalletCodeOutcome, WalletLookupKey, WalletMutationOutcome, WalletReadRepository, WalletWriteRepository, }; use sqlx::Row; - use super::SqlxWalletRepository; + use super::{effective_wallet_adjustment_amount, SqlxWalletRepository}; use crate::{PostgresPoolConfig, PostgresPoolFactory}; #[test] @@ -8575,6 +8943,7 @@ mod tests { "payment_orders", "payment_callbacks", "wallet_transactions", + "admin_user_wallet_balance_batches", "user_plan_entitlements", "redeem_code_batches", "redeem_codes", @@ -9221,6 +9590,264 @@ mod tests { pool.close().await; } + #[test] + fn bulk_adjustment_clamp_is_opt_in_and_uses_available_total() { + let input = AdjustWalletBalanceInput { + wallet_id: "wallet-1".to_string(), + amount_usd: -100.0, + balance_type: "recharge".to_string(), + operator_id: None, + description: None, + clamp_deduction_to_available_balance: true, + 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); + assert_eq!(effective_wallet_adjustment_amount(&input, -1.0), 1.0); + + let legacy_input = AdjustWalletBalanceInput { + clamp_deduction_to_available_balance: false, + ..input + }; + assert_eq!( + effective_wallet_adjustment_amount(&legacy_input, 13.0), + -100.0 + ); + } + + #[tokio::test] + #[ignore = "requires AETHER_TEST_DATABASE_URL and PostgreSQL bootstrap schema"] + async fn live_bulk_wallet_adjustment_persists_actual_delta_and_skips_zero_ledger() { + let pool = isolated_wallet_test_pool().await; + let (wallet_id, _) = seed_wallet(&pool).await; + let repository = SqlxWalletRepository::new(pool.clone()); + + let (wallet, transaction) = repository + .adjust_wallet_balance(AdjustWalletBalanceInput { + wallet_id: wallet_id.clone(), + amount_usd: -100.0, + balance_type: "recharge".to_string(), + operator_id: Some("admin-user".to_string()), + description: Some("bulk deduction".to_string()), + clamp_deduction_to_available_balance: true, + batch_context: None, + }) + .await + .expect("bulk adjustment should succeed") + .expect("wallet should exist"); + let transaction = + transaction.expect("positive available balance should create a ledger row"); + assert_eq!(transaction.amount, -13.0); + assert_eq!(transaction.balance_before, 13.0); + assert_eq!(transaction.balance_after, 0.0); + assert_eq!(wallet.balance + wallet.gift_balance, 0.0); + let persisted_amount: f64 = sqlx::query_scalar( + "SELECT amount::double precision FROM wallet_transactions WHERE id = $1", + ) + .bind(&transaction.id) + .fetch_one(&pool) + .await + .expect("ledger should store the effective deduction"); + assert_eq!(persisted_amount, -13.0); + + let (wallet, transaction) = repository + .adjust_wallet_balance(AdjustWalletBalanceInput { + wallet_id: wallet_id.clone(), + amount_usd: -1.0, + balance_type: "recharge".to_string(), + operator_id: Some("admin-user".to_string()), + description: Some("bulk deduction at zero".to_string()), + clamp_deduction_to_available_balance: true, + batch_context: None, + }) + .await + .expect("zero-balance adjustment should succeed") + .expect("wallet should still exist"); + assert_eq!(wallet.balance + wallet.gift_balance, 0.0); + assert!( + transaction.is_none(), + "zero effective delta must not create a ledger row" + ); + let transaction_count: i64 = + sqlx::query_scalar("SELECT COUNT(*) FROM wallet_transactions WHERE wallet_id = $1") + .bind(&wallet_id) + .fetch_one(&pool) + .await + .expect("ledger row count should be readable"); + assert_eq!(transaction_count, 1); + + sqlx::query("UPDATE wallets SET balance = -2, gift_balance = 1 WHERE id = $1") + .bind(&wallet_id) + .execute(&pool) + .await + .expect("legacy negative wallet balance should be seeded"); + let (wallet, transaction) = repository + .adjust_wallet_balance(AdjustWalletBalanceInput { + wallet_id: wallet_id.clone(), + amount_usd: -1.0, + balance_type: "recharge".to_string(), + operator_id: Some("admin-user".to_string()), + description: Some("bulk deduction floors a negative balance".to_string()), + clamp_deduction_to_available_balance: true, + batch_context: None, + }) + .await + .expect("legacy negative wallet should be floored at zero") + .expect("wallet should still exist"); + assert_eq!(wallet.balance + wallet.gift_balance, 0.0); + let transaction = transaction.expect("negative balance correction should be ledgered"); + assert_eq!(transaction.amount, 1.0); + assert_eq!(transaction.balance_before, -1.0); + assert_eq!(transaction.balance_after, 0.0); + let persisted_correction_amount: f64 = sqlx::query_scalar( + "SELECT amount::double precision FROM wallet_transactions WHERE id = $1", + ) + .bind(&transaction.id) + .fetch_one(&pool) + .await + .expect("ledger should store the negative-balance correction"); + assert_eq!(persisted_correction_amount, 1.0); + let transaction_count: i64 = + sqlx::query_scalar("SELECT COUNT(*) FROM wallet_transactions WHERE wallet_id = $1") + .bind(&wallet_id) + .fetch_one(&pool) + .await + .expect("ledger row count should be readable"); + assert_eq!(transaction_count, 2); + 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 25a7e8ea5..01262885f 100644 --- a/crates/aether-data/contracts/src/repository/wallet/types.rs +++ b/crates/aether-data/contracts/src/repository/wallet/types.rs @@ -2462,6 +2462,59 @@ pub struct AdjustWalletBalanceInput { pub balance_type: String, pub operator_id: Option, 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)] @@ -2893,7 +2946,51 @@ pub trait WalletWriteRepository: Send + Sync { async fn adjust_wallet_balance( &self, input: AdjustWalletBalanceInput, - ) -> Result, crate::DataLayerError>; + ) -> Result< + Option<(StoredWalletSnapshot, Option)>, + 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, 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/lifecycle/migrate/tests.rs b/crates/aether-data/runtime/src/lifecycle/migrate/tests.rs index 02fa33e24..8d375d1f5 100644 --- a/crates/aether-data/runtime/src/lifecycle/migrate/tests.rs +++ b/crates/aether-data/runtime/src/lifecycle/migrate/tests.rs @@ -1575,6 +1575,7 @@ fn pending_migrations_from_applied_skips_versions_already_applied() { 20260901000000, 20260903000000, 20260908000000, + 20260923000000, ] ); } diff --git a/crates/aether-data/runtime/src/repository/wallet/memory.rs b/crates/aether-data/runtime/src/repository/wallet/memory.rs index 08df3efc8..73a47063a 100644 --- a/crates/aether-data/runtime/src/repository/wallet/memory.rs +++ b/crates/aether-data/runtime/src/repository/wallet/memory.rs @@ -2322,8 +2322,13 @@ impl WalletWriteRepository for InMemoryWalletRepository { async fn adjust_wallet_balance( &self, _input: AdjustWalletBalanceInput, - ) -> Result, DataLayerError> - { + ) -> Result< + Option<( + StoredWalletSnapshot, + Option, + )>, + DataLayerError, + > { Ok(None) } diff --git a/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/__tests__/users.spec.ts b/frontend/src/api/__tests__/users.spec.ts index 3f2f8d486..f4dc72e52 100644 --- a/frontend/src/api/__tests__/users.spec.ts +++ b/frontend/src/api/__tests__/users.spec.ts @@ -17,7 +17,7 @@ vi.mock('@/utils/cache', () => ({ cachedRequest: cachedRequestMock, })) -import { usersApi } from '@/api/users' +import { buildUserBatchBalanceAdjustmentPayload, usersApi } from '@/api/users' describe('usersApi admin list query', () => { beforeEach(() => { @@ -104,3 +104,25 @@ describe('usersApi admin list query', () => { expect(getMock).toHaveBeenCalledWith('/api/admin/users/target-user/api-keys') }) }) + +describe('user batch wallet balance payload', () => { + it('builds an addition payload from a valid positive amount', () => { + expect(buildUserBatchBalanceAdjustmentPayload('add', '12.5')).toEqual({ + operation: 'add', + amount: 12.5, + }) + }) + + it('builds a deduction payload from a valid positive amount', () => { + expect(buildUserBatchBalanceAdjustmentPayload('deduct', 4)).toEqual({ + operation: 'deduct', + amount: 4, + }) + }) + + it('rejects zero, blank, and non-finite amounts', () => { + expect(buildUserBatchBalanceAdjustmentPayload('add', '0')).toBeNull() + expect(buildUserBatchBalanceAdjustmentPayload('deduct', '')).toBeNull() + expect(buildUserBatchBalanceAdjustmentPayload('deduct', '1e999')).toBeNull() + }) +}) diff --git a/frontend/src/api/users.ts b/frontend/src/api/users.ts index 8a9b8d6a7..0902b7775 100644 --- a/frontend/src/api/users.ts +++ b/frontend/src/api/users.ts @@ -120,9 +120,28 @@ export interface UserBatchRolePayload { role: UserRole } -export type UserBatchAction = 'enable' | 'disable' | 'update_access_control' | 'update_role' +export type UserBatchBalanceOperation = 'add' | 'deduct' -export type UserBatchActionPayload = UserBatchAccessControlPayload | UserBatchRolePayload +export interface UserBatchBalanceAdjustmentPayload { + operation: UserBatchBalanceOperation + amount: number +} + +export function buildUserBatchBalanceAdjustmentPayload( + operation: UserBatchBalanceOperation, + amountInput: string | number, +): UserBatchBalanceAdjustmentPayload | null { + const amount = typeof amountInput === 'number' ? amountInput : Number(amountInput.trim()) + if (!Number.isFinite(amount) || amount <= 0) return null + return { operation, amount } +} + +export type UserBatchAction = 'enable' | 'disable' | 'update_access_control' | 'update_role' | 'adjust_wallet_balance' + +export type UserBatchActionPayload = + | UserBatchAccessControlPayload + | UserBatchRolePayload + | UserBatchBalanceAdjustmentPayload export interface UserBatchToggleActionRequest { selection: UserBatchSelection @@ -142,10 +161,18 @@ export interface UserBatchRoleActionRequest { payload: UserBatchRolePayload } +export interface UserBatchBalanceActionRequest { + selection: UserBatchSelection + action: 'adjust_wallet_balance' + payload: UserBatchBalanceAdjustmentPayload + idempotency_key: string +} + export type UserBatchActionRequest = | UserBatchToggleActionRequest | UserBatchAccessControlActionRequest | UserBatchRoleActionRequest + | UserBatchBalanceActionRequest export interface UserBatchActionFailure { user_id: string @@ -157,6 +184,10 @@ export interface UserBatchActionResponse { success: number failed: number failures: UserBatchActionFailure[] + interrupted?: boolean + completed_user_ids?: string[] + uncertain_user_ids?: string[] + unprocessed_user_ids?: string[] warnings?: UserBatchSelectionWarning[] action?: string modified_fields?: string[] diff --git a/frontend/src/features/users/components/UserBatchActionDialog.vue b/frontend/src/features/users/components/UserBatchActionDialog.vue index fa87b7748..ac9175b37 100644 --- a/frontend/src/features/users/components/UserBatchActionDialog.vue +++ b/frontend/src/features/users/components/UserBatchActionDialog.vue @@ -2,7 +2,7 @@ +
+
+ +
+ + +
+
+ +

+ {{ legacyT('请输入大于 0 的有限金额') }} +

+

+ {{ legacyT('扣减超过单个用户可用余额时,该用户余额将归零。') }} +

+
+ + + + () const usersStore = useUsersStore() +const authStore = useAuthStore() +const walletRetryCoordinator = createUserBatchWalletRetryCoordinator({ + scope: () => authStore.user?.id ?? null, +}) const { success, warning, error } = useToast() const { legacyT, locale } = useI18n() const selectedAction = ref('enable') const targetRole = ref('user') const quotaMode = ref('skip') +const balanceOperation = ref('add') +const balanceAmount = ref('') const selectedGroupIds = ref([]) const previewLoading = ref(false) const previewItems = ref([]) 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) -const canExecute = computed(() => hasAnyTarget.value && !previewLoading.value && !executing.value) +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 && !pendingWalletReadError.value && !walletRequestMismatch.value)) +)) const selectedActionLabel = computed(() => ( USER_BATCH_ACTION_OPTIONS.find((action) => action.value === selectedAction.value)?.label ?? '批量操作' )) @@ -143,13 +274,25 @@ 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) { + return legacyT( + `批量操作中断:成功 ${lastResult.value.success} 个,结果待确认 ${lastResult.value.uncertain_user_ids?.length ?? 0} 个,尚未执行 ${lastResult.value.unprocessed_user_ids?.length ?? 0} 个`, + ) + } return legacyT(`成功 ${lastResult.value.success} 个,失败 ${lastResult.value.failed} 个`) }) const lastResultFailuresLabel = computed(() => { if (!lastResult.value || lastResult.value.failures.length === 0) return '' - const failures = lastResult.value.failures.slice(0, 3).map((item) => `${item.user_id} ${item.reason}`).join(locale.value === 'en-US' ? '; ' : ';') + const failures = lastResult.value.failures.slice(0, 3) + .map((item) => `${item.user_id} ${legacyT(item.reason)}`) + .join(locale.value === 'en-US' ? '; ' : ';') return locale.value === 'en-US' ? `: ${failures}` : `:${failures}` }) @@ -158,8 +301,15 @@ watch( (open) => { if (!open) return resetLocalState() + refreshPendingWalletBatch() void resolvePreview() }, + { immediate: true }, +) + +watch( + () => authStore.user?.id, + () => refreshPendingWalletBatch(), ) watch( @@ -177,10 +327,22 @@ function resetLocalState(): void { selectedAction.value = 'enable' targetRole.value = 'user' quotaMode.value = 'skip' + balanceOperation.value = 'add' + balanceAmount.value = '' selectedGroupIds.value = [] 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) { @@ -227,7 +389,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) { @@ -235,6 +398,12 @@ async function executeBatchAction(): Promise { return } request = { selection, action: 'update_access_control', payload } + } else if (selectedAction.value === 'adjust_wallet_balance') { + if (walletAdjustmentRequest.value === null) { + warning(legacyT('请输入大于 0 的有限金额')) + return + } + walletRequest = walletAdjustmentRequest.value } else if (selectedAction.value === 'update_role') { request = { selection, action: 'update_role', payload: buildRolePayload() } } else { @@ -243,8 +412,22 @@ 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) { + warning(`${lastResultLabel.value};${legacyT('请核对余额后再重试,勿直接重试整批')}`) + emit('completed', result) + return + } const message = legacyT(`批量操作完成:成功 ${result.success} 个,失败 ${result.failed} 个`) if (result.failed > 0) { warning(message) @@ -253,9 +436,71 @@ async function executeBatchAction(): Promise { } emit('completed', result) } catch (err) { - error(parseApiError(err, '批量操作失败'), legacyT('批量操作失败')) + if (walletRequest) { + refreshPendingWalletBatch() + if (err instanceof WalletIdempotencyPersistenceUnavailableError) { + warning(legacyT('浏览器无法安全保存钱包批量请求,本次请求未发送。')) + } else if (err instanceof WalletIdempotencyRequestInProgressError) { + warning(legacyT('另一个标签页正在处理钱包批量调整,请稍后刷新状态再试。此次未发送新请求。')) + } else if (err instanceof WalletIdempotencyCoordinationUnavailableError) { + warning(legacyT('当前浏览器无法保护跨标签页的钱包批量请求,请使用支持此功能的浏览器。请求未发送。')) + } else if ( + err instanceof WalletIdempotencyUnavailableError + || err instanceof WalletIdempotencyScopeUnavailableError + || 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 WalletIdempotencyRequestInProgressError) { + warning(legacyT('另一个标签页正在处理钱包批量调整,请稍后刷新状态再试。此次未发送新请求。')) + } else if (err instanceof WalletIdempotencyCoordinationUnavailableError) { + 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/components/UserBatchResultSummary.vue b/frontend/src/features/users/components/UserBatchResultSummary.vue index 30ec4c956..a3a68f126 100644 --- a/frontend/src/features/users/components/UserBatchResultSummary.vue +++ b/frontend/src/features/users/components/UserBatchResultSummary.vue @@ -7,11 +7,26 @@ {{ failuresLabel }} +
+ {{ legacyT('查看批次中断详情') }} +
+

{{ legacyT('已完成用户 ID') }}:{{ userIds(result.completed_user_ids) }}

+

{{ legacyT('结果待核对用户 ID') }}:{{ userIds(result.uncertain_user_ids) }}

+

{{ legacyT('尚未执行用户 ID') }}:{{ userIds(result.unprocessed_user_ids) }}

+
+