"""Read-only notification-list polling. Checkpoints advance only after DB ACK.""" from collections import Counter, deque from datetime import datetime, timezone import json import logging import math import threading import time import uuid from urllib.parse import urlencode from .douyin import DouyinError, _notice_id, normalize_notice LOG = logging.getLogger('creatorhub.gateway.notices') GROUPS = ('700', '960', '961') HISTORY_RECONCILE_SECONDS = 300 def boundary_time(value): try: result = datetime.fromisoformat(value.replace('Z', '+00:00')) if result.tzinfo is None or result.timestamp() <= 0: raise ValueError('timezone or timestamp') return result except (AttributeError, TypeError, ValueError, OverflowError) as exc: raise DouyinError('notification boundary is invalid') from exc def validate_checkpoint(boundary_at, checkpoints): if boundary_at is not None: boundary_time(boundary_at) if not isinstance(checkpoints, dict) or (checkpoints and boundary_at is None): raise DouyinError('notification checkpoints require a persisted boundary') for group, nid in checkpoints.items(): if group not in GROUPS or not isinstance(nid, str) or not _notice_id(nid): raise DouyinError('notification checkpoint is invalid') class NoticePollingSession: def __init__(self, browser, alias, uid, boundary_at=None, checkpoints=None): checkpoints = {} if checkpoints is None else checkpoints validate_checkpoint(boundary_at, checkpoints) self.browser, self.alias, self.uid = browser, alias, uid identity = browser.identity(alias, uid) self.boundary_at = boundary_at or identity.get('platform_now') self.boundary = boundary_time(self.boundary_at) self.initial = boundary_at is None self.checkpoints = dict(checkpoints) self.queue = deque() self.lock = threading.RLock() self.stopped = False self.last_full_scan_at = None def _page(self, group, cursor): url = 'https://www.douyin.com/aweme/v1/web/notice/?' + urlencode({ 'device_platform': 'webapp', 'aid': 6383, 'channel': 'channel_pc_web', 'is_new_notice': 1, 'is_mark_read': 0, 'notice_group': group, 'count': 50, 'min_time': 0, 'max_time': cursor, }) response = self.browser.get(self.alias, url) try: # Raw text, never JSON.parse in the browser: IDs exceed JS safe integers. body = json.loads(response.body) except (TypeError, ValueError) as exc: raise DouyinError(f'notification group {group} response is invalid JSON') from exc if response.status != 200 or not isinstance(body, dict) or body.get('status_code') != 0: raise DouyinError(f'notification group {group} request failed HTTP {response.status}: {body.get("status_code") if isinstance(body, dict) else "invalid body"}') if not isinstance(body.get('notice_list_v2'), list) or type(body.get('has_more')) is not int or body['has_more'] not in (0, 1): raise DouyinError(f'notification group {group} list or pagination is invalid') return body def _scan(self): self.browser.identity(self.alias, self.uid) received_at = datetime.now(timezone.utc).isoformat() next_checkpoints = dict(self.checkpoints) notices = {} gaps, ignored = [], Counter() pages = 0 full_scan = self.last_full_scan_at is None or time.monotonic() - self.last_full_scan_at >= HISTORY_RECONCILE_SECONDS for group in GROUPS: cursor, found, head = 0, False, None while True: body = self._page(group, cursor) pages += 1 for raw in body['notice_list_v2']: if not isinstance(raw, dict) or _notice_id(raw.get('user_id')) != self.uid: raise DouyinError(f'notification group {group} identity changed') nid = _notice_id(raw.get('nid_str') or raw.get('nid')) if not nid: raise DouyinError('notification ID is invalid') if head is None: head = nid if nid == self.checkpoints.get(group): found = True # Private messages belong exclusively to the IM inbox. if not any(raw.get(k) for k in ('comment', 'follow', 'digg', 'share')): ignored[str(raw.get('type') or next((k for k in ('dm', 'favorite', 'collect') if raw.get(k)), 'unknown'))] += 1 continue timestamp = raw.get('create_time') if not isinstance(timestamp, (int, float)) or isinstance(timestamp, bool) or not math.isfinite(timestamp) or timestamp <= 0: raise DouyinError('notification timestamp is invalid') normalized = normalize_notice(raw, received_at) # The start time labels history; it never excludes messages. notices[(normalized['event_type'], nid)] = (normalized, timestamp <= self.boundary.timestamp()) # Full scans recover older holes, including records behind the head. # Incremental scans still process the entire checkpoint page. if not body['has_more'] or (found and not full_scan): break following = body.get('max_time') if type(following) is not int or following <= 0 or (cursor and following >= cursor): raise DouyinError(f'notification group {group} pagination cursor did not advance') cursor = following if self.checkpoints.get(group) and not found: gaps.append(group) if head: next_checkpoints[group] = head # No partial scans enter the queue. A failed group cannot advance another. result = [] if self.initial: result.append({'kind': 'baseline', 'boundary_at': self.boundary_at}) result.extend({'kind': 'notice', 'notice': n, 'baseline': baseline} for n, baseline in notices.values()) reason = ('通知列表未找到已保存的边界,历史完整性无法确认;分组:' + ','.join(gaps)) if gaps else '' checked_at = datetime.now(timezone.utc).isoformat() result.append({'kind': 'checkpoint', 'checkpoints': next_checkpoints, 'boundary_at': self.boundary_at, 'checked_at': checked_at, 'reason': reason, 'ignored_types': dict(ignored), 'history_reconciled': full_scan}) for item in result: item['delivery_id'] = uuid.uuid4().hex LOG.info('notification scan alias=%s uid=%s mode=%s pages=%s events=%s ignored=%s gap_groups=%s', self.alias, self.uid, 'history_reconcile' if full_scan else 'incremental', pages, len(notices), dict(ignored), gaps) self.queue.extend(result) def poll(self, limit, wait_seconds): # Pacing is owned by the control plane, outside the browser alias lock. with self.lock: if self.stopped: raise DouyinError('notification polling stopped') if not self.queue: self._scan() return list(self.queue)[:limit] def pending(self): with self.lock: return list(self.queue) def ack(self, delivery_ids): ids = set(delivery_ids) with self.lock: for item in self.queue: if item['delivery_id'] in ids and item['kind'] == 'checkpoint': if any(d['delivery_id'] not in ids for d in self.queue): raise DouyinError('notification checkpoint acknowledged before its notices') self.checkpoints = dict(item['checkpoints']) self.initial = False if item['history_reconciled']: self.last_full_scan_at = time.monotonic() self.queue = deque(item for item in self.queue if item['delivery_id'] not in ids) def stop(self): with self.lock: self.stopped = True self.queue.clear() class SubscriptionManager: def __init__(self, browser): self.browser = browser self._lock = threading.RLock() self._items = {} def start(self, alias, uid, boundary_at=None, checkpoints=None): with self._lock: previous = self._items.get(alias) if previous: previous.stop() self._items.pop(alias) browser = self.browser._notification_browser(alias) item = NoticePollingSession(browser, alias, uid, boundary_at, checkpoints) self._items[alias] = item return {'connected': True, 'alias': alias, 'uid': uid, 'mode': 'notice_polling'} def _get(self, alias): with self._lock: item = self._items.get(alias) if item is None: raise DouyinError('notification polling is not running') return item def poll(self, alias, limit, wait_seconds): return self._get(alias).poll(limit, wait_seconds) def ack(self, alias, delivery_ids): self._get(alias).ack(delivery_ids) def stop(self, alias): with self._lock: item = self._items.get(alias) if item: item.stop() self._items.pop(alias) def close(self): with self._lock: aliases = list(self._items) for alias in aliases: self.stop(alias)