Files
shaguabijia-app-server/app/admin/repositories/cps.py
T
marco 277f9b16a2 feat(cps): 群发券分发与对账 + 短链点击追踪
群(sid)/活动池/转链(带sid)/美团 query_order 按 sid 对账/按群统计;
短链 /c/{code} 记点击(PV/UV)→302 跳美团,点击→下单漏斗。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-17 10:00:29 +08:00

288 lines
11 KiB
Python

"""admin CPS 数据访问:群/活动 CRUD + 美团订单拉单入库(对账) + 按群统计聚合。
美团 query_order 返回的金额是「元」字符串、时间是秒级时间戳,这里统一转「分」+ tz-aware。
统计在 Python 侧聚合(订单量级小、admin 低频),逻辑清晰、跨 PG/SQLite 无 SQL 方言坑。
"""
from __future__ import annotations
from datetime import datetime, timezone
from decimal import Decimal, InvalidOperation
from uuid import uuid4
from sqlalchemy import desc, select
from sqlalchemy.orm import Session
from app.admin.repositories.queries import _as_utc, offset_paginate
from app.integrations import meituan
from app.repositories import cps_link as cps_link_repo
from app.models.cps_activity import CpsActivity
from app.models.cps_group import CpsGroup
from app.models.cps_order import CpsOrder
# 美团订单状态:取消(4)/风控(5)不计佣金;结算(6)为佣金真正到账
_INVALID_STATUS = {"4", "5"}
_SETTLED_STATUS = "6"
# ───────────── 单位换算 ─────────────
def _yuan_to_cents(v: object) -> int | None:
"""「元」字符串/数 → 分。None / 空 / 字面 "null" → None。"""
if v is None:
return None
s = str(v).strip()
if not s or s.lower() == "null":
return None
try:
return int((Decimal(s) * 100).to_integral_value())
except (InvalidOperation, ValueError):
return None
def _ts_to_dt(ts: object) -> datetime | None:
"""秒级时间戳 → tz-aware UTC datetime(绝对时刻,前端按北京展示)。"""
if not ts:
return None
try:
return datetime.fromtimestamp(int(ts), tz=timezone.utc)
except (ValueError, OSError, TypeError):
return None
# ───────────── 群 ─────────────
def get_group(db: Session, group_id: int) -> CpsGroup | None:
return db.get(CpsGroup, group_id)
def get_group_by_sid(db: Session, sid: str) -> CpsGroup | None:
return db.execute(select(CpsGroup).where(CpsGroup.sid == sid)).scalar_one_or_none()
def list_groups(
db: Session, *, keyword: str | None = None, status: str | None = None,
limit: int = 20, cursor: int | None = None,
) -> tuple[list[CpsGroup], int | None, int]:
stmt = select(CpsGroup)
if keyword and keyword.strip():
kw = f"%{keyword.strip()}%"
stmt = stmt.where(CpsGroup.name.ilike(kw) | CpsGroup.sid.ilike(kw))
if status:
stmt = stmt.where(CpsGroup.status == status)
return offset_paginate(db, stmt, (desc(CpsGroup.id),), limit=limit, cursor=cursor)
def create_group(
db: Session, *, name: str, sid: str | None = None,
member_count: int | None = None, remark: str | None = None, commit: bool = True,
) -> CpsGroup:
"""建群。sid 留空 → 自动 g<id>(先 flush 拿 id 再回填)。"""
group = CpsGroup(
name=name,
sid=sid or f"tmp{uuid4().hex[:20]}", # 临时唯一占位,留空时下面回填 g<id>
member_count=member_count,
remark=remark,
)
db.add(group)
db.flush()
if not sid:
group.sid = f"g{group.id}"
db.flush()
if commit:
db.commit()
db.refresh(group)
return group
def update_group(
db: Session, group: CpsGroup, *, name: str | None = None,
member_count: int | None = None, status: str | None = None,
remark: str | None = None, commit: bool = True,
) -> CpsGroup:
if name is not None:
group.name = name
if member_count is not None:
group.member_count = member_count
if status is not None:
group.status = status
if remark is not None:
group.remark = remark
if commit:
db.commit()
db.refresh(group)
else:
db.flush()
return group
# ───────────── 活动 ─────────────
def get_activity(db: Session, activity_id: int) -> CpsActivity | None:
return db.get(CpsActivity, activity_id)
def list_activities(
db: Session, *, platform: str | None = None, status: str | None = None,
limit: int = 20, cursor: int | None = None,
) -> tuple[list[CpsActivity], int | None, int]:
stmt = select(CpsActivity)
if platform:
stmt = stmt.where(CpsActivity.platform == platform)
if status:
stmt = stmt.where(CpsActivity.status == status)
return offset_paginate(db, stmt, (desc(CpsActivity.id),), limit=limit, cursor=cursor)
def create_activity(
db: Session, *, name: str, platform: str = "meituan", act_id: str | None = None,
product_view_sign: str | None = None, remark: str | None = None, commit: bool = True,
) -> CpsActivity:
activity = CpsActivity(
name=name, platform=platform, act_id=act_id,
product_view_sign=product_view_sign, remark=remark,
)
db.add(activity)
if commit:
db.commit()
db.refresh(activity)
else:
db.flush()
return activity
# ───────────── 订单(对账) ─────────────
def _map_order_fields(r: dict) -> dict:
"""美团 query_order 单条 dataList → CpsOrder 字段(金额转分、时间转 datetime)。"""
pn = r.get("productName")
if pn and len(pn) > 500:
pn = pn[:500]
return {
"sid": (r.get("sid") or None),
"act_id": str(r["actId"]) if r.get("actId") is not None else None,
"biz_line": r.get("businessLine"),
"trade_type": r.get("tradeType"),
"pay_price_cents": _yuan_to_cents(r.get("payPrice")),
"commission_cents": _yuan_to_cents(r.get("profit")),
"commission_rate": str(r["commissionRate"]) if r.get("commissionRate") is not None else None,
"refund_price_cents": _yuan_to_cents(r.get("refundPrice")),
"refund_profit_cents": _yuan_to_cents(r.get("refundProfit")),
"mt_status": str(r["status"]) if r.get("status") is not None else None,
"invalid_reason": (r.get("invalidReason") or None),
"product_name": pn or None,
"pay_time": _ts_to_dt(r.get("payTime")),
"mt_update_time": _ts_to_dt(r.get("updateTime")),
"raw": r,
}
def reconcile_orders(
db: Session, *, start_time: int, end_time: int,
query_time_type: int = 1, sid: str | None = None, max_pages: int = 200,
) -> dict:
"""调美团 query_order 分页拉单 → 按 order_id upsert。返回 {fetched, inserted, updated, pages}。
订单状态会随时间变(付款→完成→结算/退款),重复拉同一单则更新。max_pages 防异常死循环。
"""
fetched = inserted = updated = pages = 0
page = 1
while page <= max_pages:
resp = meituan.query_order(
sid=sid, start_time=start_time, end_time=end_time,
query_time_type=query_time_type, page=page, limit=100,
)
rows = ((resp.get("data") or {}).get("dataList")) or []
if not rows:
break
pages += 1
for r in rows:
order_id = str(r.get("orderId") or "").strip()
if not order_id:
continue
fetched += 1
fields = _map_order_fields(r)
existing = db.execute(
select(CpsOrder).where(CpsOrder.order_id == order_id)
).scalar_one_or_none()
if existing is None:
db.add(CpsOrder(order_id=order_id, **fields))
inserted += 1
else:
for k, v in fields.items():
setattr(existing, k, v)
updated += 1
if len(rows) < 100:
break
page += 1
db.commit()
return {"fetched": fetched, "inserted": inserted, "updated": updated, "pages": pages}
def list_orders(
db: Session, *, sid: str | None = None, mt_status: str | None = None,
limit: int = 20, cursor: int | None = None,
) -> tuple[list[CpsOrder], int | None, int]:
stmt = select(CpsOrder)
if sid:
stmt = stmt.where(CpsOrder.sid == sid)
if mt_status:
stmt = stmt.where(CpsOrder.mt_status == mt_status)
return offset_paginate(db, stmt, (desc(CpsOrder.id),), limit=limit, cursor=cursor)
# ───────────── 统计(按群聚合) ─────────────
def group_stats(
db: Session, *, date_from: datetime | None = None, date_to: datetime | None = None,
) -> list[dict]:
"""按 sid 聚合订单 + join 群信息。已建但本期无单的活跃群也列出(全 0)。
未归群的 sid(历史/其它来源)单独成行 group_id=None。按预估佣金降序。
"""
stmt = select(CpsOrder)
if date_from is not None:
stmt = stmt.where(CpsOrder.pay_time >= _as_utc(date_from))
if date_to is not None:
stmt = stmt.where(CpsOrder.pay_time <= _as_utc(date_to))
orders = list(db.execute(stmt).scalars().all())
groups = {g.sid: g for g in db.execute(select(CpsGroup)).scalars().all()}
buckets: dict[str | None, list[CpsOrder]] = {}
for o in orders:
buckets.setdefault(o.sid, []).append(o)
rows: list[dict] = []
for sid, items in buckets.items():
g = groups.get(sid) if sid else None
valid = [o for o in items if o.mt_status not in _INVALID_STATUS]
settled = [o for o in items if o.mt_status == _SETTLED_STATUS]
canceled = [o for o in items if o.mt_status in _INVALID_STATUS]
rows.append({
"group_id": g.id if g else None,
"sid": sid,
"name": g.name if g else (sid or "(无 sid 归属)"),
"member_count": g.member_count if g else None,
"order_count": len(valid),
"settled_count": len(settled),
"canceled_count": len(canceled),
"gmv_cents": sum(o.pay_price_cents or 0 for o in valid),
"est_commission_cents": sum(o.commission_cents or 0 for o in valid),
"settled_commission_cents": sum(o.commission_cents or 0 for o in settled),
})
# 已建的活跃群但本期无单 → 补 0 行,让运营看到全部群
seen = set(buckets.keys())
for sid, g in groups.items():
if sid not in seen and g.status == "active":
rows.append({
"group_id": g.id, "sid": sid, "name": g.name,
"member_count": g.member_count, "order_count": 0,
"settled_count": 0, "canceled_count": 0, "gmv_cents": 0,
"est_commission_cents": 0, "settled_commission_cents": 0,
})
# 合并点击数据(按 group_id):点击只对生成过我们短链的群有,未归群行(group_id=None)恒 0
clicks = cps_link_repo.click_stats_by_group(db, date_from=date_from, date_to=date_to)
for r in rows:
c = clicks.get(r["group_id"]) if r["group_id"] is not None else None
r["click_pv"] = c["pv"] if c else 0
r["click_uv"] = c["uv"] if c else 0
rows.sort(key=lambda x: x["est_commission_cents"], reverse=True)
return rows