增加后台支付列表,修改前台有问题也能支付成功,增加支付日志
This commit is contained in:
@@ -0,0 +1,29 @@
|
||||
"""add_refund_fields_to_payment_orders
|
||||
|
||||
Revision ID: a1b2c3d4e5f6
|
||||
Revises: ed59aefc83da
|
||||
Create Date: 2026-06-10 12:00:00.000000
|
||||
"""
|
||||
from typing import Sequence, Union
|
||||
|
||||
from alembic import op
|
||||
import sqlalchemy as sa
|
||||
|
||||
|
||||
# revision identifiers, used by Alembic.
|
||||
revision: str = 'a1b2c3d4e5f6'
|
||||
down_revision: Union[str, None] = 'ed59aefc83da'
|
||||
branch_labels: Union[str, Sequence[str], None] = None
|
||||
depends_on: Union[str, Sequence[str], None] = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
op.add_column('payment_orders', sa.Column('refund_trade_no', sa.String(length=128), nullable=True))
|
||||
op.add_column('payment_orders', sa.Column('refunded_at', sa.DateTime(timezone=True), nullable=True))
|
||||
op.add_column('payment_orders', sa.Column('refund_amount', sa.Float(), nullable=True))
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
op.drop_column('payment_orders', 'refund_amount')
|
||||
op.drop_column('payment_orders', 'refunded_at')
|
||||
op.drop_column('payment_orders', 'refund_trade_no')
|
||||
@@ -440,6 +440,123 @@ async def batch_update_payment_configs(
|
||||
return {"ok": True}
|
||||
|
||||
|
||||
@router.get("/payment-stats")
|
||||
async def get_payment_stats(
|
||||
admin: User = Depends(get_admin_user),
|
||||
db: AsyncSession = Depends(get_db),
|
||||
):
|
||||
"""Return payment statistics for admin dashboard."""
|
||||
from sqlalchemy import func
|
||||
|
||||
# Status breakdown
|
||||
status_result = await db.execute(
|
||||
select(
|
||||
PaymentOrder.status,
|
||||
func.count().label("count"),
|
||||
func.coalesce(func.sum(PaymentOrder.amount), 0).label("amount"),
|
||||
).group_by(PaymentOrder.status)
|
||||
)
|
||||
by_status = {}
|
||||
for row in status_result.all():
|
||||
by_status[row.status] = {"count": row.count, "amount": float(row.amount)}
|
||||
|
||||
# Today's stats
|
||||
today_start = datetime.now().replace(hour=0, minute=0, second=0, microsecond=0)
|
||||
today_result = await db.execute(
|
||||
select(
|
||||
func.count().label("paid_count"),
|
||||
func.coalesce(func.sum(PaymentOrder.amount), 0).label("paid_amount"),
|
||||
).where(
|
||||
PaymentOrder.status == "paid",
|
||||
PaymentOrder.paid_at >= today_start,
|
||||
)
|
||||
)
|
||||
today_row = today_result.one()
|
||||
|
||||
# Recent 50 orders
|
||||
recent_result = await db.execute(
|
||||
select(PaymentOrder)
|
||||
.order_by(PaymentOrder.created_at.desc())
|
||||
.limit(50)
|
||||
)
|
||||
recent = recent_result.scalars().all()
|
||||
|
||||
return {
|
||||
"by_status": by_status,
|
||||
"today": {
|
||||
"paid_count": today_row.paid_count,
|
||||
"paid_amount": float(today_row.paid_amount),
|
||||
},
|
||||
"recent": [
|
||||
{
|
||||
"id": o.id,
|
||||
"order_no": o.order_no,
|
||||
"user_id": o.user_id,
|
||||
"amount": o.amount,
|
||||
"credits": o.credits,
|
||||
"payment_method": o.payment_method,
|
||||
"status": o.status,
|
||||
"trade_no": o.trade_no,
|
||||
"paid_at": o.paid_at.isoformat() if o.paid_at else None,
|
||||
"created_at": o.created_at.isoformat() if o.created_at else None,
|
||||
}
|
||||
for o in recent
|
||||
],
|
||||
}
|
||||
|
||||
|
||||
@router.get("/payment-orders")
|
||||
async def get_admin_payment_orders(
|
||||
method: str | None = None,
|
||||
status: str | None = None,
|
||||
page: int = 1,
|
||||
page_size: int = 20,
|
||||
admin: User = Depends(get_admin_user),
|
||||
db: AsyncSession = Depends(get_db),
|
||||
):
|
||||
"""Return paginated payment orders for admin."""
|
||||
query = select(PaymentOrder)
|
||||
if method:
|
||||
query = query.where(PaymentOrder.payment_method == method)
|
||||
if status:
|
||||
query = query.where(PaymentOrder.status == status)
|
||||
|
||||
# Count total
|
||||
count_result = await db.execute(
|
||||
select(func.count()).select_from(query.subquery())
|
||||
)
|
||||
total = count_result.scalar() or 0
|
||||
|
||||
# Paginated results
|
||||
result = await db.execute(
|
||||
query.order_by(PaymentOrder.created_at.desc())
|
||||
.offset((page - 1) * page_size)
|
||||
.limit(page_size)
|
||||
)
|
||||
orders = result.scalars().all()
|
||||
|
||||
return {
|
||||
"total": total,
|
||||
"page": page,
|
||||
"page_size": page_size,
|
||||
"items": [
|
||||
{
|
||||
"id": o.id,
|
||||
"order_no": o.order_no,
|
||||
"user_id": o.user_id,
|
||||
"amount": o.amount,
|
||||
"credits": o.credits,
|
||||
"payment_method": o.payment_method,
|
||||
"status": o.status,
|
||||
"trade_no": o.trade_no,
|
||||
"paid_at": o.paid_at.isoformat() if o.paid_at else None,
|
||||
"created_at": o.created_at.isoformat() if o.created_at else None,
|
||||
}
|
||||
for o in orders
|
||||
],
|
||||
}
|
||||
|
||||
|
||||
@router.put("/payment-configs/{config_id}")
|
||||
async def update_payment_config(
|
||||
config_id: str,
|
||||
@@ -1250,3 +1367,76 @@ async def admin_generate_video(
|
||||
await db.flush()
|
||||
|
||||
return {"message": "ok", "record_id": record_id}
|
||||
|
||||
|
||||
# ── Payment Stats ────────────────────────────────────────
|
||||
|
||||
@router.get("/payment-stats")
|
||||
async def get_payment_stats(
|
||||
admin: User = Depends(get_admin_user),
|
||||
db: AsyncSession = Depends(get_db),
|
||||
):
|
||||
"""Payment statistics for admin dashboard."""
|
||||
from app.models.payment_order import PaymentOrder
|
||||
from datetime import datetime
|
||||
|
||||
# Count and revenue by status
|
||||
rows = (await db.execute(
|
||||
select(
|
||||
PaymentOrder.status,
|
||||
PaymentOrder.payment_method,
|
||||
func.count(PaymentOrder.id).label("count"),
|
||||
func.coalesce(func.sum(PaymentOrder.amount), 0).label("total_amount"),
|
||||
).group_by(PaymentOrder.status, PaymentOrder.payment_method)
|
||||
)).all()
|
||||
|
||||
by_status: dict[str, dict] = {}
|
||||
for r in rows:
|
||||
s = r.status
|
||||
if s not in by_status:
|
||||
by_status[s] = {"count": 0, "amount": 0.0}
|
||||
by_status[s]["count"] += r.count
|
||||
by_status[s]["amount"] += float(r.total_amount)
|
||||
|
||||
# Recent orders (last 50)
|
||||
recent = (await db.execute(
|
||||
select(PaymentOrder)
|
||||
.order_by(PaymentOrder.created_at.desc())
|
||||
.limit(50)
|
||||
)).scalars().all()
|
||||
|
||||
# Today stats
|
||||
today_start = datetime.now().replace(hour=0, minute=0, second=0, microsecond=0)
|
||||
today_paid = (await db.execute(
|
||||
select(
|
||||
func.count(PaymentOrder.id),
|
||||
func.coalesce(func.sum(PaymentOrder.amount), 0),
|
||||
).where(
|
||||
PaymentOrder.status == "paid",
|
||||
PaymentOrder.paid_at >= today_start,
|
||||
)
|
||||
)).first()
|
||||
today_count, today_amount = (today_paid or (0, 0))
|
||||
|
||||
return {
|
||||
"by_status": by_status,
|
||||
"today": {
|
||||
"paid_count": int(today_count or 0),
|
||||
"paid_amount": float(today_amount or 0),
|
||||
},
|
||||
"recent": [
|
||||
{
|
||||
"id": o.id,
|
||||
"order_no": o.order_no,
|
||||
"user_id": o.user_id,
|
||||
"amount": o.amount,
|
||||
"credits": o.credits,
|
||||
"payment_method": o.payment_method,
|
||||
"status": o.status,
|
||||
"trade_no": o.trade_no,
|
||||
"created_at": _iso(o.created_at),
|
||||
"paid_at": _iso(o.paid_at),
|
||||
}
|
||||
for o in recent
|
||||
],
|
||||
}
|
||||
|
||||
@@ -89,7 +89,10 @@ async def wechat_callback(request: Request, db: AsyncSession = Depends(get_db)):
|
||||
async def alipay_callback(request: Request, db: AsyncSession = Depends(get_db)):
|
||||
form_data = await request.form()
|
||||
data = dict(form_data)
|
||||
logger.info(f"Alipay callback received: {list(data.keys())}")
|
||||
logger.info(
|
||||
f"ALIPAY_CALLBACK order_no={data.get('out_trade_no')} "
|
||||
f"trade_no={data.get('trade_no', '')} status={data.get('trade_status', '')}"
|
||||
)
|
||||
|
||||
# Verify signature first
|
||||
if not await verify_alipay_callback(data, db):
|
||||
@@ -114,9 +117,40 @@ async def list_orders(
|
||||
current_user: User = Depends(get_current_user),
|
||||
db: AsyncSession = Depends(get_db),
|
||||
):
|
||||
# Auto-expire stale pending orders before returning
|
||||
from app.services.payment import _check_and_expire_order
|
||||
result = await db.execute(
|
||||
select(PaymentOrder)
|
||||
.where(PaymentOrder.user_id == current_user.id)
|
||||
.order_by(PaymentOrder.created_at.desc())
|
||||
)
|
||||
return result.scalars().all()
|
||||
orders = result.scalars().all()
|
||||
for o in orders:
|
||||
await _check_and_expire_order(db, o)
|
||||
return orders
|
||||
|
||||
|
||||
@router.post("/orders/{order_no}/cancel")
|
||||
async def cancel_order(
|
||||
order_no: str,
|
||||
current_user: User = Depends(get_current_user),
|
||||
db: AsyncSession = Depends(get_db),
|
||||
):
|
||||
"""Cancel a pending order. Only the order owner can cancel, only if still pending."""
|
||||
result = await db.execute(
|
||||
select(PaymentOrder).where(
|
||||
PaymentOrder.order_no == order_no,
|
||||
PaymentOrder.user_id == current_user.id,
|
||||
).limit(1)
|
||||
)
|
||||
order = result.scalar_one_or_none()
|
||||
if not order:
|
||||
raise HTTPException(status_code=404, detail="订单不存在")
|
||||
if order.status != "pending":
|
||||
raise HTTPException(status_code=400, detail=f"订单状态为{order.status},无法取消")
|
||||
order.status = "cancelled"
|
||||
await db.flush()
|
||||
logger.info(
|
||||
f"ORDER_CANCELLED order_no={order_no} user={current_user.id} amount={order.amount}"
|
||||
)
|
||||
return {"ok": True}
|
||||
|
||||
@@ -57,7 +57,7 @@ class Settings(BaseSettings):
|
||||
ALIPAY_PRIVATE_KEY: str = ""
|
||||
ALIPAY_PUBLIC_KEY: str = ""
|
||||
ALIPAY_NOTIFY_URL: str = ""
|
||||
PAYMENT_MOCK: bool = True
|
||||
PAYMENT_MOCK: bool = False # Default off; use admin panel to enable for testing
|
||||
|
||||
STORAGE_TYPE: str = "local"
|
||||
STORAGE_LOCAL_PATH: str = "./storage/generate/videos"
|
||||
|
||||
@@ -21,6 +21,48 @@ from app.services.log_config import decrypt_data
|
||||
logging.basicConfig(level=logging.INFO if settings.DEBUG else logging.WARNING)
|
||||
|
||||
|
||||
def _setup_payment_logger():
|
||||
"""Configure a dedicated file logger for payment events.
|
||||
|
||||
Logs are written to logs/payment_YYYY-MM-DD.log, rotated daily.
|
||||
30 days of history are retained.
|
||||
"""
|
||||
from logging.handlers import TimedRotatingFileHandler
|
||||
log_dir = os.path.join(os.path.dirname(os.path.dirname(__file__)), "logs")
|
||||
os.makedirs(log_dir, exist_ok=True)
|
||||
log_file = os.path.join(log_dir, "payment.log")
|
||||
|
||||
payment_logger = logging.getLogger("payment")
|
||||
payment_logger.setLevel(logging.INFO)
|
||||
payment_logger.propagate = False # don't double-log to root
|
||||
|
||||
# Avoid adding duplicate handlers on reload
|
||||
if any(getattr(h, "_payment_file", False) for h in payment_logger.handlers):
|
||||
return
|
||||
|
||||
handler = TimedRotatingFileHandler(
|
||||
log_file,
|
||||
when="midnight",
|
||||
interval=1,
|
||||
backupCount=30,
|
||||
encoding="utf-8",
|
||||
utc=False, # use local time
|
||||
)
|
||||
handler.suffix = "%Y-%m-%d" # files named like payment.log.2026-06-10
|
||||
handler._payment_file = True # type: ignore[attr-defined]
|
||||
handler.setFormatter(logging.Formatter(
|
||||
"%(asctime)s [%(levelname)s] %(message)s",
|
||||
datefmt="%Y-%m-%d %H:%M:%S",
|
||||
))
|
||||
payment_logger.addHandler(handler)
|
||||
# Mirror to console in DEBUG mode
|
||||
if settings.DEBUG:
|
||||
payment_logger.addHandler(logging.StreamHandler())
|
||||
|
||||
|
||||
_setup_payment_logger()
|
||||
|
||||
|
||||
@asynccontextmanager
|
||||
async def lifespan(app: FastAPI):
|
||||
from app.models import async_session
|
||||
@@ -35,13 +77,31 @@ async def lifespan(app: FastAPI):
|
||||
from app.services.video_queue import task_queue
|
||||
await task_queue.recover()
|
||||
queue_task = asyncio.create_task(task_queue.run())
|
||||
|
||||
|
||||
# Background task: auto-expire pending payment orders
|
||||
async def _order_expiry_loop():
|
||||
from app.services.payment import expire_all_pending_orders
|
||||
from logging import getLogger
|
||||
bg_logger = getLogger("payment")
|
||||
while True:
|
||||
try:
|
||||
async with async_session() as db:
|
||||
n = await expire_all_pending_orders(db)
|
||||
if n > 0:
|
||||
bg_logger.info(f"Auto-expired {n} pending payment order(s)")
|
||||
except Exception as e:
|
||||
bg_logger.error(f"Order expiry loop error: {e}")
|
||||
await asyncio.sleep(60) # check every minute
|
||||
|
||||
expiry_task = asyncio.create_task(_order_expiry_loop())
|
||||
|
||||
app.state.db_session_factory = async_session
|
||||
|
||||
yield
|
||||
|
||||
task_queue.stop()
|
||||
await queue_task
|
||||
expiry_task.cancel()
|
||||
await close_database()
|
||||
await close_redis()
|
||||
|
||||
|
||||
@@ -22,3 +22,9 @@ class PaymentOrder(Base, TimestampMixin):
|
||||
DateTime(timezone=True), nullable=True
|
||||
)
|
||||
trade_no: Mapped[str | None] = mapped_column(String(128), nullable=True)
|
||||
# Refund fields
|
||||
refund_trade_no: Mapped[str | None] = mapped_column(String(128), nullable=True)
|
||||
refunded_at: Mapped[datetime | None] = mapped_column(
|
||||
DateTime(timezone=True), nullable=True
|
||||
)
|
||||
refund_amount: Mapped[float | None] = mapped_column(Float, nullable=True)
|
||||
|
||||
@@ -1,5 +1,7 @@
|
||||
import logging
|
||||
from datetime import datetime
|
||||
import os
|
||||
from datetime import datetime, timedelta
|
||||
from logging.handlers import TimedRotatingFileHandler
|
||||
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
@@ -10,7 +12,32 @@ from app.models.system_config import SystemConfig
|
||||
from app.services.credits import add_credits
|
||||
from app.utils.id_gen import generate_id, generate_order_no
|
||||
|
||||
logger = logging.getLogger("videogen")
|
||||
# ---------------------------------------------------------------------------
|
||||
# Payment logger → logs/payment/YYYY-MM-DD.log (one file per day, keep 30 days)
|
||||
# ---------------------------------------------------------------------------
|
||||
logger = logging.getLogger("payment")
|
||||
logger.setLevel(logging.INFO)
|
||||
|
||||
_log_dir = os.path.join(os.path.dirname(os.path.dirname(__file__)), "logs", "payment")
|
||||
os.makedirs(_log_dir, exist_ok=True)
|
||||
|
||||
_file_handler = TimedRotatingFileHandler(
|
||||
os.path.join(_log_dir, "payment.log"),
|
||||
when="midnight",
|
||||
interval=1,
|
||||
backupCount=0,
|
||||
encoding="utf-8",
|
||||
utc=False,
|
||||
)
|
||||
_file_handler.suffix = "%Y-%m-%d"
|
||||
_file_handler.setFormatter(logging.Formatter(
|
||||
"[%(asctime)s] %(levelname)s %(message)s", datefmt="%Y-%m-%d %H:%M:%S"
|
||||
))
|
||||
if not logger.handlers:
|
||||
logger.addHandler(_file_handler)
|
||||
|
||||
# Orders pending payment for longer than this are auto-cancelled
|
||||
ORDER_EXPIRE_MINUTES = 5
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
@@ -26,6 +53,46 @@ async def _get_payment_configs(db: AsyncSession) -> dict[str, str]:
|
||||
return {c.key: c.value for c in result.scalars().all()}
|
||||
|
||||
|
||||
async def _check_and_expire_order(db: AsyncSession, order: PaymentOrder) -> bool:
|
||||
"""If a pending order has passed its expiry, mark it cancelled.
|
||||
Returns True if the order was expired.
|
||||
"""
|
||||
if order.status != "pending":
|
||||
return False
|
||||
expiry = order.created_at + timedelta(minutes=ORDER_EXPIRE_MINUTES)
|
||||
if datetime.now(order.created_at.tzinfo) >= expiry:
|
||||
order.status = "cancelled"
|
||||
await db.flush()
|
||||
logger.info(
|
||||
f"ORDER_EXPIRED order_no={order.order_no} user={order.user_id} "
|
||||
f"amount={order.amount} created_at={order.created_at.isoformat()}"
|
||||
)
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
async def expire_all_pending_orders(db: AsyncSession) -> int:
|
||||
"""Background task: mark all expired pending orders as cancelled.
|
||||
Returns the number of orders expired.
|
||||
"""
|
||||
threshold = datetime.now() - timedelta(minutes=ORDER_EXPIRE_MINUTES)
|
||||
result = await db.execute(
|
||||
select(PaymentOrder).where(
|
||||
PaymentOrder.status == "pending",
|
||||
PaymentOrder.created_at <= threshold,
|
||||
)
|
||||
)
|
||||
orders = result.scalars().all()
|
||||
for o in orders:
|
||||
o.status = "cancelled"
|
||||
logger.info(
|
||||
f"ORDER_EXPIRED order_no={o.order_no} user={o.user_id} amount={o.amount}"
|
||||
)
|
||||
if orders:
|
||||
await db.flush()
|
||||
return len(orders)
|
||||
|
||||
|
||||
def _is_mock_mode(db_configs: dict[str, str]) -> bool:
|
||||
"""Check if payment mock mode is enabled (from DB or env)."""
|
||||
db_val = db_configs.get("payment_mock", "")
|
||||
@@ -125,6 +192,10 @@ async def create_recharge_order(
|
||||
)
|
||||
db.add(order)
|
||||
await db.flush()
|
||||
logger.info(
|
||||
f"ORDER_CREATED order_no={order.order_no} user={user_id} "
|
||||
f"amount={price} credits={total_credits} method={method} mock={mock_mode}"
|
||||
)
|
||||
|
||||
if mock_mode:
|
||||
# Mock: immediately complete payment
|
||||
@@ -371,4 +442,7 @@ async def process_payment_success_by_order_no(db: AsyncSession, order_no: str, t
|
||||
related_id=order.id,
|
||||
)
|
||||
await db.flush()
|
||||
logger.info(f"Payment success processed: order_no={order_no}, trade_no={trade_no}")
|
||||
logger.info(
|
||||
f"PAYMENT_SUCCESS order_no={order_no} user={order.user_id} "
|
||||
f"amount={order.amount} credits={order.credits} trade_no={trade_no}"
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user