"""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, platforms: list[str], sid: str | None = None, member_count: int | None = None, remark: str | None = None, commit: bool = True, ) -> CpsGroup: """建群。含 meituan 才有 sid(留空自动 g);纯淘宝/京东 sid=None(那俩不支持 sid)。""" has_meituan = "meituan" in (platforms or []) # 含美团:用填的 sid,或临时占位待回填 g;不含美团:sid 恒 None seed_sid = (sid or f"tmp{uuid4().hex[:20]}") if has_meituan else None group = CpsGroup( name=name, platforms=list(platforms or []), sid=seed_sid, member_count=member_count, remark=remark, ) db.add(group) db.flush() if has_meituan and 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, platforms: list[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 platforms is not None: group.platforms = list(platforms) 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 list_activity_images(db: Session) -> list[str]: """已用过的活动落地页图(distinct 非空 image_url),供新建活动复用选择。""" rows = ( db.execute( select(CpsActivity.image_url) .where(CpsActivity.image_url.is_not(None)) .distinct() ) .scalars() .all() ) return [r for r in rows if r] def create_activity( db: Session, *, name: str, platform: str = "meituan", act_id: str | None = None, product_view_sign: str | None = None, payload: str | None = None, image_url: 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, payload=payload, image_url=image_url, 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。按预估佣金降序。 """ ostmt = select(CpsOrder) if date_from is not None: ostmt = ostmt.where(CpsOrder.pay_time >= _as_utc(date_from)) if date_to is not None: ostmt = ostmt.where(CpsOrder.pay_time <= _as_utc(date_to)) orders_by_sid: dict[str | None, list[CpsOrder]] = {} for o in db.execute(ostmt).scalars().all(): orders_by_sid.setdefault(o.sid, []).append(o) clicks = cps_link_repo.click_stats_by_group(db, date_from=date_from, date_to=date_to) groups = list(db.execute(select(CpsGroup)).scalars().all()) def _order_agg(items: list[CpsOrder]) -> dict: 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] return { "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), } # 淘宝/京东无法对账 → 对账字段全 None(前端显示 "-") _no_recon = { "order_count": None, "settled_count": None, "canceled_count": None, "gmv_cents": None, "est_commission_cents": None, "settled_commission_cents": None, } rows: list[dict] = [] seen_sids: set[str] = set() for g in groups: if g.status != "active": continue c = clicks.get(g.id) or {} row = { "group_id": g.id, "sid": g.sid, "name": g.name, "platforms": list(g.platforms or []), "member_count": g.member_count, "click_pv": c.get("pv", 0), "click_uv": c.get("uv", 0), "copy_pv": c.get("copy_pv", 0), "copy_uv": c.get("copy_uv", 0), } # 对账只对美团群(有 sid);淘宝/京东 → None if "meituan" in (g.platforms or []) and g.sid: row.update(_order_agg(orders_by_sid.get(g.sid, []))) seen_sids.add(g.sid) else: row.update(_no_recon) rows.append(row) # 未归群的历史美团 sid(如 wonderableai):单列,有对账无点击 for sid, items in orders_by_sid.items(): if sid is None or sid in seen_sids: continue row = { "group_id": None, "sid": sid, "name": sid, "platforms": ["meituan"], "member_count": None, "click_pv": 0, "click_uv": 0, "copy_pv": 0, "copy_uv": 0, } row.update(_order_agg(items)) rows.append(row) # 佣金降序(None 当 0),次按点击 rows.sort(key=lambda x: ((x["est_commission_cents"] or 0), x["click_pv"]), reverse=True) return rows