5d6a4675fb
CI/CD Pipeline / Check if frontend-only change (push) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (push) Successful in 1s
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 3s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 4s
CI/CD Pipeline / Check push changed paths (push) Successful in 15s
CI/CD Pipeline / Frontend Lint (push) Has been skipped
CI/CD Pipeline / PR Build API Image (push) Has been skipped
CI/CD Pipeline / PR Build Web Image (push) Has been skipped
CI/CD Pipeline / PR Build Worker Image (push) Has been skipped
CI/CD Pipeline / Build Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Validate - Style (pull_request) Has been skipped
CI/CD Pipeline / Validate - Security (pull_request) Has been skipped
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Has been skipped
CI/CD Pipeline / Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / PR Build Web Image (pull_request) Has been skipped
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 52s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 54s
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 31s
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / CI Gate (pull_request) Successful in 2s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 1m49s
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / ACR Image Cleanup (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Web Image (push) Successful in 55s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m17s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been skipped
CI/CD Pipeline / Build Staging API Image (push) Successful in 3m23s
CI/CD Pipeline / Retag skipped Staging API Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Has been skipped
CI/CD Pipeline / Integration Tests (push) Successful in 4m19s
CI/CD Pipeline / Frontend Unit Tests (push) Successful in 4m30s
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 4m47s
CI/CD Pipeline / Validate - Style (push) Successful in 5m10s
CI/CD Pipeline / Validate - Security (push) Has been cancelled
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Has been cancelled
CI/CD Pipeline / Staging E2E Tests (push) Has been cancelled
CI/CD Pipeline / Staging API Integration Tests (push) Has been cancelled
CI/CD Pipeline / Build Production API Image (push) Has been cancelled
CI/CD Pipeline / Build Production Web Image (push) Has been cancelled
CI/CD Pipeline / Build Production Worker Image (push) Has been cancelled
CI/CD Pipeline / Deploy Production (push) Has been cancelled
CI/CD Pipeline / Production Browser E2E (push) Has been cancelled
CI/CD Pipeline / ACR Image Cleanup (push) Has been cancelled
CI/CD Pipeline / Canary Release to Production (push) Has been cancelled
CI/CD Pipeline / CI Gate (push) Has been cancelled
CI/CD Pipeline / Unit Tests (push) Has been cancelled
AI Code Review / AI Code Review (pull_request) Has been cancelled
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com> Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
541 lines
19 KiB
Python
541 lines
19 KiB
Python
"""积分服务层 — 积分账户、扣减、充值、流水、每日免费额度 (#1895)
|
||
|
||
直接操作 SQLAlchemy session,不走 Repository 抽象层,简化事务处理。
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
import uuid
|
||
from datetime import UTC, datetime, timedelta
|
||
from typing import Any
|
||
|
||
from sqlalchemy.orm import Session
|
||
|
||
from packages.domain.points_rules import (
|
||
POINTS_PACKAGES,
|
||
)
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
|
||
# ── 延迟导入模型(避免循环/顺序依赖) ────────────────────────────────────
|
||
|
||
|
||
def _get_models():
|
||
"""延迟获取积分相关模型类。"""
|
||
from packages.adapters.sqlalchemy_impl.models import (
|
||
DailyUsageRecordModel,
|
||
PointsAccountModel,
|
||
PointsOrderModel,
|
||
PointsTransactionModel,
|
||
UserModel,
|
||
)
|
||
|
||
return (
|
||
PointsAccountModel,
|
||
PointsTransactionModel,
|
||
PointsOrderModel,
|
||
DailyUsageRecordModel,
|
||
UserModel,
|
||
)
|
||
|
||
|
||
def _get_redis_client():
|
||
"""获取 Redis 客户端,用于每日额度缓存。"""
|
||
try:
|
||
import redis as redis_lib
|
||
from app.config import settings
|
||
|
||
return redis_lib.from_url(settings.REDIS_URL, decode_responses=True)
|
||
except Exception:
|
||
return None
|
||
|
||
|
||
class PointsService:
|
||
"""积分核心服务。"""
|
||
|
||
# ──────────────── 账户管理 ────────────────
|
||
|
||
def get_or_create_account(self, user_id: str, db: Session) -> dict[str, Any]:
|
||
"""获取或创建积分账户,返回账户快照。"""
|
||
PointsAccountModel, _, _, _, _ = _get_models()
|
||
|
||
account = db.query(PointsAccountModel).filter(PointsAccountModel.user_id == user_id).first()
|
||
if account is None:
|
||
account = PointsAccountModel(
|
||
id=uuid.uuid4().hex,
|
||
user_id=user_id,
|
||
balance=0,
|
||
total_earned=0,
|
||
total_spent=0,
|
||
)
|
||
db.add(account)
|
||
db.flush()
|
||
|
||
return {
|
||
"id": account.id,
|
||
"user_id": account.user_id,
|
||
"balance": account.balance,
|
||
"total_earned": account.total_earned,
|
||
"total_spent": account.total_spent,
|
||
}
|
||
|
||
# ──────────────── 余额检查 ────────────────
|
||
|
||
def check_balance(self, user_id: str, amount: float, db: Session) -> dict[str, Any]:
|
||
"""检查余额是否足够。"""
|
||
account_data = self.get_or_create_account(user_id, db)
|
||
balance = account_data["balance"]
|
||
return {
|
||
"sufficient": balance >= amount,
|
||
"balance": balance,
|
||
"required": amount,
|
||
"remaining_after": balance - amount,
|
||
}
|
||
|
||
# ──────────────── 积分扣减(事务性) ────────────────
|
||
|
||
def deduct_points(
|
||
self,
|
||
user_id: str,
|
||
amount: float,
|
||
source: str,
|
||
db: Session,
|
||
description: str = "",
|
||
ref_id: str = "",
|
||
) -> dict[str, Any]:
|
||
"""扣减积分(事务性:SELECT FOR UPDATE → 检查余额 → 扣减 → 流水 → 同步用户表)。
|
||
|
||
Returns:
|
||
{"success": True/False, "balance": float, "transaction_id": str|None}
|
||
"""
|
||
PointsAccountModel, PointsTransactionModel, _, _, UserModel = _get_models()
|
||
|
||
try:
|
||
# 1. 行锁获取账户
|
||
account = (
|
||
db.query(PointsAccountModel).filter(PointsAccountModel.user_id == user_id).with_for_update().first()
|
||
)
|
||
if account is None:
|
||
account = PointsAccountModel(
|
||
id=uuid.uuid4().hex,
|
||
user_id=user_id,
|
||
balance=0,
|
||
total_earned=0,
|
||
total_spent=0,
|
||
)
|
||
db.add(account)
|
||
db.flush()
|
||
|
||
# 2. 检查余额
|
||
if account.balance < amount:
|
||
return {
|
||
"success": False,
|
||
"balance": account.balance,
|
||
"transaction_id": None,
|
||
}
|
||
|
||
# 3. 扣减余额
|
||
account.balance -= amount
|
||
account.total_spent += amount
|
||
|
||
# 4. 创建流水
|
||
txn_id = uuid.uuid4().hex
|
||
txn = PointsTransactionModel(
|
||
id=txn_id,
|
||
user_id=user_id,
|
||
account_id=account.id,
|
||
type="deduct",
|
||
source=source,
|
||
amount=amount,
|
||
balance_after=account.balance,
|
||
description=description or f"积分扣减: {source}",
|
||
ref_id=ref_id,
|
||
)
|
||
db.add(txn)
|
||
|
||
# 5. 同步用户表 points_balance
|
||
db.execute(
|
||
UserModel.__table__.update()
|
||
.where(UserModel.__table__.c.id == user_id)
|
||
.values(points_balance=account.balance)
|
||
)
|
||
|
||
db.commit()
|
||
return {
|
||
"success": True,
|
||
"balance": account.balance,
|
||
"transaction_id": txn_id,
|
||
}
|
||
|
||
except Exception:
|
||
db.rollback()
|
||
logger.exception(
|
||
"积分扣减失败: user_id=%s, amount=%.2f, source=%s",
|
||
user_id,
|
||
amount,
|
||
source,
|
||
)
|
||
return {"success": False, "balance": 0, "transaction_id": None}
|
||
|
||
# ──────────────── 积分增加 ────────────────
|
||
|
||
def add_points(
|
||
self,
|
||
user_id: str,
|
||
amount: float,
|
||
source: str,
|
||
db: Session,
|
||
description: str = "",
|
||
ref_id: str = "",
|
||
) -> dict[str, Any]:
|
||
"""增加积分(充值/赠送/退款)。"""
|
||
PointsAccountModel, PointsTransactionModel, _, _, UserModel = _get_models()
|
||
|
||
try:
|
||
account = (
|
||
db.query(PointsAccountModel).filter(PointsAccountModel.user_id == user_id).with_for_update().first()
|
||
)
|
||
if account is None:
|
||
account = PointsAccountModel(
|
||
id=uuid.uuid4().hex,
|
||
user_id=user_id,
|
||
balance=0,
|
||
total_earned=0,
|
||
total_spent=0,
|
||
)
|
||
db.add(account)
|
||
db.flush()
|
||
|
||
account.balance += amount
|
||
account.total_earned += amount
|
||
|
||
txn_id = uuid.uuid4().hex
|
||
txn = PointsTransactionModel(
|
||
id=txn_id,
|
||
user_id=user_id,
|
||
account_id=account.id,
|
||
type="add",
|
||
source=source,
|
||
amount=amount,
|
||
balance_after=account.balance,
|
||
description=description or f"积分增加: {source}",
|
||
ref_id=ref_id,
|
||
)
|
||
db.add(txn)
|
||
|
||
db.execute(
|
||
UserModel.__table__.update()
|
||
.where(UserModel.__table__.c.id == user_id)
|
||
.values(points_balance=account.balance)
|
||
)
|
||
|
||
db.commit()
|
||
return {
|
||
"success": True,
|
||
"balance": account.balance,
|
||
"transaction_id": txn_id,
|
||
}
|
||
|
||
except Exception:
|
||
db.rollback()
|
||
logger.exception(
|
||
"积分增加失败: user_id=%s, amount=%.2f, source=%s",
|
||
user_id,
|
||
amount,
|
||
source,
|
||
)
|
||
return {"success": False, "balance": 0, "transaction_id": None}
|
||
|
||
# ──────────────── 积分退还 ────────────────
|
||
|
||
def refund_points(
|
||
self,
|
||
user_id: str,
|
||
amount: float,
|
||
source: str,
|
||
db: Session,
|
||
ref_id: str = "",
|
||
description: str = "",
|
||
) -> dict[str, Any]:
|
||
"""退还积分(业务失败回退)。内部调用 add_points,source 前缀 refund:。"""
|
||
return self.add_points(
|
||
user_id=user_id,
|
||
amount=amount,
|
||
source=f"refund:{source}",
|
||
db=db,
|
||
description=description or f"积分退还: {source}",
|
||
ref_id=ref_id,
|
||
)
|
||
|
||
# ──────────────── 爆款视频(viral_video)动态定价 ────────────────
|
||
|
||
def deduct_viral_video(self, user_id: str, credits: float, job_id: str, db: Session) -> dict[str, Any]:
|
||
"""爆款视频预扣积分(confirm-copy 阶段)。"""
|
||
return self.deduct_points(
|
||
user_id=user_id,
|
||
amount=float(credits or 0),
|
||
source="viral_video",
|
||
db=db,
|
||
description="爆款视频生成",
|
||
ref_id=job_id,
|
||
)
|
||
|
||
def settle_viral_video(
|
||
self,
|
||
user_id: str,
|
||
estimated: float,
|
||
actual: float,
|
||
txn_id: str,
|
||
db: Session,
|
||
) -> dict[str, Any]:
|
||
"""爆款视频完成后按实际 tokens 结算(多退少补)。
|
||
|
||
- actual < estimated: 退差额
|
||
- actual > estimated: 补扣差额(余额不足时记 warning,不阻塞完成)
|
||
- |diff| < 0.01: 不动
|
||
"""
|
||
diff = round(float(actual or 0) - float(estimated or 0), 2)
|
||
if abs(diff) < 0.01:
|
||
return {"success": True, "action": "none", "diff": 0.0}
|
||
if diff < 0:
|
||
refund = round(-diff, 2)
|
||
try:
|
||
res = self.refund_points(
|
||
user_id=user_id,
|
||
amount=refund,
|
||
source="viral_video",
|
||
db=db,
|
||
ref_id=txn_id,
|
||
description="爆款视频结算退费",
|
||
)
|
||
return {"success": bool(res.get("success")), "action": "refund", "diff": -refund, "amount": refund}
|
||
except Exception:
|
||
logger.exception("[viral_video] 结算退费异常 user_id=%s refund=%.2f", user_id, refund)
|
||
return {"success": False, "action": "refund", "diff": -refund}
|
||
else:
|
||
extra = round(diff, 2)
|
||
try:
|
||
res = self.deduct_points(
|
||
user_id=user_id,
|
||
amount=extra,
|
||
source="viral_video",
|
||
db=db,
|
||
description="爆款视频结算补扣",
|
||
ref_id=txn_id,
|
||
)
|
||
if not res.get("success"):
|
||
logger.warning(
|
||
"[viral_video] 结算补扣余额不足 user_id=%s extra=%.2f balance=%s (不阻塞任务完成)",
|
||
user_id,
|
||
extra,
|
||
res.get("balance"),
|
||
)
|
||
return {"success": bool(res.get("success")), "action": "deduct", "diff": extra, "amount": extra}
|
||
except Exception:
|
||
logger.exception("[viral_video] 结算补扣异常 user_id=%s extra=%.2f", user_id, extra)
|
||
return {"success": False, "action": "deduct", "diff": extra}
|
||
|
||
def refund_viral_video(self, user_id: str, credits: float, txn_id: str, db: Session) -> dict[str, Any]:
|
||
"""爆款视频失败全额退款。"""
|
||
amount = float(credits or 0)
|
||
if amount <= 0:
|
||
return {"success": True, "action": "none", "amount": 0.0}
|
||
try:
|
||
return self.refund_points(
|
||
user_id=user_id,
|
||
amount=amount,
|
||
source="viral_video",
|
||
db=db,
|
||
ref_id=txn_id,
|
||
description="爆款视频失败退款",
|
||
)
|
||
except Exception:
|
||
logger.exception("[viral_video] 失败退款异常 user_id=%s amount=%.2f", user_id, amount)
|
||
return {"success": False, "action": "refund", "amount": amount}
|
||
|
||
# ──────────────── 流水查询 ────────────────
|
||
|
||
def get_transactions(
|
||
self,
|
||
user_id: str,
|
||
db: Session,
|
||
page: int = 1,
|
||
page_size: int = 20,
|
||
type_filter: str | None = None,
|
||
source_filter: str | None = None,
|
||
start_date: datetime | None = None,
|
||
end_date: datetime | None = None,
|
||
) -> dict[str, Any]:
|
||
"""查询积分流水(分页+筛选)。"""
|
||
_, PointsTransactionModel, _, _, _ = _get_models()
|
||
|
||
query = db.query(PointsTransactionModel).filter(PointsTransactionModel.user_id == user_id)
|
||
|
||
if type_filter:
|
||
query = query.filter(PointsTransactionModel.type == type_filter)
|
||
if source_filter:
|
||
query = query.filter(PointsTransactionModel.source == source_filter)
|
||
if start_date:
|
||
query = query.filter(PointsTransactionModel.created_at >= start_date)
|
||
if end_date:
|
||
query = query.filter(PointsTransactionModel.created_at <= end_date)
|
||
|
||
total = query.count()
|
||
items = (
|
||
query.order_by(PointsTransactionModel.created_at.desc())
|
||
.offset((page - 1) * page_size)
|
||
.limit(page_size)
|
||
.all()
|
||
)
|
||
|
||
return {
|
||
"items": [
|
||
{
|
||
"id": item.id,
|
||
"type": item.type,
|
||
"source": item.source,
|
||
"amount": item.amount,
|
||
"balance_after": item.balance_after,
|
||
"description": item.description,
|
||
"ref_id": item.ref_id,
|
||
"created_at": (item.created_at.isoformat() if item.created_at else None),
|
||
}
|
||
for item in items
|
||
],
|
||
"total": total,
|
||
"page": page,
|
||
"page_size": page_size,
|
||
}
|
||
|
||
# ──────────────── 每日免费混剪额度(已下线:智能混剪全免费) ────────────────
|
||
|
||
def get_daily_usage(self, user_id: str, db: Session) -> dict[str, Any]:
|
||
"""查询今日免费额度使用情况(智能混剪已全免费,返回 unlimited)。"""
|
||
now = datetime.now(UTC)
|
||
tomorrow = (now + timedelta(days=1)).replace(hour=0, minute=0, second=0, microsecond=0)
|
||
return {
|
||
"free_clips_used": 0,
|
||
"free_clips_limit": -1, # -1 表示 unlimited
|
||
"free_clips_remaining": -1,
|
||
"reset_at": tomorrow.isoformat(),
|
||
}
|
||
|
||
# ──────────────── 订单管理 ────────────────
|
||
|
||
def create_order(
|
||
self,
|
||
user_id: str,
|
||
order_type: str,
|
||
product_code: str,
|
||
db: Session,
|
||
) -> dict[str, Any]:
|
||
"""创建积分充值或会员购买订单。
|
||
|
||
Args:
|
||
order_type: "points" 或 "membership"
|
||
product_code: 积分包 code (如 "starter_pack") 或会员类型 (如 "monthly")
|
||
"""
|
||
_, _, PointsOrderModel, _, _ = _get_models()
|
||
|
||
amount_cents = 0
|
||
points_amount = 0
|
||
if order_type == "points":
|
||
package = POINTS_PACKAGES.get(product_code)
|
||
if not package:
|
||
raise ValueError(f"Unknown points package: {product_code}")
|
||
amount_cents = package["price_cents"]
|
||
points_amount = package["points"]
|
||
elif order_type == "membership":
|
||
from packages.domain.points_rules import MEMBERSHIP_PRICES
|
||
|
||
membership = MEMBERSHIP_PRICES.get(product_code)
|
||
if not membership:
|
||
raise ValueError(f"Unknown membership type: {product_code}")
|
||
amount_cents = membership["price_cents"]
|
||
else:
|
||
raise ValueError(f"Unknown order type: {order_type}")
|
||
|
||
order = PointsOrderModel(
|
||
id=uuid.uuid4().hex,
|
||
user_id=user_id,
|
||
order_type=order_type,
|
||
product_code=product_code,
|
||
amount_cents=amount_cents,
|
||
original_amount_cents=amount_cents,
|
||
points_amount=points_amount,
|
||
status="pending",
|
||
)
|
||
db.add(order)
|
||
db.commit()
|
||
|
||
return {
|
||
"id": order.id,
|
||
"order_type": order.order_type,
|
||
"product_code": order.product_code,
|
||
"amount_cents": order.amount_cents,
|
||
"status": order.status,
|
||
"created_at": (order.created_at.isoformat() if order.created_at else None),
|
||
}
|
||
|
||
def confirm_payment(
|
||
self,
|
||
order_id: str,
|
||
payment_id: str,
|
||
db: Session,
|
||
) -> dict[str, Any]:
|
||
"""确认支付 → 更新订单状态 → 发放积分或会员。"""
|
||
_, _, PointsOrderModel, _, UserModel = _get_models()
|
||
|
||
try:
|
||
order = db.query(PointsOrderModel).filter(PointsOrderModel.id == order_id).with_for_update().first()
|
||
if order is None:
|
||
return {"success": False, "message": "订单不存在"}
|
||
if order.status != "pending":
|
||
return {"success": False, "message": f"订单状态异常: {order.status}"}
|
||
|
||
# 更新订单状态
|
||
order.status = "paid"
|
||
order.payment_id = payment_id
|
||
order.paid_at = datetime.now(UTC)
|
||
|
||
if order.order_type == "points":
|
||
# 发放积分
|
||
self.add_points(
|
||
user_id=order.user_id,
|
||
amount=order.points_amount,
|
||
source=f"recharge:{order.product_code}",
|
||
db=db,
|
||
description=f"积分充值: {order.product_code}",
|
||
ref_id=order.id,
|
||
)
|
||
elif order.order_type == "membership":
|
||
# 激活会员
|
||
from packages.domain.points_rules import MEMBERSHIP_PRICES
|
||
|
||
membership = MEMBERSHIP_PRICES.get(order.product_code, {})
|
||
duration_days = membership.get("duration_days", 30)
|
||
|
||
user = db.query(UserModel).filter(UserModel.id == order.user_id).first()
|
||
if user:
|
||
now = datetime.now(UTC)
|
||
current_expires = user.member_expires_at or now
|
||
if current_expires < now:
|
||
current_expires = now
|
||
user.member_expires_at = current_expires + timedelta(days=duration_days)
|
||
user.member_type = order.product_code
|
||
user.is_member = True
|
||
|
||
db.commit()
|
||
return {
|
||
"success": True,
|
||
"message": "支付确认成功",
|
||
"order_id": order_id,
|
||
}
|
||
|
||
except Exception:
|
||
db.rollback()
|
||
logger.exception("确认支付失败: order_id=%s", order_id)
|
||
return {"success": False, "message": "确认支付异常"}
|