"""排行榜奖励 — 俱乐部分榜结算、懒结算、领取入账(严格幂等)。""" from __future__ import annotations import logging from datetime import date, datetime, timedelta from decimal import Decimal from typing import Optional from django.db import IntegrityError, transaction from django.db.models import F, Sum from django.utils import timezone from backend.models import ( LeaderDailyStats, ManagerDailyStats, MerchantDailyStats, PlayerDailyStats, ) from jituan.services.club_user import get_user_club_id from rank.reward_models import ( RankRewardClaim, RankRewardScheme, RankRewardSettlement, RankRewardTier, ) from users.business_models import User from users.models import UserDashou, UserGuanshi, UserShangjia, UserZuzhang logger = logging.getLogger(__name__) MAX_RANK = 50 RIQI_TO_PERIOD = { '昨日': RankRewardScheme.PERIOD_DAY, '上周': RankRewardScheme.PERIOD_WEEK, '上月': RankRewardScheme.PERIOD_MONTH, } OPEN_RIQI = frozenset({'今日', '本周', '本月', '总榜'}) OPEN_RIQI_PERIOD_MAP = { '今日': RankRewardScheme.PERIOD_DAY, '本周': RankRewardScheme.PERIOD_WEEK, '本月': RankRewardScheme.PERIOD_MONTH, } DEFAULT_SORT_FIELD = { 'dashou': 'chengjiao_zongliang', 'guanshi': 'chongzhi_dashou_shu', 'zuzhang': 'yaoqing_guanshi_shu', 'shangjia': 'jiesuan_jine', } # shenfen -> sort_field_key -> (model, id_field, db_field, is_int) SORT_FIELD_REGISTRY = { 'dashou': { 'chengjiao_zongliang': (PlayerDailyStats, 'PlayerID', 'CompletedOrderTotal', True), 'chengjiao_zonge': (PlayerDailyStats, 'PlayerID', 'CompletedAmount', False), 'jiedan_zongliang': (PlayerDailyStats, 'PlayerID', 'AcceptedOrderTotal', True), 'jiedan_zonge': (PlayerDailyStats, 'PlayerID', 'AcceptedAmount', False), }, 'guanshi': { 'chongzhi_dashou_shu': (ManagerDailyStats, 'ManagerID', 'RechargedPlayerCount', True), 'shouru_zonge': (ManagerDailyStats, 'ManagerID', 'TotalIncome', False), 'yaoqing_dashou_shu': (ManagerDailyStats, 'ManagerID', 'InvitedPlayerCount', True), }, 'zuzhang': { 'yaoqing_guanshi_shu': (LeaderDailyStats, 'LeaderID', 'InvitedManagerCount', True), 'shouru_zonge': (LeaderDailyStats, 'LeaderID', 'TotalIncome', False), 'fenyong_jine': (LeaderDailyStats, 'LeaderID', 'CommissionAmount', False), }, 'shangjia': { 'jiesuan_dingdan_shu': (MerchantDailyStats, 'MerchantID', 'SettledOrderCount', True), 'jiesuan_jine': (MerchantDailyStats, 'MerchantID', 'SettledAmount', False), 'paifa_dingdan_shu': (MerchantDailyStats, 'MerchantID', 'AssignedOrderCount', True), 'paifa_jine': (MerchantDailyStats, 'MerchantID', 'AssignedAmount', False), }, } SORT_FIELD_LABELS = { 'chengjiao_zongliang': '成交量', 'chengjiao_zonge': '成交总额', 'jiedan_zongliang': '接单量', 'jiedan_zonge': '接单总额', 'chongzhi_dashou_shu': '有效打手数', 'shouru_zonge': '收益总额', 'yaoqing_dashou_shu': '邀请打手数', 'yaoqing_guanshi_shu': '邀请管事数', 'fenyong_jine': '分佣金额', 'jiesuan_dingdan_shu': '结算单量', 'jiesuan_jine': '结算金额', 'paifa_dingdan_shu': '派单量', 'paifa_jine': '派单流水', } BALANCE_FIELD_MAP = { 'dashou': 'yue', 'guanshi': 'yue', 'zuzhang': 'ketixian_jine', 'shangjia': 'yue', } def sort_field_options_for_shenfen(shenfen: str) -> list: reg = SORT_FIELD_REGISTRY.get(shenfen) or {} return [ {'key': k, 'label': SORT_FIELD_LABELS.get(k, k)} for k in reg.keys() ] def resolve_sort_field_for_display(club_id: str, shenfen: str, riqi: str) -> tuple[str, str]: """按后台奖励方案决定展示/排序字段(方案存在即生效,与是否启用奖励无关)。""" period_type = RIQI_TO_PERIOD.get(riqi) or OPEN_RIQI_PERIOD_MAP.get(riqi) sort_field = None if period_type: scheme = RankRewardScheme.objects.filter( club_id=club_id, shenfen=shenfen, period_type=period_type, ).first() if scheme: reg = SORT_FIELD_REGISTRY.get(shenfen) or {} if scheme.sort_field in reg: sort_field = scheme.sort_field if not sort_field: sort_field = DEFAULT_SORT_FIELD.get(shenfen, 'chengjiao_zongliang') return sort_field, SORT_FIELD_LABELS.get(sort_field, sort_field) def club_member_uids(club_id: str) -> list: """仅统计归属该俱乐部的用户(未开业/无用户则空榜)。""" if not club_id: return [] return list( User.objects.filter(ClubID=club_id).values_list('UserUID', flat=True) ) def filter_rank_rows_by_user_club(rows: list, club_id: str) -> list: if not rows or not club_id: return rows or [] club_map = load_user_club_map([r['yonghuid'] for r in rows]) return [r for r in rows if club_map.get(r['yonghuid']) == club_id] def resolve_period_dates(riqi: str) -> Optional[tuple[date, date, str]]: today = date.today() if riqi == '昨日': d = today - timedelta(days=1) return d, d, f'day:{d.isoformat()}' if riqi == '上周': this_monday = today - timedelta(days=today.weekday()) last_sunday = this_monday - timedelta(days=1) last_monday = last_sunday - timedelta(days=6) iso = last_monday.isocalendar() return last_monday, last_sunday, f'week:{iso[0]}-W{iso[1]:02d}' if riqi == '上月': first_this = today.replace(day=1) last_day = first_this - timedelta(days=1) first_prev = last_day.replace(day=1) return first_prev, last_day, f'month:{first_prev.strftime("%Y-%m")}' return None def _resolve_date_range_for_display(riqi: str): """与 paihang_views 一致,用于集团/俱乐部展示榜。""" today = date.today() if riqi == '今日': return today, today if riqi == '昨日': d = today - timedelta(days=1) return d, d if riqi == '本周': monday = today - timedelta(days=today.weekday()) return monday, today if riqi == '上周': this_monday = today - timedelta(days=today.weekday()) last_sunday = this_monday - timedelta(days=1) last_monday = last_sunday - timedelta(days=6) return last_monday, last_sunday if riqi == '本月': return today.replace(day=1), today if riqi == '上月': first_this = today.replace(day=1) last_day = first_this - timedelta(days=1) return last_day.replace(day=1), last_day return None def query_club_rank_rows(club_id: str, shenfen: str, sort_field: str, riqi: str, limit: int = MAX_RANK): """按 club 归属用户 + 日统计 club_id 双过滤(分奖与展示唯一依据)。""" reg = SORT_FIELD_REGISTRY.get(shenfen, {}).get(sort_field) if not reg: return [] model, id_field, db_field, is_int = reg member_uids = club_member_uids(club_id) if not member_uids: return [] fetch_limit = min(max(limit * 5, 100), 250) def _pack(qs_or_rows, is_values=False): out = [] if is_values: for r in qs_or_rows: uid = r[id_field] metric = r.get('metric') or 0 if metric and uid in member_uids: out.append({'yonghuid': uid, 'metric': metric, 'is_int': is_int}) else: for r in qs_or_rows: uid = getattr(r, id_field) metric = getattr(r, db_field) or 0 if metric and uid in member_uids: out.append({'yonghuid': uid, 'metric': metric, 'is_int': is_int}) out = filter_rank_rows_by_user_club(out, club_id) out.sort(key=lambda x: x['metric'], reverse=True) return out[:limit] if riqi == '总榜': rows = ( model.objects.filter(club_id=club_id, **{f'{id_field}__in': member_uids}) .values(id_field) .annotate(metric=Sum(db_field)) .order_by(f'-metric')[:fetch_limit] ) return _pack(rows, is_values=True) date_range = _resolve_date_range_for_display(riqi) if not date_range: return [] start_date, end_date = date_range if start_date == end_date: qs = ( model.objects.filter( club_id=club_id, Date=start_date, **{f'{id_field}__in': member_uids}, ) .order_by(f'-{db_field}')[:fetch_limit] ) return _pack(qs, is_values=False) rows = ( model.objects.filter( club_id=club_id, Date__gte=start_date, Date__lte=end_date, **{f'{id_field}__in': member_uids}, ) .values(id_field) .annotate(metric=Sum(db_field)) .order_by('-metric')[:fetch_limit] ) return _pack(rows, is_values=True) def query_group_rank_rows(shenfen: str, sort_field: str, riqi: str, limit: int = MAX_RANK): """集团总榜(仅展示,不过滤 club)。""" return _query_group_rank_rows(shenfen, sort_field, riqi, limit) def _query_group_rank_rows(shenfen: str, sort_field: str, riqi: str, limit: int = MAX_RANK): reg = SORT_FIELD_REGISTRY.get(shenfen, {}).get(sort_field) if not reg: return [] model, id_field, db_field, is_int = reg if riqi == '总榜': rows = ( model.objects.values(id_field) .annotate(metric=Sum(db_field)) .order_by('-metric')[:limit] ) return [ {'yonghuid': r[id_field], 'metric': r['metric'] or 0, 'is_int': is_int} for r in rows if r.get('metric') ] date_range = _resolve_date_range_for_display(riqi) if not date_range: return [] start_date, end_date = date_range if start_date == end_date: qs = model.objects.filter(Date=start_date).order_by(f'-{db_field}')[:limit] return [ {'yonghuid': getattr(r, id_field), 'metric': getattr(r, db_field) or 0, 'is_int': is_int} for r in qs if getattr(r, db_field) ] rows = ( model.objects.filter(Date__gte=start_date, Date__lte=end_date) .values(id_field) .annotate(metric=Sum(db_field)) .order_by('-metric')[:limit] ) return [ {'yonghuid': r[id_field], 'metric': r['metric'] or 0, 'is_int': is_int} for r in rows if r.get('metric') ] def _fmt_stat_value(val, is_int: bool): if val is None: return 0 if is_int else 0.0 if is_int: return int(val) if isinstance(val, Decimal): return float(val.quantize(Decimal('0.01'))) return float(val) def _enrich_today_profiles(shenfen: str, yonghuids: list) -> dict: """今日日统计为空时,从 Profile 补全展示字段(与 phbhqsj 回退一致)。""" from users.models import UserDashou, UserGuanshi, UserShangjia, UserZuzhang from users.paihang_views import ROLE_CONFIG cfg = ROLE_CONFIG[shenfen] response_field_map = cfg['response_field_map'] int_fields = cfg['int_fields'] id_field = cfg['id_field'] profile_cfg = { 'dashou': (UserDashou, { 'AcceptedOrderTotal': 'jinrijiedan', 'CompletedOrderTotal': 'jinrijiedan', 'AcceptedAmount': 'jinrishouyi', 'CompletedAmount': 'jinrishouyi', }), 'guanshi': (UserGuanshi, { 'InvitedPlayerCount': None, 'RechargedPlayerCount': 'jinrichongzhi', 'TotalIncome': None, }), 'zuzhang': (UserZuzhang, { 'InvitedManagerCount': None, 'CommissionAmount': 'jinri_fenyong', 'TotalIncome': 'jinri_fenyong', }), 'shangjia': (UserShangjia, { 'AssignedOrderCount': 'jinridingdan', 'AssignedAmount': 'jinriliushui', 'SettledOrderCount': None, 'SettledAmount': None, }), } model_cls, field_map = profile_cfg.get(shenfen, (None, {})) if not model_cls: return {} out = {} qs = model_cls.objects.filter(user__UserUID__in=yonghuids).select_related('user') for p in qs: uid = p.user.UserUID row = {} for db_field, profile_field in field_map.items(): key = response_field_map.get(db_field, db_field) if profile_field: val = getattr(p, profile_field, 0) or 0 else: val = 0 row[key] = _fmt_stat_value(val, db_field in int_fields) out[uid] = row return out def enrich_rank_stats(shenfen: str, yonghuids: list, riqi: str, club_id: str = '') -> dict: """加载榜单行完整统计(排序指标 + 其它单量/金额),避免前端只看到一个数。""" if not yonghuids: return {} from users.paihang_views import ROLE_CONFIG cfg = ROLE_CONFIG[shenfen] model = cfg['model'] id_field = cfg['id_field'] fields = cfg['fields'] int_fields = cfg['int_fields'] response_field_map = cfg['response_field_map'] base_filter = {f'{id_field}__in': yonghuids} if club_id and hasattr(model, 'club_id'): base_filter['club_id'] = club_id result = {uid: {} for uid in yonghuids} def _fill_from_row(uid, row_data): for f in fields: key = response_field_map.get(f, f) result[uid][key] = _fmt_stat_value(row_data.get(f), f in int_fields) if riqi == '总榜': rows = ( model.objects.filter(**base_filter) .values(id_field) .annotate(**{f: Sum(f) for f in fields}) ) for r in rows: uid = r[id_field] if uid in result: _fill_from_row(uid, r) return result date_range = _resolve_date_range_for_display(riqi) if not date_range: return result start_date, end_date = date_range if start_date == end_date: qs = model.objects.filter(Date=start_date, **base_filter) if qs.exists(): for r in qs: uid = getattr(r, id_field) if uid not in result: continue row_data = {f: getattr(r, f) for f in fields} _fill_from_row(uid, row_data) return result if riqi == '今日': return _enrich_today_profiles(shenfen, yonghuids) return result rows = ( model.objects.filter(Date__gte=start_date, Date__lte=end_date, **base_filter) .values(id_field) .annotate(**{f: Sum(f) for f in fields}) ) for r in rows: uid = r[id_field] if uid in result: _fill_from_row(uid, r) return result def load_user_club_map(yonghuids: list) -> dict: if not yonghuids: return {} users = User.objects.filter(UserUID__in=yonghuids) return {u.UserUID: get_user_club_id(u) or getattr(u, 'ClubID', None) or '' for u in users} def user_has_active_shenfen(yonghuid: str, shenfen: str) -> bool: if shenfen == 'dashou': return UserDashou.objects.filter(user__UserUID=yonghuid, zhanghaozhuangtai=1).exists() if shenfen == 'guanshi': return UserGuanshi.objects.filter(user__UserUID=yonghuid, zhuangtai=1).exists() if shenfen == 'zuzhang': return UserZuzhang.objects.filter(user__UserUID=yonghuid, zhuangtai=1).exists() if shenfen == 'shangjia': return UserShangjia.objects.filter(user__UserUID=yonghuid, zhuangtai=1).exists() return False def tier_amount_for_rank(tiers, mingci: int) -> Optional[Decimal]: for t in tiers: if t.rank_from <= mingci <= t.rank_to: return t.reward_amount return None def build_tier_rewards_preview(scheme) -> dict: """按方案档位生成各名次预计奖金(进行中周期展示用)。""" if not scheme: return {} tiers = list(RankRewardTier.query.filter(scheme_id=scheme.id).order_by('rank_from')) if not tiers: return {} preview = {} for mingci in range(1, MAX_RANK + 1): amt = tier_amount_for_rank(tiers, mingci) if amt and amt > 0: preview[mingci] = float(amt) return preview def maybe_settle_club_rewards(club_id: str, shenfen: str, riqi: str) -> Optional[RankRewardSettlement]: """懒结算:仅已结束周期;幂等。""" if riqi in OPEN_RIQI: return None period_type = RIQI_TO_PERIOD.get(riqi) if not period_type: return None period_info = resolve_period_dates(riqi) if not period_info: return None period_start, period_end, period_key = period_info scheme = RankRewardScheme.query.filter( club_id=club_id, shenfen=shenfen, period_type=period_type, enabled=True, ).first() if not scheme: return None if scheme.sort_field not in (SORT_FIELD_REGISTRY.get(shenfen) or {}): logger.warning('invalid sort_field %s for %s', scheme.sort_field, shenfen) return None existing = RankRewardSettlement.query.filter( club_id=club_id, scheme_id=scheme.id, period_key=period_key, ).first() if existing: return existing rows = query_club_rank_rows(club_id, shenfen, scheme.sort_field, riqi) tiers = list(RankRewardTier.query.filter(scheme_id=scheme.id).order_by('rank_from')) if not tiers: return None min_metric = scheme.min_metric or Decimal('0') club_map = load_user_club_map([r['yonghuid'] for r in rows]) expire_at = timezone.now() + timedelta(days=scheme.claim_days or 30) balance_field = BALANCE_FIELD_MAP[shenfen] claims_to_create = [] for idx, row in enumerate(rows): mingci = idx + 1 metric = Decimal(str(row['metric'] or 0)) if metric < min_metric: continue amount = tier_amount_for_rank(tiers, mingci) if not amount or amount <= 0: continue uid = row['yonghuid'] if club_map.get(uid) != club_id: continue if not user_has_active_shenfen(uid, shenfen): continue idem = f'{club_id}|{scheme.id}|{period_key}|{uid}' claims_to_create.append(RankRewardClaim( club_id=club_id, yonghuid=uid, shenfen=shenfen, mingci=mingci, metric_value=metric, reward_amount=amount, status=RankRewardClaim.STATUS_PENDING, expire_at=expire_at, balance_field=balance_field, idempotent_key=idem, )) try: with transaction.atomic(): settlement = RankRewardSettlement.query.create( scheme_id=scheme.id, club_id=club_id, period_key=period_key, period_start=period_start, period_end=period_end, riqi_label=riqi, sort_field=scheme.sort_field, claim_count=0, ) created = 0 for claim in claims_to_create: claim.settlement_id = settlement.id try: claim.save() created += 1 except IntegrityError: pass settlement.claim_count = created settlement.save(update_fields=['claim_count']) return settlement except IntegrityError: return RankRewardSettlement.query.filter( club_id=club_id, scheme_id=scheme.id, period_key=period_key, ).first() def build_reward_payload(club_id: str, shenfen: str, riqi: str, yonghuid: str = '') -> dict: """奖励规则 + 我的待领取/已领取状态。""" period_type = RIQI_TO_PERIOD.get(riqi) or OPEN_RIQI_PERIOD_MAP.get(riqi) display_sort, display_label = resolve_sort_field_for_display(club_id, shenfen, riqi) scheme = None if period_type: scheme = RankRewardScheme.query.filter( club_id=club_id, shenfen=shenfen, period_type=period_type, ).first() settlement = None if scheme and riqi not in OPEN_RIQI: maybe_settle_club_rewards(club_id, shenfen, riqi) period_info = resolve_period_dates(riqi) if period_info: _, _, period_key = period_info settlement = RankRewardSettlement.query.filter( club_id=club_id, scheme_id=scheme.id, period_key=period_key, ).first() tiers = [] if scheme: tiers = [ { 'rank_from': t.rank_from, 'rank_to': t.rank_to, 'amount': float(t.reward_amount), 'label': t.label or '', } for t in RankRewardTier.query.filter(scheme_id=scheme.id).order_by('rank_from') ] rank_rewards = {} if scheme and settlement: for c in RankRewardClaim.query.filter(settlement_id=settlement.id): rank_rewards[c.mingci] = float(c.reward_amount) tier_rewards = build_tier_rewards_preview(scheme) if scheme else {} my_claim = None my_claimed = None if yonghuid and settlement: pending = RankRewardClaim.query.filter( settlement_id=settlement.id, yonghuid=yonghuid, status=RankRewardClaim.STATUS_PENDING, ).first() if pending: if pending.expire_at < timezone.now(): pending.status = RankRewardClaim.STATUS_EXPIRED pending.save(update_fields=['status', 'UpdateTime']) else: my_claim = { 'claim_id': pending.id, 'amount': float(pending.reward_amount), 'mingci': pending.mingci, 'status': 'pending', 'status_text': '待领取', 'expire_at': pending.expire_at.isoformat(), 'club_id': club_id, } claimed = RankRewardClaim.query.filter( settlement_id=settlement.id, yonghuid=yonghuid, status=RankRewardClaim.STATUS_CLAIMED, ).first() if claimed: my_claimed = { 'claim_id': claimed.id, 'amount': float(claimed.reward_amount), 'mingci': claimed.mingci, 'status': 'claimed', 'status_text': '已领取', 'claimed_at': claimed.claimed_at.isoformat() if claimed.claimed_at else '', } return { 'club_id': club_id, 'shenfen': shenfen, 'riqi': riqi, 'scheme_enabled': bool(scheme and scheme.enabled), 'period_status': 'open' if riqi in OPEN_RIQI else 'closed', 'period_label': riqi, 'sort_field': (scheme.sort_field if scheme and scheme.enabled else display_sort), 'sort_label': ( SORT_FIELD_LABELS.get(scheme.sort_field, display_label) if scheme and scheme.enabled else display_label ), 'title': scheme.title if scheme else '', 'description': scheme.description if scheme else '', 'tiers': tiers, 'tier_rewards': tier_rewards, 'rank_rewards': rank_rewards, 'my_claim': my_claim, 'my_claimed': my_claimed, } def list_my_pending_claims(yonghuid: str, club_id: str = '') -> list: qs = RankRewardClaim.query.filter( yonghuid=yonghuid, status=RankRewardClaim.STATUS_PENDING, ).order_by('-CreateTime') if club_id: qs = qs.filter(club_id=club_id) now = timezone.now() result = [] for c in qs[:20]: if c.expire_at < now: c.status = RankRewardClaim.STATUS_EXPIRED c.save(update_fields=['status', 'UpdateTime']) continue result.append({ 'claim_id': c.id, 'club_id': c.club_id, 'shenfen': c.shenfen, 'amount': float(c.reward_amount), 'mingci': c.mingci, 'status': 'pending', 'status_text': '待领取', 'expire_at': c.expire_at.isoformat(), }) return result def claim_reward(claim_id: int, user) -> dict: uid = user.UserUID user_club = get_user_club_id(user) or getattr(user, 'ClubID', None) or '' with transaction.atomic(): claim = RankRewardClaim.objects.select_for_update().filter(id=claim_id).first() if not claim: return {'ok': False, 'msg': '奖励记录不存在'} if claim.yonghuid != uid: return {'ok': False, 'msg': '无权领取该奖励'} if claim.club_id != user_club: return {'ok': False, 'msg': '俱乐部不匹配,无法领取'} if claim.status == RankRewardClaim.STATUS_CLAIMED: return {'ok': True, 'msg': '已领取', 'amount': float(claim.reward_amount), 'duplicate': True} if claim.status != RankRewardClaim.STATUS_PENDING: return {'ok': False, 'msg': '奖励不可领取'} if claim.expire_at < timezone.now(): claim.status = RankRewardClaim.STATUS_EXPIRED claim.save(update_fields=['status', 'UpdateTime']) return {'ok': False, 'msg': '奖励已过期'} if not user_has_active_shenfen(uid, claim.shenfen): return {'ok': False, 'msg': '身份状态异常,无法领取'} amount = claim.reward_amount _credit_balance(claim.shenfen, uid, amount) claim.status = RankRewardClaim.STATUS_CLAIMED claim.claimed_at = timezone.now() claim.save(update_fields=['status', 'claimed_at', 'UpdateTime']) return { 'ok': True, 'msg': '领取成功', 'amount': float(amount), 'shenfen': claim.shenfen, 'balance_field': claim.balance_field, } def _credit_balance(shenfen: str, yonghuid: str, amount: Decimal): amount = Decimal(str(amount)) if shenfen == 'dashou': p = UserDashou.objects.select_for_update().get(user__UserUID=yonghuid) UserDashou.objects.filter(pk=p.pk).update(yue=F('yue') + amount) elif shenfen == 'guanshi': p = UserGuanshi.objects.select_for_update().get(user__UserUID=yonghuid) UserGuanshi.objects.filter(pk=p.pk).update(yue=F('yue') + amount) elif shenfen == 'zuzhang': p = UserZuzhang.objects.select_for_update().get(user__UserUID=yonghuid) UserZuzhang.objects.filter(pk=p.pk).update(ketixian_jine=F('ketixian_jine') + amount) elif shenfen == 'shangjia': p = UserShangjia.objects.select_for_update().get(user__UserUID=yonghuid) UserShangjia.objects.filter(pk=p.pk).update(yue=F('yue') + amount) else: raise ValueError(f'unsupported shenfen: {shenfen}')