Files
Genarrative/server-rs/crates/spacetime-module/src/payment_storage.rs
T
kdletters 4b529a8952
Project CI / AI game creator shell Rust crates (pull_request) Successful in 2m59s
Project CI / Backend tests (pull_request) Failing after 4m8s
Project CI / AI game creator shell Rust lane 1/2 (pull_request) Successful in 5m9s
Project CI / AI game creator shell Rust lane 2/2 (pull_request) Successful in 6m20s
Project CI / Frontend tests (pull_request) Successful in 3m17s
Project CI / AI game creator shell web tests (pull_request) Successful in 2m48s
Project CI / Repository checks (pull_request) Failing after 4m15s
Project CI / Native shell tests (pull_request) Successful in 6m52s
新增支付服务接入与订单收银台
新增支付应用、API Key 与外部产品订单接口

接入微信 Native 下单、二维码收银台与支付结果回调

新增后台支付订单与回调投递页面

补齐 SpacetimeDB 支付表、生成绑定、OpenAPI 与使用说明
2026-10-03 17:18:13 +08:00

1200 lines
40 KiB
Rust

use crate::*;
#[spacetimedb::table(
accessor = payment_app,
index(accessor = by_payment_app_owner_user_id, btree(columns = [owner_user_id]))
)]
#[derive(Clone)]
pub struct PaymentApp {
#[primary_key]
pub app_id: String,
pub owner_user_id: String,
pub name: String,
pub callback_url: Option<String>,
pub enabled: bool,
pub created_at: Timestamp,
pub updated_at: Timestamp,
}
#[spacetimedb::table(
accessor = payment_api_key,
index(accessor = by_payment_api_key_app_id, btree(columns = [app_id])),
index(accessor = by_payment_api_key_owner_user_id, btree(columns = [owner_user_id]))
)]
pub struct PaymentApiKey {
#[primary_key]
pub key_id: String,
pub app_id: String,
pub owner_user_id: String,
pub name: String,
pub key_prefix: String,
#[unique]
pub key_hash: String,
pub scopes_json: String,
pub created_at: Timestamp,
pub last_used_at: Option<Timestamp>,
pub revoked_at: Option<Timestamp>,
pub updated_at: Timestamp,
}
#[spacetimedb::table(
accessor = payment_order,
index(accessor = by_payment_order_app_id, btree(columns = [app_id])),
index(accessor = by_payment_order_owner_user_id, btree(columns = [owner_user_id])),
index(accessor = by_payment_order_status, btree(columns = [status]))
)]
#[derive(Clone)]
pub struct PaymentOrder {
#[primary_key]
pub order_id: String,
pub app_id: String,
pub owner_user_id: String,
pub merchant_order_id: String,
#[unique]
pub idempotency_lookup: String,
#[unique]
pub merchant_order_lookup: String,
pub title: String,
pub items_json: String,
pub amount_cents: u64,
pub currency: String,
pub provider: PaymentProvider,
pub provider_trade_no: Option<String>,
#[unique]
pub checkout_token: String,
pub status: PaymentOrderStatus,
pub created_at: Timestamp,
pub expires_at: Timestamp,
pub paid_at: Option<Timestamp>,
pub updated_at: Timestamp,
pub provider_code_url: Option<String>,
pub callback_url: Option<String>,
pub callback_signing_key_hash: Option<String>,
#[default(None::<Timestamp>)]
pub provider_attempted_at: Option<Timestamp>,
}
#[spacetimedb::table(
accessor = payment_webhook_delivery,
index(accessor = by_payment_webhook_delivery_status_available, btree(columns = [status, available_at])),
index(accessor = by_payment_webhook_delivery_owner_user_id, btree(columns = [owner_user_id])),
index(accessor = by_payment_webhook_delivery_order_id, btree(columns = [order_id]))
)]
#[derive(Clone)]
pub struct PaymentWebhookDelivery {
#[primary_key]
pub delivery_id: String,
pub order_id: String,
pub app_id: String,
pub owner_user_id: String,
pub callback_url: String,
pub payload_json: String,
pub signature: String,
pub status: String,
pub attempt_count: u32,
pub available_at: Timestamp,
pub lease_worker_id: Option<String>,
pub lease_expires_at: Option<Timestamp>,
pub last_error_message: Option<String>,
pub created_at: Timestamp,
pub updated_at: Timestamp,
}
const PAYMENT_WEBHOOK_STATUS_PENDING: &str = "pending";
const PAYMENT_WEBHOOK_STATUS_PROCESSING: &str = "processing";
const PAYMENT_WEBHOOK_STATUS_DELIVERED: &str = "delivered";
const PAYMENT_WEBHOOK_STATUS_FAILED: &str = "failed";
#[spacetimedb::procedure]
pub fn create_payment_app_and_return(
ctx: &mut ProcedureContext,
input: PaymentAppCreateInput,
) -> PaymentAppProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::require_editor_generation_runtime_service_identity(tx, caller)?;
create_payment_app(tx, input.clone())
}) {
Ok(app) => PaymentAppProcedureResult {
ok: true,
app: Some(app),
key: None,
apps: Vec::new(),
keys: Vec::new(),
error_message: None,
},
Err(error_message) => payment_app_error(error_message),
}
}
#[spacetimedb::procedure]
pub fn list_payment_apps_and_return(
ctx: &mut ProcedureContext,
input: PaymentAppListInput,
) -> PaymentAppProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::require_editor_generation_runtime_service_identity(tx, caller)?;
list_payment_apps(tx, input.clone())
}) {
Ok(apps) => PaymentAppProcedureResult {
ok: true,
app: None,
key: None,
apps,
keys: Vec::new(),
error_message: None,
},
Err(error_message) => payment_app_error(error_message),
}
}
#[spacetimedb::procedure]
pub fn update_payment_app_and_return(
ctx: &mut ProcedureContext,
input: PaymentAppUpdateInput,
) -> PaymentAppProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::require_editor_generation_runtime_service_identity(tx, caller)?;
update_payment_app(tx, input.clone())
}) {
Ok(app) => PaymentAppProcedureResult {
ok: true,
app: Some(app),
key: None,
apps: Vec::new(),
keys: Vec::new(),
error_message: None,
},
Err(error_message) => payment_app_error(error_message),
}
}
#[spacetimedb::procedure]
pub fn list_payment_api_keys_and_return(
ctx: &mut ProcedureContext,
input: PaymentApiKeyListInput,
) -> PaymentAppProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::require_editor_generation_runtime_service_identity(tx, caller)?;
list_payment_api_keys(tx, input.clone())
}) {
Ok(keys) => PaymentAppProcedureResult {
ok: true,
app: None,
key: None,
apps: Vec::new(),
keys,
error_message: None,
},
Err(error_message) => payment_app_error(error_message),
}
}
#[spacetimedb::procedure]
pub fn create_payment_api_key_and_return(
ctx: &mut ProcedureContext,
input: PaymentApiKeyCreateInput,
) -> PaymentAppProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::require_editor_generation_runtime_service_identity(tx, caller)?;
create_payment_api_key(tx, input.clone())
}) {
Ok(key) => PaymentAppProcedureResult {
ok: true,
app: None,
key: Some(key),
apps: Vec::new(),
keys: Vec::new(),
error_message: None,
},
Err(error_message) => payment_app_error(error_message),
}
}
#[spacetimedb::procedure]
pub fn revoke_payment_api_key_and_return(
ctx: &mut ProcedureContext,
input: PaymentApiKeyRevokeInput,
) -> PaymentAppProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::require_editor_generation_runtime_service_identity(tx, caller)?;
revoke_payment_api_key(tx, input.clone())
}) {
Ok(key) => PaymentAppProcedureResult {
ok: true,
app: None,
key: Some(key),
apps: Vec::new(),
keys: Vec::new(),
error_message: None,
},
Err(error_message) => payment_app_error(error_message),
}
}
#[spacetimedb::procedure]
pub fn authenticate_payment_api_key_and_return(
ctx: &mut ProcedureContext,
input: PaymentApiKeyAuthenticateInput,
) -> PaymentAppProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::require_editor_generation_runtime_service_identity(tx, caller)?;
authenticate_payment_api_key(tx, input.clone())
}) {
Ok((app, key)) => PaymentAppProcedureResult {
ok: true,
app: Some(app),
key: Some(key),
apps: Vec::new(),
keys: Vec::new(),
error_message: None,
},
Err(error_message) => payment_app_error(error_message),
}
}
#[spacetimedb::procedure]
pub fn create_payment_order_and_return(
ctx: &mut ProcedureContext,
input: PaymentOrderCreateInput,
) -> PaymentOrderProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::require_editor_generation_runtime_service_identity(tx, caller)?;
create_payment_order(tx, input.clone())
}) {
Ok(order) => payment_order_single_ok(order),
Err(error_message) => payment_order_error(error_message),
}
}
#[spacetimedb::procedure]
pub fn get_payment_order_and_return(
ctx: &mut ProcedureContext,
input: PaymentOrderGetInput,
) -> PaymentOrderProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::require_editor_generation_runtime_service_identity(tx, caller)?;
get_payment_order(tx, input.clone())
}) {
Ok(order) => payment_order_single_ok(order),
Err(error_message) => payment_order_error(error_message),
}
}
#[spacetimedb::procedure]
pub fn get_payment_order_by_checkout_token_and_return(
ctx: &mut ProcedureContext,
input: PaymentOrderCheckoutGetInput,
) -> PaymentOrderProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::require_editor_generation_runtime_service_identity(tx, caller)?;
get_payment_order_by_checkout_token(tx, input.clone())
}) {
Ok(order) => payment_order_single_ok(order),
Err(error_message) => payment_order_error(error_message),
}
}
#[spacetimedb::procedure]
pub fn list_payment_orders_and_return(
ctx: &mut ProcedureContext,
input: PaymentOrderListInput,
) -> PaymentOrderProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::require_editor_generation_runtime_service_identity(tx, caller)?;
list_payment_orders(tx, input.clone())
}) {
Ok(orders) => PaymentOrderProcedureResult {
ok: true,
order: None,
orders,
error_message: None,
},
Err(error_message) => payment_order_error(error_message),
}
}
#[spacetimedb::procedure]
pub fn mark_payment_order_paid_and_return(
ctx: &mut ProcedureContext,
input: PaymentOrderMarkPaidInput,
) -> PaymentOrderProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::require_editor_generation_runtime_service_identity(tx, caller)?;
mark_payment_order_paid(tx, input.clone())
}) {
Ok(order) => payment_order_single_ok(order),
Err(error_message) => payment_order_error(error_message),
}
}
#[spacetimedb::procedure]
pub fn set_payment_order_provider_code_and_return(
ctx: &mut ProcedureContext,
input: PaymentOrderProviderCodeInput,
) -> PaymentOrderProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::require_editor_generation_runtime_service_identity(tx, caller)?;
set_payment_order_provider_code(tx, input.clone())
}) {
Ok(order) => payment_order_single_ok(order),
Err(error_message) => payment_order_error(error_message),
}
}
#[spacetimedb::procedure]
pub fn mark_payment_order_provider_attempted_and_return(
ctx: &mut ProcedureContext,
input: PaymentOrderProviderAttemptInput,
) -> PaymentOrderProcedureResult {
match ctx.try_with_tx(|tx| mark_payment_order_provider_attempted(tx, input.clone())) {
Ok(order) => payment_order_single_ok(order),
Err(error_message) => payment_order_error(error_message),
}
}
#[spacetimedb::procedure]
pub fn enqueue_payment_webhook_and_return(
ctx: &mut ProcedureContext,
input: PaymentWebhookEnqueueInput,
) -> PaymentWebhookProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::require_editor_generation_runtime_service_identity(tx, caller)?;
enqueue_payment_webhook(tx, input.clone())
}) {
Ok(delivery) => PaymentWebhookProcedureResult {
ok: true,
delivery,
deliveries: Vec::new(),
error_message: None,
},
Err(error_message) => payment_webhook_error(error_message),
}
}
#[spacetimedb::procedure]
pub fn claim_payment_webhooks_and_return(
ctx: &mut ProcedureContext,
input: PaymentWebhookClaimInput,
) -> PaymentWebhookProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::require_editor_generation_runtime_service_identity(tx, caller)?;
claim_payment_webhooks(tx, input.clone())
}) {
Ok(deliveries) => PaymentWebhookProcedureResult {
ok: true,
delivery: None,
deliveries,
error_message: None,
},
Err(error_message) => payment_webhook_error(error_message),
}
}
#[spacetimedb::procedure]
pub fn complete_payment_webhook_and_return(
ctx: &mut ProcedureContext,
input: PaymentWebhookCompleteInput,
) -> PaymentWebhookProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::require_editor_generation_runtime_service_identity(tx, caller)?;
complete_payment_webhook(tx, input.clone())
}) {
Ok(delivery) => PaymentWebhookProcedureResult {
ok: true,
delivery: Some(delivery),
deliveries: Vec::new(),
error_message: None,
},
Err(error_message) => payment_webhook_error(error_message),
}
}
fn create_payment_app(
ctx: &ReducerContext,
input: PaymentAppCreateInput,
) -> Result<PaymentAppSnapshot, String> {
let input = build_payment_app_create_input(
input.app_id,
input.owner_user_id,
input.name,
input.callback_url,
input.now_micros,
)?;
if ctx.db.payment_app().app_id().find(&input.app_id).is_some() {
return Err("支付应用 ID 已存在".to_string());
}
let now = Timestamp::from_micros_since_unix_epoch(input.now_micros);
ctx.db.payment_app().insert(PaymentApp {
app_id: input.app_id.clone(),
owner_user_id: input.owner_user_id,
name: input.name,
callback_url: input.callback_url,
enabled: true,
created_at: now,
updated_at: now,
});
ctx.db
.payment_app()
.app_id()
.find(&input.app_id)
.map(|row| payment_app_snapshot(&row))
.ok_or_else(|| "支付应用创建后读取失败".to_string())
}
fn list_payment_apps(
ctx: &ReducerContext,
input: PaymentAppListInput,
) -> Result<Vec<PaymentAppSnapshot>, String> {
let owner_user_id = required(input.owner_user_id, "payment_app.owner_user_id")?;
let mut apps = ctx
.db
.payment_app()
.by_payment_app_owner_user_id()
.filter(&owner_user_id)
.map(|row| payment_app_snapshot(&row))
.collect::<Vec<_>>();
apps.sort_by(|left, right| right.created_at_micros.cmp(&left.created_at_micros));
Ok(apps)
}
fn update_payment_app(
ctx: &ReducerContext,
input: PaymentAppUpdateInput,
) -> Result<PaymentAppSnapshot, String> {
let app_id = required(input.app_id, "payment_app.app_id")?;
let owner_user_id = required(input.owner_user_id, "payment_app.owner_user_id")?;
let row = ctx
.db
.payment_app()
.app_id()
.find(&app_id)
.ok_or_else(|| "支付应用不存在".to_string())?;
if row.owner_user_id != owner_user_id {
return Err("无权访问支付应用".to_string());
}
let now = Timestamp::from_micros_since_unix_epoch(input.updated_at_micros);
ctx.db.payment_app().app_id().delete(&app_id);
ctx.db.payment_app().insert(PaymentApp {
enabled: input.enabled,
updated_at: now,
..row
});
ctx.db
.payment_app()
.app_id()
.find(&app_id)
.map(|row| payment_app_snapshot(&row))
.ok_or_else(|| "支付应用更新后读取失败".to_string())
}
fn list_payment_api_keys(
ctx: &ReducerContext,
input: PaymentApiKeyListInput,
) -> Result<Vec<PaymentApiKeySnapshot>, String> {
let owner_user_id = required(input.owner_user_id, "payment_api_key.owner_user_id")?;
let mut keys = ctx
.db
.payment_api_key()
.by_payment_api_key_owner_user_id()
.filter(&owner_user_id)
.map(|row| payment_api_key_snapshot(&row))
.collect::<Vec<_>>();
keys.sort_by(|left, right| right.created_at_micros.cmp(&left.created_at_micros));
Ok(keys)
}
fn create_payment_api_key(
ctx: &ReducerContext,
input: PaymentApiKeyCreateInput,
) -> Result<PaymentApiKeySnapshot, String> {
let input = build_payment_api_key_create_input(
input.key_id,
input.app_id,
input.owner_user_id,
input.name,
input.key_prefix,
input.key_hash,
input.scopes_json,
input.now_micros,
)?;
let app = ctx
.db
.payment_app()
.app_id()
.find(&input.app_id)
.ok_or_else(|| "支付应用不存在".to_string())?;
if app.owner_user_id != input.owner_user_id || !app.enabled {
return Err("支付应用不可用".to_string());
}
if ctx
.db
.payment_api_key()
.key_hash()
.find(&input.key_hash)
.is_some()
{
return Err("支付 API Key 摘要已存在".to_string());
}
let now = Timestamp::from_micros_since_unix_epoch(input.now_micros);
ctx.db.payment_api_key().insert(PaymentApiKey {
key_id: input.key_id.clone(),
app_id: input.app_id,
owner_user_id: input.owner_user_id,
name: input.name,
key_prefix: input.key_prefix,
key_hash: input.key_hash,
scopes_json: input.scopes_json,
created_at: now,
last_used_at: None,
revoked_at: None,
updated_at: now,
});
ctx.db
.payment_api_key()
.key_id()
.find(&input.key_id)
.map(|row| payment_api_key_snapshot(&row))
.ok_or_else(|| "支付 API Key 创建后读取失败".to_string())
}
fn revoke_payment_api_key(
ctx: &ReducerContext,
input: PaymentApiKeyRevokeInput,
) -> Result<PaymentApiKeySnapshot, String> {
let key_id = required(input.key_id, "payment_api_key.key_id")?;
let owner_user_id = required(input.owner_user_id, "payment_api_key.owner_user_id")?;
let row = ctx
.db
.payment_api_key()
.key_id()
.find(&key_id)
.ok_or_else(|| "支付 API Key 不存在".to_string())?;
if row.owner_user_id != owner_user_id {
return Err("无权访问支付 API Key".to_string());
}
let now = Timestamp::from_micros_since_unix_epoch(input.revoked_at_micros);
ctx.db.payment_api_key().key_id().delete(&key_id);
ctx.db.payment_api_key().insert(PaymentApiKey {
revoked_at: row.revoked_at.or(Some(now)),
updated_at: now,
..row
});
ctx.db
.payment_api_key()
.key_id()
.find(&key_id)
.map(|row| payment_api_key_snapshot(&row))
.ok_or_else(|| "支付 API Key 撤销后读取失败".to_string())
}
fn authenticate_payment_api_key(
ctx: &ReducerContext,
input: PaymentApiKeyAuthenticateInput,
) -> Result<(PaymentAppSnapshot, PaymentApiKeySnapshot), String> {
let key_hash = required(input.key_hash, "payment_api_key.key_hash")?;
let key = ctx
.db
.payment_api_key()
.key_hash()
.find(&key_hash)
.ok_or_else(|| "支付 API Key 不存在或已失效".to_string())?;
if key.revoked_at.is_some() {
return Err("支付 API Key 不存在或已失效".to_string());
}
let app = ctx
.db
.payment_app()
.app_id()
.find(&key.app_id)
.ok_or_else(|| "支付应用不存在".to_string())?;
if !app.enabled || app.owner_user_id != key.owner_user_id {
return Err("支付应用不可用".to_string());
}
let now = Timestamp::from_micros_since_unix_epoch(input.used_at_micros);
let key_id = key.key_id.clone();
ctx.db.payment_api_key().key_id().delete(&key_id);
ctx.db.payment_api_key().insert(PaymentApiKey {
last_used_at: Some(now),
updated_at: now,
..key
});
let key = ctx
.db
.payment_api_key()
.key_id()
.find(&key_id)
.ok_or_else(|| "支付 API Key 更新失败".to_string())?;
Ok((payment_app_snapshot(&app), payment_api_key_snapshot(&key)))
}
fn create_payment_order(
ctx: &ReducerContext,
input: PaymentOrderCreateInput,
) -> Result<PaymentOrderSnapshot, String> {
let input = build_payment_order_create_input(
input.order_id,
input.app_id,
input.owner_user_id,
input.merchant_order_id,
input.idempotency_lookup,
input.title,
input.items_json,
input.amount_cents,
input.currency,
input.provider,
input.checkout_token,
input.now_micros,
input.expires_at_micros,
input.callback_url.clone(),
input.callback_signing_key_hash.clone(),
)?;
let app = ctx
.db
.payment_app()
.app_id()
.find(&input.app_id)
.ok_or_else(|| "支付应用不存在".to_string())?;
if !app.enabled || app.owner_user_id != input.owner_user_id {
return Err("支付应用不可用".to_string());
}
if let Some(existing) = ctx
.db
.payment_order()
.idempotency_lookup()
.find(&input.idempotency_lookup)
{
if existing.app_id != input.app_id
|| existing.amount_cents != input.amount_cents
|| existing.title != input.title
|| existing.items_json != input.items_json
{
return Err("相同幂等键对应的订单字段不一致".to_string());
}
return Ok(payment_order_snapshot(&existing));
}
let merchant_order_lookup = format!("{}:{}", input.app_id, input.merchant_order_id);
if ctx
.db
.payment_order()
.merchant_order_lookup()
.find(&merchant_order_lookup)
.is_some()
{
return Err("接入方订单号已存在".to_string());
}
if ctx
.db
.payment_order()
.checkout_token()
.find(&input.checkout_token)
.is_some()
{
return Err("收银台 token 已存在".to_string());
}
let now = Timestamp::from_micros_since_unix_epoch(input.now_micros);
let expires_at = Timestamp::from_micros_since_unix_epoch(input.expires_at_micros);
ctx.db.payment_order().insert(PaymentOrder {
order_id: input.order_id.clone(),
app_id: input.app_id,
owner_user_id: input.owner_user_id,
merchant_order_id: input.merchant_order_id,
idempotency_lookup: input.idempotency_lookup,
merchant_order_lookup,
title: input.title,
items_json: input.items_json,
amount_cents: input.amount_cents,
currency: input.currency,
provider: input.provider,
provider_trade_no: None,
checkout_token: input.checkout_token,
status: PaymentOrderStatus::Pending,
created_at: now,
expires_at,
paid_at: None,
updated_at: now,
provider_code_url: None,
callback_url: input.callback_url,
callback_signing_key_hash: input.callback_signing_key_hash,
provider_attempted_at: None,
});
ctx.db
.payment_order()
.order_id()
.find(&input.order_id)
.map(|row| payment_order_snapshot(&row))
.ok_or_else(|| "支付订单创建后读取失败".to_string())
}
fn get_payment_order(
ctx: &ReducerContext,
input: PaymentOrderGetInput,
) -> Result<PaymentOrderSnapshot, String> {
let order_id = required(input.order_id, "payment_order.order_id")?;
let order = ctx
.db
.payment_order()
.order_id()
.find(&order_id)
.ok_or_else(|| "支付订单不存在".to_string())?;
if let Some(owner_user_id) = input.owner_user_id {
if order.owner_user_id != required(owner_user_id, "payment_order.owner_user_id")? {
return Err("无权访问支付订单".to_string());
}
}
Ok(payment_order_snapshot(&order))
}
fn get_payment_order_by_checkout_token(
ctx: &ReducerContext,
input: PaymentOrderCheckoutGetInput,
) -> Result<PaymentOrderSnapshot, String> {
let token = required(input.checkout_token, "payment_order.checkout_token")?;
ctx.db
.payment_order()
.checkout_token()
.find(&token)
.map(|order| payment_order_snapshot(&order))
.ok_or_else(|| "支付订单不存在".to_string())
}
fn list_payment_orders(
ctx: &ReducerContext,
input: PaymentOrderListInput,
) -> Result<Vec<PaymentOrderSnapshot>, String> {
let owner_user_id = required(input.owner_user_id, "payment_order.owner_user_id")?;
let limit = input.limit.clamp(1, 100) as usize;
let mut orders = match input.app_id {
Some(app_id) => ctx
.db
.payment_order()
.by_payment_order_app_id()
.filter(&required(app_id, "payment_order.app_id")?)
.filter(|order| order.owner_user_id == owner_user_id)
.map(|row| payment_order_snapshot(&row))
.collect::<Vec<_>>(),
None => ctx
.db
.payment_order()
.by_payment_order_owner_user_id()
.filter(&owner_user_id)
.map(|row| payment_order_snapshot(&row))
.collect::<Vec<_>>(),
};
orders.sort_by(|left, right| right.created_at_micros.cmp(&left.created_at_micros));
orders.truncate(limit);
Ok(orders)
}
fn mark_payment_order_paid(
ctx: &ReducerContext,
input: PaymentOrderMarkPaidInput,
) -> Result<PaymentOrderSnapshot, String> {
let input = build_payment_order_mark_paid_input(
input.order_id,
input.provider_trade_no,
input.amount_cents,
input.paid_at_micros,
)?;
let order = ctx
.db
.payment_order()
.order_id()
.find(&input.order_id)
.ok_or_else(|| "支付订单不存在".to_string())?;
if order.amount_cents != input.amount_cents {
return Err("支付金额与本地订单不一致".to_string());
}
if order.status == PaymentOrderStatus::Paid {
if order.provider_trade_no.as_deref() != Some(input.provider_trade_no.as_str()) {
return Err("已支付订单的平台交易号不一致".to_string());
}
return Ok(payment_order_snapshot(&order));
}
if matches!(
order.status,
PaymentOrderStatus::Closed | PaymentOrderStatus::Refunded
) {
return Err("当前支付订单状态不能确认支付".to_string());
}
let now = Timestamp::from_micros_since_unix_epoch(input.paid_at_micros);
ctx.db.payment_order().order_id().delete(&order.order_id);
ctx.db.payment_order().insert(PaymentOrder {
status: PaymentOrderStatus::Paid,
provider_trade_no: Some(input.provider_trade_no),
paid_at: Some(now),
updated_at: now,
..order
});
ctx.db
.payment_order()
.order_id()
.find(&input.order_id)
.map(|row| payment_order_snapshot(&row))
.ok_or_else(|| "支付订单确认后读取失败".to_string())
}
fn set_payment_order_provider_code(
ctx: &ReducerContext,
input: PaymentOrderProviderCodeInput,
) -> Result<PaymentOrderSnapshot, String> {
let order_id = required(input.order_id, "payment_order.order_id")?;
let provider_code_url = required(input.provider_code_url, "payment_order.provider_code_url")?;
let order = ctx
.db
.payment_order()
.order_id()
.find(&order_id)
.ok_or_else(|| "支付订单不存在".to_string())?;
if matches!(
order.status,
PaymentOrderStatus::Paid | PaymentOrderStatus::Closed
) {
return Ok(payment_order_snapshot(&order));
}
let updated_at = Timestamp::from_micros_since_unix_epoch(input.updated_at_micros);
ctx.db.payment_order().order_id().delete(&order_id);
ctx.db.payment_order().insert(PaymentOrder {
provider_code_url: Some(provider_code_url),
updated_at,
..order
});
ctx.db
.payment_order()
.order_id()
.find(&order_id)
.map(|row| payment_order_snapshot(&row))
.ok_or_else(|| "支付订单 provider 结果写入后读取失败".to_string())
}
fn mark_payment_order_provider_attempted(
ctx: &ReducerContext,
input: PaymentOrderProviderAttemptInput,
) -> Result<PaymentOrderSnapshot, String> {
let order_id = required(input.order_id, "payment_order.order_id")?;
let order = ctx
.db
.payment_order()
.order_id()
.find(&order_id)
.ok_or_else(|| "支付订单不存在".to_string())?;
if order.provider_code_url.is_some() || order.status == PaymentOrderStatus::Paid {
return Ok(payment_order_snapshot(&order));
}
if order.provider_attempted_at.is_some() {
return Err("PAYMENT_PROVIDER_PENDING".to_string());
}
let attempted_at = Timestamp::from_micros_since_unix_epoch(input.attempted_at_micros);
ctx.db.payment_order().order_id().delete(&order_id);
ctx.db.payment_order().insert(PaymentOrder {
status: PaymentOrderStatus::Paying,
provider_attempted_at: Some(attempted_at),
updated_at: attempted_at,
..order
});
ctx.db
.payment_order()
.order_id()
.find(&order_id)
.map(|row| payment_order_snapshot(&row))
.ok_or_else(|| "支付订单 provider claim 后读取失败".to_string())
}
fn enqueue_payment_webhook(
ctx: &ReducerContext,
input: PaymentWebhookEnqueueInput,
) -> Result<Option<PaymentWebhookDeliverySnapshot>, String> {
let order_id = required(input.order_id, "payment_webhook.order_id")?;
let payload_json = required(input.payload_json, "payment_webhook.payload_json")?;
let signature = required(input.signature, "payment_webhook.signature")?;
if payload_json.len() > 64 * 1024 {
return Err("payment_webhook.payload_json 超出大小限制".to_string());
}
let order = ctx
.db
.payment_order()
.order_id()
.find(&order_id)
.ok_or_else(|| "支付订单不存在".to_string())?;
if order.status != PaymentOrderStatus::Paid {
return Err("只有已支付订单可以创建外部回调".to_string());
}
let Some(callback_url) = order.callback_url.clone() else {
return Ok(None);
};
let delivery_id = format!("payment-webhook:{order_id}");
if let Some(existing) = ctx
.db
.payment_webhook_delivery()
.delivery_id()
.find(&delivery_id)
{
if existing.payload_json != payload_json || existing.signature != signature {
return Err("相同支付订单的回调内容不一致".to_string());
}
return Ok(Some(payment_webhook_snapshot(&existing)));
}
let now = Timestamp::from_micros_since_unix_epoch(input.now_micros);
ctx.db
.payment_webhook_delivery()
.insert(PaymentWebhookDelivery {
delivery_id: delivery_id.clone(),
order_id: order.order_id,
app_id: order.app_id,
owner_user_id: order.owner_user_id,
callback_url,
payload_json,
signature,
status: PAYMENT_WEBHOOK_STATUS_PENDING.to_string(),
attempt_count: 0,
available_at: now,
lease_worker_id: None,
lease_expires_at: None,
last_error_message: None,
created_at: now,
updated_at: now,
});
ctx.db
.payment_webhook_delivery()
.delivery_id()
.find(&delivery_id)
.map(|row| payment_webhook_snapshot(&row))
.map(Some)
.ok_or_else(|| "支付回调写入后读取失败".to_string())
}
fn claim_payment_webhooks(
ctx: &ReducerContext,
input: PaymentWebhookClaimInput,
) -> Result<Vec<PaymentWebhookDeliverySnapshot>, String> {
let worker_id = required(input.worker_id, "payment_webhook.worker_id")?;
let now = Timestamp::from_micros_since_unix_epoch(input.now_micros);
let lease_expires_at = Timestamp::from_micros_since_unix_epoch(input.lease_expires_at_micros);
let limit = input.limit.clamp(1, 50) as usize;
let mut rows = ctx
.db
.payment_webhook_delivery()
.iter()
.filter(|row| {
(row.status == PAYMENT_WEBHOOK_STATUS_PENDING && row.available_at <= now)
|| (row.status == PAYMENT_WEBHOOK_STATUS_PROCESSING
&& row.lease_expires_at.is_some_and(|value| value <= now))
})
.collect::<Vec<_>>();
rows.sort_by(|left, right| left.available_at.cmp(&right.available_at));
rows.truncate(limit);
let mut claimed = Vec::with_capacity(rows.len());
for row in rows {
let delivery_id = row.delivery_id.clone();
ctx.db
.payment_webhook_delivery()
.delivery_id()
.delete(&delivery_id);
ctx.db
.payment_webhook_delivery()
.insert(PaymentWebhookDelivery {
status: PAYMENT_WEBHOOK_STATUS_PROCESSING.to_string(),
lease_worker_id: Some(worker_id.clone()),
lease_expires_at: Some(lease_expires_at),
updated_at: now,
..row
});
if let Some(updated) = ctx
.db
.payment_webhook_delivery()
.delivery_id()
.find(&delivery_id)
{
claimed.push(payment_webhook_snapshot(&updated));
}
}
Ok(claimed)
}
fn complete_payment_webhook(
ctx: &ReducerContext,
input: PaymentWebhookCompleteInput,
) -> Result<PaymentWebhookDeliverySnapshot, String> {
let delivery_id = required(input.delivery_id, "payment_webhook.delivery_id")?;
let worker_id = required(input.worker_id, "payment_webhook.worker_id")?;
let mut row = ctx
.db
.payment_webhook_delivery()
.delivery_id()
.find(&delivery_id)
.ok_or_else(|| "支付回调投递不存在".to_string())?;
if row.status != PAYMENT_WEBHOOK_STATUS_PROCESSING
|| row.lease_worker_id.as_deref() != Some(worker_id.as_str())
{
return Err("支付回调投递租约不属于当前 worker".to_string());
}
let now = Timestamp::from_micros_since_unix_epoch(input.now_micros);
row.lease_worker_id = None;
row.lease_expires_at = None;
row.updated_at = now;
if input.success {
row.status = PAYMENT_WEBHOOK_STATUS_DELIVERED.to_string();
row.last_error_message = None;
} else {
row.attempt_count = row.attempt_count.saturating_add(1);
row.last_error_message = input
.error_message
.and_then(|value| normalize_optional_text(Some(value)));
if row.attempt_count >= PAYMENT_WEBHOOK_MAX_ATTEMPTS {
row.status = PAYMENT_WEBHOOK_STATUS_FAILED.to_string();
} else {
row.status = PAYMENT_WEBHOOK_STATUS_PENDING.to_string();
row.available_at =
Timestamp::from_micros_since_unix_epoch(input.next_available_at_micros);
}
}
ctx.db
.payment_webhook_delivery()
.delivery_id()
.delete(&delivery_id);
ctx.db.payment_webhook_delivery().insert(row.clone());
Ok(payment_webhook_snapshot(&row))
}
fn payment_app_snapshot(row: &PaymentApp) -> PaymentAppSnapshot {
PaymentAppSnapshot {
app_id: row.app_id.clone(),
owner_user_id: row.owner_user_id.clone(),
name: row.name.clone(),
callback_url: row.callback_url.clone(),
enabled: row.enabled,
created_at_micros: row.created_at.to_micros_since_unix_epoch(),
updated_at_micros: row.updated_at.to_micros_since_unix_epoch(),
}
}
fn payment_api_key_snapshot(row: &PaymentApiKey) -> PaymentApiKeySnapshot {
PaymentApiKeySnapshot {
key_id: row.key_id.clone(),
app_id: row.app_id.clone(),
owner_user_id: row.owner_user_id.clone(),
name: row.name.clone(),
key_prefix: row.key_prefix.clone(),
scopes_json: row.scopes_json.clone(),
created_at_micros: row.created_at.to_micros_since_unix_epoch(),
last_used_at_micros: row
.last_used_at
.map(|value| value.to_micros_since_unix_epoch()),
revoked_at_micros: row
.revoked_at
.map(|value| value.to_micros_since_unix_epoch()),
updated_at_micros: row.updated_at.to_micros_since_unix_epoch(),
}
}
fn payment_order_snapshot(row: &PaymentOrder) -> PaymentOrderSnapshot {
PaymentOrderSnapshot {
order_id: row.order_id.clone(),
app_id: row.app_id.clone(),
owner_user_id: row.owner_user_id.clone(),
merchant_order_id: row.merchant_order_id.clone(),
idempotency_lookup: row.idempotency_lookup.clone(),
title: row.title.clone(),
items_json: row.items_json.clone(),
amount_cents: row.amount_cents,
currency: row.currency.clone(),
provider: row.provider,
provider_trade_no: row.provider_trade_no.clone(),
provider_code_url: row.provider_code_url.clone(),
callback_url: row.callback_url.clone(),
callback_signing_key_hash: row.callback_signing_key_hash.clone(),
provider_attempted_at_micros: row
.provider_attempted_at
.map(|value| value.to_micros_since_unix_epoch()),
checkout_token: row.checkout_token.clone(),
status: row.status,
created_at_micros: row.created_at.to_micros_since_unix_epoch(),
expires_at_micros: row.expires_at.to_micros_since_unix_epoch(),
paid_at_micros: row.paid_at.map(|value| value.to_micros_since_unix_epoch()),
updated_at_micros: row.updated_at.to_micros_since_unix_epoch(),
}
}
fn payment_webhook_snapshot(row: &PaymentWebhookDelivery) -> PaymentWebhookDeliverySnapshot {
PaymentWebhookDeliverySnapshot {
delivery_id: row.delivery_id.clone(),
order_id: row.order_id.clone(),
app_id: row.app_id.clone(),
owner_user_id: row.owner_user_id.clone(),
callback_url: row.callback_url.clone(),
payload_json: row.payload_json.clone(),
signature: row.signature.clone(),
status: row.status.clone(),
attempt_count: row.attempt_count,
available_at_micros: row.available_at.to_micros_since_unix_epoch(),
lease_worker_id: row.lease_worker_id.clone(),
lease_expires_at_micros: row
.lease_expires_at
.map(|value| value.to_micros_since_unix_epoch()),
last_error_message: row.last_error_message.clone(),
created_at_micros: row.created_at.to_micros_since_unix_epoch(),
updated_at_micros: row.updated_at.to_micros_since_unix_epoch(),
}
}
fn payment_app_error(error_message: String) -> PaymentAppProcedureResult {
PaymentAppProcedureResult {
ok: false,
app: None,
key: None,
apps: Vec::new(),
keys: Vec::new(),
error_message: Some(error_message),
}
}
fn payment_order_single_ok(order: PaymentOrderSnapshot) -> PaymentOrderProcedureResult {
PaymentOrderProcedureResult {
ok: true,
order: Some(order),
orders: Vec::new(),
error_message: None,
}
}
fn payment_order_error(error_message: String) -> PaymentOrderProcedureResult {
PaymentOrderProcedureResult {
ok: false,
order: None,
orders: Vec::new(),
error_message: Some(error_message),
}
}
fn payment_webhook_error(error_message: String) -> PaymentWebhookProcedureResult {
PaymentWebhookProcedureResult {
ok: false,
delivery: None,
deliveries: Vec::new(),
error_message: Some(error_message),
}
}
fn required(value: String, field: &str) -> Result<String, String> {
let value = value.trim().to_string();
if value.is_empty() {
return Err(format!("{field} 不能为空"));
}
Ok(value)
}