"""统一直播统计 SQLite 存储层。 所有写入经由单个 asyncio.Queue 消费协程串行执行。业务侧 record_* 方法均为 best-effort:队列未启动、已关闭、数据不合法或数据库暂时不可用时不会向主业务抛出异常。 """ from __future__ import annotations import asyncio import ctypes import hashlib import json import logging import os import sqlite3 import time from collections.abc import Mapping, Sequence from datetime import UTC, datetime, timedelta, timezone from pathlib import Path from typing import Any from urllib.parse import urlsplit, urlunsplit try: from .core.runtime_paths import APP_ROOT except ImportError: # 允许作为独立模块加载 APP_ROOT = Path(__file__).resolve().parents[1] LOGGER = logging.getLogger(__name__) BEIJING_TZ = timezone(timedelta(hours=8), name="Asia/Shanghai") DEFAULT_DATABASE_PATH = APP_ROOT / "data" / "statistics.sqlite3" _SCHEMA_VERSION = 2 _STOP = object() _FLUSH = object() _LIFECYCLE_KEYS = { "login_sessions": ("login_session_id", "started_at_utc", "business_date"), "group_runs": ("group_run_id", "started_at_utc", "business_date"), "song_requests": ("request_id", "requested_at_utc", "business_date"), "playback_sessions": ("playback_id", "started_at_utc", "business_date"), "broadcast_requests": ("request_id", "requested_at_utc", "business_date"), "tts_requests": ("request_id", "requested_at_utc", "business_date"), "bilibili_connections": ("connection_id", "connected_at_utc", "business_date"), } _SENSITIVE_PARTS = ( "cookie", "token", "authorization", "api_key", "apikey", "access_key", "secret", "password", "passwd", "music_u", "session_key", "csrf", "buvid", ) _AUDIO_URL_PARTS = ("audio_url", "temp_url", "temporary_url", "voice_url", "tts_url") def utc_now() -> str: """返回可排序的 UTC ISO-8601 时间。""" return datetime.now(UTC).isoformat(timespec="milliseconds").replace("+00:00", "Z") def business_date(value: datetime | str | None = None) -> str: """按北京时间计算业务日;无时区输入按 UTC 处理。""" if isinstance(value, str): try: current = datetime.fromisoformat(value.replace("Z", "+00:00")) except ValueError: current = datetime.now(UTC) else: current = value or datetime.now(UTC) if current.tzinfo is None: current = current.replace(tzinfo=UTC) return (current.astimezone(BEIJING_TZ) - timedelta(hours=4)).date().isoformat() def _safe_url(value: str) -> str: """仅保留 URL 的来源与路径,不保存查询串、片段和用户信息。""" try: parsed = urlsplit(value) except ValueError: return "[REDACTED_URL]" if parsed.scheme not in {"http", "https"} or not parsed.netloc: return "[REDACTED_URL]" hostname = parsed.hostname or "" port = f":{parsed.port}" if parsed.port else "" return urlunsplit((parsed.scheme, f"{hostname}{port}", parsed.path, "", "")) def redact_payload(value: Any, key: str = "") -> Any: """递归复制并脱敏可序列化 payload。""" lowered = key.casefold() if any(part in lowered for part in _SENSITIVE_PARTS): return "[REDACTED]" if value is None or isinstance(value, (bool, int, float)): return value if isinstance(value, str): if any(part in lowered for part in _AUDIO_URL_PARTS): return "[REDACTED_AUDIO_URL]" if lowered.endswith("url") and ("?" in value or "#" in value): return _safe_url(value) return value if isinstance(value, bytes): return f"[BYTES:{len(value)}]" if isinstance(value, Mapping): return {str(item_key): redact_payload(item_value, str(item_key)) for item_key, item_value in value.items()} if isinstance(value, Sequence): return [redact_payload(item, key) for item in value] return str(value) def _json(value: Any) -> str: return json.dumps(redact_payload(value), ensure_ascii=False, separators=(",", ":"), default=str) def _utc_iso(value: datetime | str) -> str: if isinstance(value, str): parsed = datetime.fromisoformat(value.replace("Z", "+00:00")) else: parsed = value if parsed.tzinfo is None: parsed = parsed.replace(tzinfo=UTC) return parsed.astimezone(UTC).isoformat(timespec="milliseconds").replace("+00:00", "Z") def _clean_fields(fields: Mapping[str, Any]) -> dict[str, Any]: clean = redact_payload(dict(fields)) assert isinstance(clean, dict) for name, value in tuple(clean.items()): if name.endswith("_at_utc") and value is not None: clean[name] = _utc_iso(value) elif isinstance(value, (dict, list, tuple)): clean[name] = _json(value) elif isinstance(value, bool): clean[name] = int(value) return clean _SCHEMA = """ CREATE TABLE IF NOT EXISTS schema_migrations ( version INTEGER PRIMARY KEY, applied_at_utc TEXT NOT NULL ); CREATE TABLE IF NOT EXISTS process_sessions ( id INTEGER PRIMARY KEY AUTOINCREMENT, session_id TEXT NOT NULL UNIQUE, started_at_utc TEXT NOT NULL, ended_at_utc TEXT, business_date TEXT NOT NULL, pid INTEGER, host TEXT, app_version TEXT, exit_reason TEXT, metadata_json TEXT NOT NULL DEFAULT '{}' ); CREATE TABLE IF NOT EXISTS users ( id INTEGER PRIMARY KEY AUTOINCREMENT, platform TEXT NOT NULL, platform_user_id TEXT NOT NULL, display_name TEXT, avatar_url TEXT, user_level INTEGER, is_admin INTEGER NOT NULL DEFAULT 0, first_seen_at_utc TEXT NOT NULL, last_seen_at_utc TEXT NOT NULL, snapshot_json TEXT NOT NULL DEFAULT '{}', points INTEGER NOT NULL DEFAULT 0, UNIQUE(platform, platform_user_id) ); CREATE TABLE IF NOT EXISTS signin_events ( id INTEGER PRIMARY KEY AUTOINCREMENT, event_id TEXT UNIQUE, user_id INTEGER REFERENCES users(id), platform TEXT, platform_user_id TEXT, signed_at_utc TEXT NOT NULL, business_date TEXT NOT NULL, points_awarded INTEGER NOT NULL DEFAULT 0, streak_days INTEGER, status TEXT, payload_json TEXT NOT NULL DEFAULT '{}' ); CREATE TABLE IF NOT EXISTS point_transactions ( id INTEGER PRIMARY KEY AUTOINCREMENT, transaction_id TEXT NOT NULL UNIQUE, user_id INTEGER REFERENCES users(id), platform TEXT, platform_user_id TEXT, occurred_at_utc TEXT NOT NULL, business_date TEXT NOT NULL, amount INTEGER NOT NULL, balance_after INTEGER, reason TEXT, reference_type TEXT, reference_id TEXT, payload_json TEXT NOT NULL DEFAULT '{}' ); CREATE TABLE IF NOT EXISTS queue_events ( id INTEGER PRIMARY KEY AUTOINCREMENT, event_id TEXT UNIQUE, occurred_at_utc TEXT NOT NULL, business_date TEXT NOT NULL, queue_name TEXT NOT NULL, action TEXT NOT NULL, item_id TEXT, user_id INTEGER REFERENCES users(id), position INTEGER, queue_size INTEGER, wait_ms INTEGER, payload_json TEXT NOT NULL DEFAULT '{}' ); CREATE TABLE IF NOT EXISTS login_sessions ( id INTEGER PRIMARY KEY AUTOINCREMENT, login_session_id TEXT NOT NULL UNIQUE, user_id INTEGER REFERENCES users(id), platform TEXT, platform_user_id TEXT, started_at_utc TEXT NOT NULL, ended_at_utc TEXT, business_date TEXT NOT NULL, status TEXT, client_kind TEXT, metadata_json TEXT NOT NULL DEFAULT '{}' ); CREATE TABLE IF NOT EXISTS group_runs ( id INTEGER PRIMARY KEY AUTOINCREMENT, group_run_id TEXT NOT NULL UNIQUE, group_name TEXT, started_at_utc TEXT NOT NULL, ended_at_utc TEXT, business_date TEXT NOT NULL, status TEXT, requested_count INTEGER, completed_count INTEGER, failed_count INTEGER, payload_json TEXT NOT NULL DEFAULT '{}' ); CREATE TABLE IF NOT EXISTS song_requests ( id INTEGER PRIMARY KEY AUTOINCREMENT, request_id TEXT NOT NULL UNIQUE, user_id INTEGER REFERENCES users(id), platform TEXT, platform_user_id TEXT, requested_at_utc TEXT NOT NULL, business_date TEXT NOT NULL, song_id TEXT, song_name TEXT, artist TEXT, source TEXT, status TEXT, points_cost INTEGER NOT NULL DEFAULT 0, queue_position INTEGER, payload_json TEXT NOT NULL DEFAULT '{}' ); CREATE TABLE IF NOT EXISTS playback_sessions ( id INTEGER PRIMARY KEY AUTOINCREMENT, playback_id TEXT NOT NULL UNIQUE, request_id TEXT REFERENCES song_requests(request_id), song_id TEXT, song_name TEXT, started_at_utc TEXT NOT NULL, ended_at_utc TEXT, business_date TEXT NOT NULL, status TEXT, duration_ms INTEGER, played_ms INTEGER, stop_reason TEXT, payload_json TEXT NOT NULL DEFAULT '{}' ); CREATE TABLE IF NOT EXISTS gift_events ( id INTEGER PRIMARY KEY AUTOINCREMENT, event_id TEXT NOT NULL UNIQUE, user_id INTEGER REFERENCES users(id), platform TEXT, platform_user_id TEXT, occurred_at_utc TEXT NOT NULL, business_date TEXT NOT NULL, gift_id TEXT, gift_name TEXT, quantity INTEGER NOT NULL DEFAULT 1, unit_value REAL NOT NULL DEFAULT 0, total_value REAL NOT NULL DEFAULT 0, currency TEXT, combo_id TEXT, payload_json TEXT NOT NULL DEFAULT '{}' ); CREATE TABLE IF NOT EXISTS gift_aggregates ( id INTEGER PRIMARY KEY AUTOINCREMENT, business_date TEXT NOT NULL, platform TEXT NOT NULL DEFAULT '', platform_user_id TEXT NOT NULL DEFAULT '', gift_id TEXT NOT NULL DEFAULT '', gift_name TEXT, quantity INTEGER NOT NULL DEFAULT 0, total_value REAL NOT NULL DEFAULT 0, updated_at_utc TEXT NOT NULL, UNIQUE(business_date, platform, platform_user_id, gift_id) ); CREATE TABLE IF NOT EXISTS broadcast_requests ( id INTEGER PRIMARY KEY AUTOINCREMENT, request_id TEXT NOT NULL UNIQUE, user_id INTEGER REFERENCES users(id), requested_at_utc TEXT NOT NULL, completed_at_utc TEXT, business_date TEXT NOT NULL, channel TEXT, content_length INTEGER, status TEXT, result_code TEXT, payload_json TEXT NOT NULL DEFAULT '{}' ); CREATE TABLE IF NOT EXISTS tts_requests ( id INTEGER PRIMARY KEY AUTOINCREMENT, request_id TEXT NOT NULL UNIQUE, user_id INTEGER REFERENCES users(id), requested_at_utc TEXT NOT NULL, completed_at_utc TEXT, business_date TEXT NOT NULL, voice TEXT, text_length INTEGER, status TEXT, duration_ms INTEGER, result_code TEXT, payload_json TEXT NOT NULL DEFAULT '{}' ); CREATE TABLE IF NOT EXISTS bilibili_connections ( id INTEGER PRIMARY KEY AUTOINCREMENT, connection_id TEXT NOT NULL UNIQUE, room_id TEXT, connected_at_utc TEXT NOT NULL, disconnected_at_utc TEXT, business_date TEXT NOT NULL, status TEXT, reconnect_count INTEGER NOT NULL DEFAULT 0, disconnect_reason TEXT, payload_json TEXT NOT NULL DEFAULT '{}' ); CREATE TABLE IF NOT EXISTS danmu_events ( id INTEGER PRIMARY KEY AUTOINCREMENT, event_id TEXT NOT NULL UNIQUE, user_id INTEGER REFERENCES users(id), platform TEXT NOT NULL DEFAULT 'bilibili', platform_user_id TEXT, room_id TEXT, occurred_at_utc TEXT NOT NULL, business_date TEXT NOT NULL, message_type TEXT NOT NULL DEFAULT 'danmu', content_length INTEGER, command TEXT, handled INTEGER NOT NULL DEFAULT 0, payload_json TEXT NOT NULL DEFAULT '{}' ); CREATE TABLE IF NOT EXISTS service_state_events ( id INTEGER PRIMARY KEY AUTOINCREMENT, event_id TEXT UNIQUE, service_name TEXT NOT NULL, occurred_at_utc TEXT NOT NULL, business_date TEXT NOT NULL, old_state TEXT, new_state TEXT NOT NULL, reason TEXT, payload_json TEXT NOT NULL DEFAULT '{}' ); CREATE TABLE IF NOT EXISTS admin_audit_events ( id INTEGER PRIMARY KEY AUTOINCREMENT, event_id TEXT NOT NULL UNIQUE, occurred_at_utc TEXT NOT NULL, business_date TEXT NOT NULL, actor TEXT, action TEXT NOT NULL, target_type TEXT, target_id TEXT, success INTEGER NOT NULL DEFAULT 1, remote_address_hash TEXT, payload_json TEXT NOT NULL DEFAULT '{}' ); CREATE TABLE IF NOT EXISTS import_sources ( id INTEGER PRIMARY KEY AUTOINCREMENT, source_key TEXT NOT NULL UNIQUE, source_type TEXT NOT NULL, source_fingerprint TEXT, first_imported_at_utc TEXT NOT NULL, last_imported_at_utc TEXT NOT NULL, status TEXT, metadata_json TEXT NOT NULL DEFAULT '{}' ); CREATE TABLE IF NOT EXISTS import_checkpoints ( id INTEGER PRIMARY KEY AUTOINCREMENT, source_id INTEGER NOT NULL REFERENCES import_sources(id) ON DELETE CASCADE, checkpoint_key TEXT NOT NULL, checkpoint_value TEXT, updated_at_utc TEXT NOT NULL, payload_json TEXT NOT NULL DEFAULT '{}', UNIQUE(source_id, checkpoint_key) ); CREATE TABLE IF NOT EXISTS events ( id INTEGER PRIMARY KEY AUTOINCREMENT, event_id TEXT NOT NULL UNIQUE, event_type TEXT NOT NULL, category TEXT, occurred_at_utc TEXT NOT NULL, business_date TEXT NOT NULL, process_session_id TEXT, user_id INTEGER REFERENCES users(id), platform TEXT, platform_user_id TEXT, correlation_id TEXT, payload_json TEXT NOT NULL DEFAULT '{}' ); CREATE INDEX IF NOT EXISTS idx_events_date_type ON events(business_date, event_type); CREATE INDEX IF NOT EXISTS idx_events_user_time ON events(platform, platform_user_id, occurred_at_utc); CREATE INDEX IF NOT EXISTS idx_signin_date ON signin_events(business_date); CREATE INDEX IF NOT EXISTS idx_points_date ON point_transactions(business_date); CREATE INDEX IF NOT EXISTS idx_song_date_status ON song_requests(business_date, status); CREATE INDEX IF NOT EXISTS idx_playback_date_status ON playback_sessions(business_date, status); CREATE INDEX IF NOT EXISTS idx_gifts_date ON gift_events(business_date); CREATE INDEX IF NOT EXISTS idx_danmu_date ON danmu_events(business_date); CREATE INDEX IF NOT EXISTS idx_tts_date_status ON tts_requests(business_date, status); CREATE INDEX IF NOT EXISTS idx_queue_date_action ON queue_events(business_date, action); CREATE INDEX IF NOT EXISTS idx_group_date_status ON group_runs(business_date, status); CREATE INDEX IF NOT EXISTS idx_bilibili_date_status ON bilibili_connections(business_date, status); CREATE INDEX IF NOT EXISTS idx_broadcast_date_status ON broadcast_requests(business_date, status); CREATE INDEX IF NOT EXISTS idx_service_date_state ON service_state_events(business_date, new_state); CREATE INDEX IF NOT EXISTS idx_admin_audit_date_action ON admin_audit_events(business_date, action); CREATE TABLE IF NOT EXISTS live_sessions ( id INTEGER PRIMARY KEY AUTOINCREMENT, session_id TEXT NOT NULL UNIQUE, kind TEXT NOT NULL, started_at_utc TEXT NOT NULL, ended_at_utc TEXT, status TEXT, duration_ms INTEGER, source TEXT, payload_json TEXT NOT NULL DEFAULT '{}', updated_at_utc TEXT NOT NULL ); CREATE TABLE IF NOT EXISTS live_session_daily_durations ( session_id TEXT NOT NULL REFERENCES live_sessions(session_id) ON DELETE CASCADE, business_date TEXT NOT NULL, duration_ms INTEGER NOT NULL, PRIMARY KEY(session_id, business_date) ); CREATE INDEX IF NOT EXISTS idx_live_sessions_kind_time ON live_sessions(kind, started_at_utc); CREATE INDEX IF NOT EXISTS idx_live_daily_date ON live_session_daily_durations(business_date); """ _DAILY_STATISTICS_VIEW = """ CREATE VIEW daily_live_statistics AS WITH dates AS ( SELECT business_date FROM live_session_daily_durations UNION SELECT business_date FROM events UNION SELECT business_date FROM process_sessions UNION SELECT business_date FROM signin_events UNION SELECT business_date FROM point_transactions UNION SELECT business_date FROM queue_events UNION SELECT business_date FROM group_runs UNION SELECT business_date FROM song_requests UNION SELECT business_date FROM playback_sessions UNION SELECT business_date FROM gift_events UNION SELECT business_date FROM danmu_events UNION SELECT business_date FROM tts_requests UNION SELECT business_date FROM broadcast_requests UNION SELECT business_date FROM bilibili_connections UNION SELECT business_date FROM service_state_events UNION SELECT business_date FROM admin_audit_events ) SELECT d.business_date, COALESCE((SELECT SUM(x.duration_ms) FROM live_session_daily_durations x JOIN live_sessions s ON s.session_id=x.session_id WHERE x.business_date=d.business_date AND s.kind IN ('live','stream','streaming')),0) AS live_duration_ms, COALESCE((SELECT SUM(x.duration_ms) FROM live_session_daily_durations x JOIN live_sessions s ON s.session_id=x.session_id WHERE x.business_date=d.business_date AND s.kind='process'),0) AS process_duration_ms, COALESCE((SELECT SUM(x.duration_ms) FROM live_session_daily_durations x JOIN live_sessions s ON s.session_id=x.session_id WHERE x.business_date=d.business_date AND s.kind='bilibili_connection'),0) AS bilibili_connection_duration_ms, COALESCE((SELECT SUM(reconnect_count) FROM bilibili_connections x WHERE x.business_date=d.business_date),0) AS bilibili_reconnect_count, (SELECT COUNT(*) FROM danmu_events x WHERE x.business_date=d.business_date AND x.message_type='danmu') AS danmu_count, (SELECT COUNT(*) FROM danmu_events x WHERE x.business_date=d.business_date AND (x.command IS NOT NULL AND x.command<>'') ) AS command_count, (SELECT COUNT(DISTINCT platform || ':' || COALESCE(platform_user_id,'')) FROM danmu_events x WHERE x.business_date=d.business_date AND COALESCE(platform_user_id,'')<>'') AS active_users, (SELECT COUNT(*) FROM signin_events x WHERE x.business_date=d.business_date) AS signin_count, (SELECT COUNT(DISTINCT platform || ':' || COALESCE(platform_user_id,'')) FROM signin_events x WHERE x.business_date=d.business_date AND COALESCE(platform_user_id,'')<>'') AS unique_signin_users, COALESCE((SELECT SUM(amount) FROM point_transactions x WHERE x.business_date=d.business_date),0) AS point_net_change, (SELECT COUNT(*) FROM queue_events x WHERE x.business_date=d.business_date AND x.action IN ('join','enqueue','entered')) AS queue_join_count, (SELECT COUNT(DISTINCT COALESCE(CAST(user_id AS TEXT), json_extract(payload_json,'$.platform_user_id'), item_id)) FROM queue_events x WHERE x.business_date=d.business_date AND x.action IN ('join','enqueue','entered')) AS unique_queue_users, (SELECT COUNT(*) FROM queue_events x WHERE x.business_date=d.business_date AND x.action IN ('leave','dequeue','removed')) AS leave_count, (SELECT COUNT(*) FROM queue_events x WHERE x.business_date=d.business_date AND x.action IN ('timeout','expired')) AS timeout_count, COALESCE((SELECT MAX(queue_size) FROM queue_events x WHERE x.business_date=d.business_date),0) AS max_queue_size, (SELECT AVG(wait_ms) FROM queue_events x WHERE x.business_date=d.business_date AND x.action='promoted' AND wait_ms IS NOT NULL) AS avg_wait_ms, (SELECT COUNT(*) FROM queue_events x WHERE x.business_date=d.business_date AND x.action='promoted' AND wait_ms IS NOT NULL) AS wait_sample_count, (SELECT COUNT(*) FROM group_runs x WHERE x.business_date=d.business_date) AS group_run_count, (SELECT COUNT(*) FROM group_runs x WHERE x.business_date=d.business_date AND x.status IN ('success','completed','finished')) AS group_run_success_count, (SELECT COUNT(*) FROM group_runs x WHERE x.business_date=d.business_date AND x.status IN ('failed','error')) AS group_run_failed_count, COALESCE((SELECT SUM(CASE WHEN ended_at_utc IS NOT NULL THEN MAX(0, CAST((julianday(ended_at_utc)-julianday(started_at_utc))*86400000 AS INTEGER)) ELSE 0 END) FROM group_runs x WHERE x.business_date=d.business_date),0) AS group_run_total_duration_ms, (SELECT COUNT(*) FROM song_requests x WHERE x.business_date=d.business_date) AS song_request_count, (SELECT COUNT(DISTINCT platform || ':' || COALESCE(platform_user_id,'')) FROM song_requests x WHERE x.business_date=d.business_date AND COALESCE(platform_user_id,'')<>'') AS unique_song_users, (SELECT COUNT(*) FROM song_requests x WHERE x.business_date=d.business_date AND x.status IN ('played','completed','finished','success')) AS played_count, (SELECT COUNT(*) FROM song_requests x WHERE x.business_date=d.business_date AND x.status IN ('skipped','skip')) AS skipped_count, (SELECT COUNT(*) FROM song_requests x WHERE x.business_date=d.business_date AND x.status IN ('error','failed')) AS error_count, (SELECT COUNT(*) FROM playback_sessions x WHERE x.business_date=d.business_date AND json_extract(x.payload_json,'$.source')='background') AS background_playback_count, (SELECT COUNT(*) FROM gift_events x WHERE x.business_date=d.business_date) AS gift_event_count, COALESCE((SELECT SUM(quantity) FROM gift_events x WHERE x.business_date=d.business_date),0) AS gift_quantity, (SELECT COUNT(DISTINCT platform || ':' || COALESCE(platform_user_id,'')) FROM gift_events x WHERE x.business_date=d.business_date AND COALESCE(platform_user_id,'')<>'') AS unique_gifters, COALESCE((SELECT SUM(total_value) FROM gift_events x WHERE x.business_date=d.business_date AND x.currency='CNY'),0) AS gift_total_value, (SELECT COUNT(*) FROM tts_requests x WHERE x.business_date=d.business_date) AS tts_request_count, (SELECT COUNT(*) FROM tts_requests x WHERE x.business_date=d.business_date AND x.status IN ('success','completed','finished')) AS tts_success_count, (SELECT COUNT(*) FROM tts_requests x WHERE x.business_date=d.business_date AND x.status IN ('failure','failed','error')) AS tts_failure_count, COALESCE((SELECT AVG(duration_ms) FROM tts_requests x WHERE x.business_date=d.business_date AND duration_ms IS NOT NULL),0) AS tts_avg_duration_ms, (SELECT COUNT(*) FROM broadcast_requests x WHERE x.business_date=d.business_date) AS broadcast_request_count, (SELECT COUNT(*) FROM broadcast_requests x WHERE x.business_date=d.business_date AND x.status IN ('success','completed','finished')) AS broadcast_success_count, (SELECT COUNT(*) FROM broadcast_requests x WHERE x.business_date=d.business_date AND x.status IN ('failure','failed','error')) AS broadcast_failure_count, (SELECT COUNT(*) FROM service_state_events x WHERE x.business_date=d.business_date AND x.new_state IN ('error','failed','failure')) AS service_error_count, (SELECT COUNT(*) FROM service_state_events x WHERE x.business_date=d.business_date AND x.new_state IN ('restart','restarted','starting') ) AS service_restart_count, (SELECT COUNT(*) FROM admin_audit_events x WHERE x.business_date=d.business_date) AS admin_action_count, (SELECT COUNT(*) FROM events x WHERE x.business_date=d.business_date) AS generic_event_count FROM dates d; """ class StatsStore: """基于 SQLite 的异步单写统计存储。""" def __init__( self, database_path: str | Path = DEFAULT_DATABASE_PATH, *, queue_size: int = 10_000, batch_size: int = 100, commit_interval: float = 0.5, ) -> None: self.database_path = Path(database_path) self.queue: asyncio.Queue[Any] = asyncio.Queue(maxsize=max(1, queue_size)) self.batch_size = max(1, batch_size) self.commit_interval = max(0.05, commit_interval) self._writer_task: asyncio.Task[None] | None = None self._started = False self._closing = False self._connection: sqlite3.Connection | None = None self._last_error: str | None = None self._written = 0 self._dropped = 0 self._process_session_id: str | None = None async def start(self) -> bool: if self._started and self._writer_task and not self._writer_task.done(): return True if self._closing: return False try: self.database_path.parent.mkdir(parents=True, exist_ok=True) connection = sqlite3.connect(self.database_path, timeout=5.0) connection.row_factory = sqlite3.Row connection.execute("PRAGMA journal_mode=WAL") connection.execute("PRAGMA synchronous=NORMAL") connection.execute("PRAGMA foreign_keys=ON") connection.execute("PRAGMA busy_timeout=5000") self._migrate(connection) gift_values_repaired = self._repair_gift_value_history(connection) if gift_values_repaired: LOGGER.warning("已修复统计库历史礼物金额: %s 条", gift_values_repaired) repaired = self._repair_startup_reconciliation_history(connection) if any(repaired.values()): LOGGER.warning("已修复统计库历史异常生命周期: %s", repaired) recovered = self._reconcile_stale_lifecycles(connection) if any(recovered.values()): LOGGER.warning("已收口统计库遗留生命周期: %s", recovered) self._connection = connection self._started = True self._writer_task = asyncio.create_task(self._writer(), name="stats-store-writer") return True except Exception as exc: self._last_error = f"{type(exc).__name__}: {exc}" LOGGER.exception("无法启动统计存储") try: connection.close() except (UnboundLocalError, sqlite3.Error): pass self._connection = None self._started = False return False async def close(self) -> None: if not self._started: return self._closing = True try: await self.flush() await self.queue.put(_STOP) if self._writer_task: await self._writer_task except Exception as exc: self._last_error = f"{type(exc).__name__}: {exc}" LOGGER.exception("关闭统计存储失败") finally: self._started = False self._writer_task = None if self._connection: try: self._connection.close() except sqlite3.Error: pass self._connection = None self._closing = False async def flush(self) -> bool: writer = self._writer_task if not self._started or not writer or writer.done(): return False loop = asyncio.get_running_loop() completed = loop.create_future() try: await self.queue.put((_FLUSH, completed)) done, _ = await asyncio.wait((completed, writer), return_when=asyncio.FIRST_COMPLETED) if completed in done: return bool(completed.result()) self._last_error = "统计 writer 在 flush 完成前退出" return False except Exception as exc: self._last_error = f"{type(exc).__name__}: {exc}" return False async def execute(self, func, *args) -> Any: """在 writer 持有的唯一连接上执行同步函数并返回其结果。 供其它模块(如兑换码)复用同一 sqlite 连接、串行化写入, 消除多连接写同一文件导致的 "database is locked" 写争用。 """ writer = self._writer_task if not self._started or not writer or writer.done(): raise RuntimeError("统计存储未启动") loop = asyncio.get_running_loop() completed = loop.create_future() try: await self.queue.put(("__execute__", func, args, completed)) except Exception as exc: completed.cancel() raise RuntimeError("统计存储写入队列已满") from exc return await completed def emit_nowait( self, event_type: str, payload: Mapping[str, Any] | None = None, **fields: Any, ) -> bool: """无等待旁路写入;队列满或未启动时返回 False。""" if not self._started or self._closing or not event_type: self._dropped += 1 return False try: operation = self._event_operation(event_type, payload, fields) self.queue.put_nowait(operation) return True except (asyncio.QueueFull, Exception) as exc: self._dropped += 1 self._last_error = f"{type(exc).__name__}: {exc}" return False async def emit( self, event_type: str, payload: Mapping[str, Any] | None = None, **fields: Any, ) -> bool: """等待队列容量后旁路写入;失败不抛出。""" if not self._started or self._closing or not event_type: self._dropped += 1 return False try: await self.queue.put(self._event_operation(event_type, payload, fields)) return True except Exception as exc: self._dropped += 1 self._last_error = f"{type(exc).__name__}: {exc}" return False def _event_operation( self, event_type: str, payload: Mapping[str, Any] | None, fields: Mapping[str, Any], ) -> tuple[str, dict[str, Any]]: now = _utc_iso(fields.get("occurred_at_utc") or utc_now()) event_id = str(fields.get("event_id") or self.make_event_id(event_type, now, fields, payload or {})) return ( "events", { "event_id": event_id, "event_type": event_type, "category": fields.get("category"), "occurred_at_utc": now, "business_date": fields.get("business_date") or business_date(now), "process_session_id": fields.get("process_session_id") or self._process_session_id, "user_id": fields.get("user_id"), "platform": fields.get("platform"), "platform_user_id": fields.get("platform_user_id"), "correlation_id": fields.get("correlation_id"), "payload_json": _json(payload or {}), }, ) def _enqueue_record(self, table: str, fields: Mapping[str, Any]) -> bool: if not self._started or self._closing: self._dropped += 1 return False try: self.queue.put_nowait((table, _clean_fields(fields))) return True except Exception as exc: self._dropped += 1 self._last_error = f"{type(exc).__name__}: {exc}" return False async def _writer(self) -> None: connection = self._connection if connection is None: return pending = 0 while True: try: try: item = await asyncio.wait_for(self.queue.get(), timeout=self.commit_interval) except TimeoutError: if pending: connection.commit() pending = 0 continue if item is _STOP: if pending: connection.commit() self.queue.task_done() return if isinstance(item, tuple) and item and item[0] is _FLUSH: if pending: connection.commit() pending = 0 future = item[1] if not future.done(): future.set_result(True) self.queue.task_done() continue if isinstance(item, tuple) and len(item) == 4 and item[0] == "__execute__": if pending: connection.commit() pending = 0 _, func, args, future = item try: result = func(connection, *args) if not future.done(): future.set_result(result) except Exception as exc: try: connection.rollback() except Exception: pass if not future.done(): future.set_exception(exc) self.queue.task_done() continue table, fields = item if table == "__process_start__": self._insert(connection, "process_sessions", fields) self._upsert_live_session(connection, { "session_id": fields["session_id"], "kind": "process", "started_at_utc": fields["started_at_utc"], "status": "running", "source": "process_sessions", "payload_json": fields.get("metadata_json", "{}"), }) elif table == "__process_end__": connection.execute( "UPDATE process_sessions SET ended_at_utc=?, exit_reason=? WHERE session_id=?", (fields["ended_at_utc"], fields.get("exit_reason"), fields["session_id"]), ) row = connection.execute( "SELECT started_at_utc, metadata_json FROM process_sessions WHERE session_id=?", (fields["session_id"],), ).fetchone() if row: self._upsert_live_session(connection, { "session_id": fields["session_id"], "kind": "process", "started_at_utc": row[0], "ended_at_utc": fields["ended_at_utc"], "status": "ended", "source": "process_sessions", "payload_json": row[1], }) elif table == "__bilibili_connection__": self._upsert_lifecycle(connection, "bilibili_connections", fields) if fields.get("disconnected_at_utc"): row = connection.execute( "SELECT connected_at_utc, disconnected_at_utc, status, payload_json " "FROM bilibili_connections WHERE connection_id=?", (fields["connection_id"],), ).fetchone() if row and row[1]: self._upsert_live_session(connection, { "session_id": fields["connection_id"], "kind": "bilibili_connection", "started_at_utc": row[0], "ended_at_utc": row[1], "status": row[2], "source": "bilibili_connections", "payload_json": row[3], }) elif table == "__user_upsert__": self._upsert_user(connection, fields) elif table == "__import__": self._run_import(connection, fields) pending = 0 elif table == "__live_session_upsert__": self._upsert_live_session(connection, fields) else: if table in _LIFECYCLE_KEYS: inserted = self._upsert_lifecycle(connection, table, fields) else: inserted = self._insert(connection, table, fields) if table == "gift_events" and inserted: self._aggregate_gift(connection, fields) pending += 1 self._written += 1 self.queue.task_done() if pending >= self.batch_size: connection.commit() pending = 0 except Exception as exc: connection.rollback() pending = 0 self._last_error = f"{type(exc).__name__}: {exc}" LOGGER.exception("统计写入失败") try: self.queue.task_done() except ValueError: pass @staticmethod def _insert(connection: sqlite3.Connection, table: str, fields: Mapping[str, Any]) -> bool: allowed = { "process_sessions", "signin_events", "point_transactions", "queue_events", "login_sessions", "group_runs", "song_requests", "playback_sessions", "gift_events", "broadcast_requests", "tts_requests", "bilibili_connections", "danmu_events", "service_state_events", "admin_audit_events", "events", } if table not in allowed: raise ValueError(f"不允许写入表: {table}") columns = tuple(fields) sql = f"INSERT OR IGNORE INTO {table} ({','.join(columns)}) VALUES ({','.join('?' for _ in columns)})" cursor = connection.execute(sql, tuple(fields[column] for column in columns)) return cursor.rowcount > 0 @staticmethod def _upsert_lifecycle(connection: sqlite3.Connection, table: str, fields: Mapping[str, Any]) -> bool: definition = _LIFECYCLE_KEYS.get(table) if definition is None: raise ValueError(f"不允许生命周期更新表: {table}") key_column, start_column, business_column = definition if key_column not in fields: raise ValueError(f"{table} 缺少稳定 ID: {key_column}") columns = tuple(fields) placeholders = ",".join("?" for _ in columns) immutable = {key_column, start_column, business_column} updates = tuple(column for column in columns if column not in immutable) if updates: update_sql = ",".join(f"{column}=excluded.{column}" for column in updates) conflict_sql = f"DO UPDATE SET {update_sql}" else: conflict_sql = "DO NOTHING" sql = ( f"INSERT INTO {table} ({','.join(columns)}) VALUES ({placeholders}) " f"ON CONFLICT({key_column}) {conflict_sql}" ) cursor = connection.execute(sql, tuple(fields[column] for column in columns)) return cursor.rowcount > 0 @staticmethod def _upsert_user(connection: sqlite3.Connection, fields: Mapping[str, Any]) -> None: connection.execute( """INSERT INTO users ( platform, platform_user_id, display_name, avatar_url, user_level, is_admin, first_seen_at_utc, last_seen_at_utc, snapshot_json ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(platform, platform_user_id) DO UPDATE SET display_name=CASE WHEN excluded.last_seen_at_utc >= users.last_seen_at_utc THEN COALESCE(excluded.display_name, users.display_name) ELSE users.display_name END, avatar_url=CASE WHEN excluded.last_seen_at_utc >= users.last_seen_at_utc THEN COALESCE(excluded.avatar_url, users.avatar_url) ELSE users.avatar_url END, user_level=CASE WHEN excluded.last_seen_at_utc >= users.last_seen_at_utc THEN COALESCE(excluded.user_level, users.user_level) ELSE users.user_level END, is_admin=CASE WHEN excluded.last_seen_at_utc >= users.last_seen_at_utc THEN excluded.is_admin ELSE users.is_admin END, first_seen_at_utc=MIN(users.first_seen_at_utc, excluded.first_seen_at_utc), last_seen_at_utc=MAX(users.last_seen_at_utc, excluded.last_seen_at_utc), snapshot_json=CASE WHEN excluded.last_seen_at_utc >= users.last_seen_at_utc THEN excluded.snapshot_json ELSE users.snapshot_json END""", tuple(fields[name] for name in ( "platform", "platform_user_id", "display_name", "avatar_url", "user_level", "is_admin", "first_seen_at_utc", "last_seen_at_utc", "snapshot_json", )), ) @staticmethod def _upsert_live_session(connection: sqlite3.Connection, fields: Mapping[str, Any]) -> None: existing = connection.execute( "SELECT kind, started_at_utc, ended_at_utc, status, duration_ms, source, payload_json FROM live_sessions WHERE session_id=?", (fields["session_id"],), ).fetchone() started = fields.get("started_at_utc") or (existing[1] if existing else None) if not started: raise ValueError("live session 缺少 started_at_utc") ended = fields.get("ended_at_utc") if "ended_at_utc" in fields else (existing[2] if existing else None) kind = fields.get("kind") or (existing[0] if existing else "live") status = fields.get("status") if "status" in fields else (existing[3] if existing else None) source = fields.get("source") if "source" in fields else (existing[5] if existing else None) payload_json = fields.get("payload_json") if "payload_json" in fields else (existing[6] if existing else "{}") duration_ms = fields.get("duration_ms") if duration_ms is None and ended: duration_ms = StatsStore._duration_ms(started, ended) connection.execute( """INSERT INTO live_sessions ( session_id, kind, started_at_utc, ended_at_utc, status, duration_ms, source, payload_json, updated_at_utc ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(session_id) DO UPDATE SET kind=excluded.kind, started_at_utc=excluded.started_at_utc, ended_at_utc=excluded.ended_at_utc, status=excluded.status, duration_ms=excluded.duration_ms, source=excluded.source, payload_json=excluded.payload_json, updated_at_utc=excluded.updated_at_utc""", (fields["session_id"], kind, started, ended, status, duration_ms, source, payload_json, utc_now()), ) connection.execute("DELETE FROM live_session_daily_durations WHERE session_id=?", (fields["session_id"],)) if ended: for day, allocated_ms in StatsStore._split_session_duration(started, ended, duration_ms): connection.execute( "INSERT INTO live_session_daily_durations(session_id, business_date, duration_ms) VALUES (?, ?, ?)", (fields["session_id"], day, allocated_ms), ) @staticmethod def _parse_utc(value: str) -> datetime: parsed = datetime.fromisoformat(str(value).replace("Z", "+00:00")) if parsed.tzinfo is None: parsed = parsed.replace(tzinfo=UTC) return parsed.astimezone(UTC) @staticmethod def _duration_ms(started_at_utc: str, ended_at_utc: str) -> int: return max(0, round((StatsStore._parse_utc(ended_at_utc) - StatsStore._parse_utc(started_at_utc)).total_seconds() * 1000)) @staticmethod def _split_session_duration(started_at_utc: str, ended_at_utc: str, duration_ms: Any = None) -> list[tuple[str, int]]: started = StatsStore._parse_utc(started_at_utc) ended = StatsStore._parse_utc(ended_at_utc) if ended <= started: return [] exact_total = (ended - started).total_seconds() * 1000 target_total = max(0, int(duration_ms)) if duration_ms is not None else round(exact_total) segments: list[tuple[str, float]] = [] cursor = started while cursor < ended: local = cursor.astimezone(BEIJING_TZ) boundary_date = local.date() if local.hour < 4 else local.date() + timedelta(days=1) boundary = datetime.combine(boundary_date, datetime.min.time(), BEIJING_TZ).replace(hour=4).astimezone(UTC) segment_end = min(ended, boundary) segments.append((business_date(cursor), (segment_end - cursor).total_seconds() * 1000)) cursor = segment_end allocated: list[tuple[str, int]] = [] consumed = 0 for index, (day, milliseconds) in enumerate(segments): value = target_total - consumed if index == len(segments) - 1 else round(target_total * milliseconds / exact_total) value = max(0, value) allocated.append((day, value)) consumed += value return allocated @staticmethod def _business_day_end(value: datetime) -> datetime: local = value.astimezone(BEIJING_TZ) boundary_date = local.date() if local.hour < 4 else local.date() + timedelta(days=1) return datetime.combine(boundary_date, datetime.min.time(), BEIJING_TZ).replace(hour=4).astimezone(UTC) @classmethod def _merged_daily_interval_durations( cls, intervals: Sequence[tuple[Any, Any]], ) -> dict[str, int]: daily_intervals: dict[str, list[tuple[datetime, datetime]]] = {} for started_value, ended_value in intervals: if not started_value or not ended_value: continue try: started = cls._parse_utc(str(started_value)) ended = cls._parse_utc(str(ended_value)) except (TypeError, ValueError): continue if ended <= started: continue cursor = started while cursor < ended: segment_end = min(ended, cls._business_day_end(cursor)) daily_intervals.setdefault(business_date(cursor), []).append((cursor, segment_end)) cursor = segment_end durations: dict[str, int] = {} for day, segments in daily_intervals.items(): merged_ms = 0.0 current_start: datetime | None = None current_end: datetime | None = None for segment_start, segment_end in sorted(segments): if current_start is None: current_start, current_end = segment_start, segment_end elif segment_start <= current_end: current_end = max(current_end, segment_end) else: merged_ms += (current_end - current_start).total_seconds() * 1000 current_start, current_end = segment_start, segment_end if current_start is not None and current_end is not None: merged_ms += (current_end - current_start).total_seconds() * 1000 durations[day] = min(86_400_000, max(0, round(merged_ms))) return durations @staticmethod def _aggregate_gift(connection: sqlite3.Connection, fields: Mapping[str, Any]) -> None: connection.execute( """INSERT INTO gift_aggregates ( business_date, platform, platform_user_id, gift_id, gift_name, quantity, total_value, updated_at_utc ) VALUES (?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(business_date, platform, platform_user_id, gift_id) DO UPDATE SET gift_name=COALESCE(excluded.gift_name, gift_aggregates.gift_name), quantity=gift_aggregates.quantity + excluded.quantity, total_value=gift_aggregates.total_value + excluded.total_value, updated_at_utc=excluded.updated_at_utc""", ( fields["business_date"], fields.get("platform") or "", fields.get("platform_user_id") or "", fields.get("gift_id") or "", fields.get("gift_name"), fields.get("quantity", 1), fields.get("total_value", 0), fields.get("occurred_at_utc") or utc_now(), ), ) def _run_import(self, connection: sqlite3.Connection, fields: Mapping[str, Any]) -> None: future = fields["future"] try: existing = connection.execute( "SELECT id, source_fingerprint, status FROM import_sources WHERE source_key=?", (fields["source_key"],), ).fetchone() if existing and existing[1] == fields.get("fingerprint") and existing[2] == "completed": connection.commit() if not future.done(): future.set_result(True) return now = utc_now() connection.execute( """INSERT INTO import_sources ( source_key, source_type, source_fingerprint, first_imported_at_utc, last_imported_at_utc, status, metadata_json ) VALUES (?, ?, ?, ?, ?, 'running', ?) ON CONFLICT(source_key) DO UPDATE SET source_type=excluded.source_type, source_fingerprint=excluded.source_fingerprint, last_imported_at_utc=excluded.last_imported_at_utc, status='running', metadata_json=excluded.metadata_json""", (fields["source_key"], fields["source_type"], fields.get("fingerprint"), now, now, fields["metadata_json"]), ) source_id = connection.execute( "SELECT id FROM import_sources WHERE source_key=?", (fields["source_key"],) ).fetchone()[0] for table, record in fields["records"]: if table == "users": self._upsert_user(connection, record) else: if table in _LIFECYCLE_KEYS: inserted = self._upsert_lifecycle(connection, table, record) else: inserted = self._insert(connection, table, record) if table == "gift_events" and inserted: self._aggregate_gift(connection, record) connection.execute( """INSERT INTO import_checkpoints ( source_id, checkpoint_key, checkpoint_value, updated_at_utc, payload_json ) VALUES (?, ?, ?, ?, '{}') ON CONFLICT(source_id, checkpoint_key) DO UPDATE SET checkpoint_value=excluded.checkpoint_value, updated_at_utc=excluded.updated_at_utc""", (source_id, fields["checkpoint_key"], fields.get("checkpoint_value"), now), ) connection.execute( "UPDATE import_sources SET status='completed', last_imported_at_utc=? WHERE id=?", (now, source_id), ) connection.commit() if not future.done(): future.set_result(True) except Exception as exc: connection.rollback() if not future.done(): future.set_result(False) raise exc @staticmethod def _migrate(connection: sqlite3.Connection) -> None: current_version = int(connection.execute("PRAGMA user_version").fetchone()[0]) if current_version > _SCHEMA_VERSION: raise RuntimeError(f"数据库版本 {current_version} 高于程序支持版本 {_SCHEMA_VERSION}") connection.executescript(_SCHEMA) connection.execute("DROP VIEW IF EXISTS daily_live_statistics") connection.executescript(_DAILY_STATISTICS_VIEW) # 积分权威列迁移:旧库补充 points 列,并从历史 snapshot_json 一次性回填。 user_columns = {row[1] for row in connection.execute("PRAGMA table_info(users)")} if "points" not in user_columns: connection.execute("ALTER TABLE users ADD COLUMN points INTEGER NOT NULL DEFAULT 0") connection.execute( "UPDATE users SET points = CAST(COALESCE(json_extract(snapshot_json, '$.points'), 0) AS INTEGER)" ) for version in range(max(1, current_version + 1), _SCHEMA_VERSION + 1): connection.execute( "INSERT OR IGNORE INTO schema_migrations(version, applied_at_utc) VALUES (?, ?)", (version, utc_now()), ) connection.execute(f"PRAGMA user_version={_SCHEMA_VERSION}") connection.commit() @staticmethod def _pid_is_running(pid: Any) -> bool: try: value = int(pid) except (TypeError, ValueError): return False if value <= 0: return False if os.name == "nt": process = ctypes.windll.kernel32.OpenProcess(0x1000, False, value) if not process: return False try: exit_code = ctypes.c_ulong() if not ctypes.windll.kernel32.GetExitCodeProcess(process, ctypes.byref(exit_code)): return False return exit_code.value == 259 finally: ctypes.windll.kernel32.CloseHandle(process) try: os.kill(value, 0) return True except PermissionError: return True except (OSError, ProcessLookupError): return False @staticmethod def _next_lifecycle_start( connection: sqlite3.Connection, table: str, start_column: str, row_id: int, started_at_utc: str, ) -> str | None: row = connection.execute( f"SELECT {start_column} FROM {table} " f"WHERE {start_column}>? OR ({start_column}=? AND id>?) " f"ORDER BY {start_column}, id LIMIT 1", (started_at_utc, started_at_utc, row_id), ).fetchone() return str(row[0]) if row else None @classmethod def _reconciliation_end( cls, connection: sqlite3.Connection, table: str, start_column: str, row_id: int, started_at_utc: str, now: str, ) -> str: next_started = cls._next_lifecycle_start( connection, table, start_column, row_id, started_at_utc ) if next_started and cls._parse_utc(next_started) < cls._parse_utc(now): return next_started return now @staticmethod def _repair_gift_value_history(connection: sqlite3.Connection) -> int: cursor = connection.execute( """UPDATE gift_events SET total_value=MAX(0, CAST(json_extract(payload_json, '$.raw_total_coin') AS REAL)) / 1000.0, unit_value=CASE WHEN quantity > 0 THEN (MAX(0, CAST(json_extract(payload_json, '$.raw_total_coin') AS REAL)) / 1000.0) / quantity ELSE 0 END, payload_json=json_set( payload_json, '$.value_rule', 'bilibili_gold_coin_1000_to_cny_1_v2' ) WHERE currency='CNY' AND LOWER(COALESCE(json_extract(payload_json, '$.coin_type'), ''))='gold' AND COALESCE(json_extract(payload_json, '$.value_rule'), '')<>'bilibili_gold_coin_1000_to_cny_1_v2' AND ( json_extract(payload_json, '$.value_rule')='gold_battery_10_to_cny_1_v1' OR json_extract(payload_json, '$.value_migration')='bilibili_coin_to_cny_v1' ) AND json_extract(payload_json, '$.raw_total_coin') IS NOT NULL""" ) repaired = max(0, int(cursor.rowcount)) if repaired: connection.execute("DELETE FROM gift_aggregates") connection.execute( """INSERT INTO gift_aggregates ( business_date, platform, platform_user_id, gift_id, gift_name, quantity, total_value, updated_at_utc ) SELECT business_date, COALESCE(platform, ''), COALESCE(platform_user_id, ''), COALESCE(gift_id, ''), MAX(gift_name), SUM(quantity), SUM(total_value), MAX(occurred_at_utc) FROM gift_events GROUP BY business_date, COALESCE(platform, ''), COALESCE(platform_user_id, ''), COALESCE(gift_id, '')""" ) connection.commit() return repaired @classmethod def _repair_startup_reconciliation_history(cls, connection: sqlite3.Connection) -> dict[str, int]: repaired = {"process_sessions": 0, "bilibili_connections": 0, "group_runs": 0} process_rows = connection.execute( "SELECT id, session_id, started_at_utc, ended_at_utc, exit_reason, metadata_json " "FROM process_sessions ORDER BY started_at_utc, id" ).fetchall() for index, row in enumerate(process_rows[:-1]): next_started = str(process_rows[index + 1][2]) if row[4] != "startup_reconciliation" or not row[3]: continue if cls._parse_utc(str(row[3])) <= cls._parse_utc(next_started): continue connection.execute( "UPDATE process_sessions SET ended_at_utc=? WHERE id=?", (next_started, row[0]), ) cls._upsert_live_session(connection, { "session_id": row[1], "kind": "process", "started_at_utc": row[2], "ended_at_utc": next_started, "status": "interrupted", "source": "process_sessions", "payload_json": row[5] or "{}", }) repaired["process_sessions"] += 1 connection_rows = connection.execute( "SELECT id, connection_id, connected_at_utc, disconnected_at_utc, disconnect_reason, payload_json " "FROM bilibili_connections ORDER BY connected_at_utc, id" ).fetchall() for index, row in enumerate(connection_rows[:-1]): next_started = str(connection_rows[index + 1][2]) if row[4] != "startup_reconciliation" or not row[3]: continue if cls._parse_utc(str(row[3])) <= cls._parse_utc(next_started): continue connection.execute( "UPDATE bilibili_connections SET disconnected_at_utc=? WHERE id=?", (next_started, row[0]), ) cls._upsert_live_session(connection, { "session_id": row[1], "kind": "bilibili_connection", "started_at_utc": row[2], "ended_at_utc": next_started, "status": "failed", "source": "bilibili_connections", "payload_json": row[5] or "{}", }) repaired["bilibili_connections"] += 1 group_rows = connection.execute( "SELECT id, started_at_utc, ended_at_utc, payload_json " "FROM group_runs ORDER BY started_at_utc, id" ).fetchall() for index, row in enumerate(group_rows[:-1]): if not row[2]: continue try: payload = json.loads(row[3] or "{}") except (TypeError, ValueError, json.JSONDecodeError): continue next_started = str(group_rows[index + 1][1]) if payload.get("reason") != "startup_reconciliation": continue if cls._parse_utc(str(row[2])) <= cls._parse_utc(next_started): continue connection.execute( "UPDATE group_runs SET ended_at_utc=? WHERE id=?", (next_started, row[0]), ) repaired["group_runs"] += 1 connection.commit() return repaired @classmethod def _reconcile_stale_lifecycles(cls, connection: sqlite3.Connection) -> dict[str, int]: """启动时收口异常退出遗留的进程、B站连接和配置组会话。""" recovered = {"process_sessions": 0, "bilibili_connections": 0, "group_runs": 0} now = utc_now() stale_process_ids: list[str] = [] for row in connection.execute( "SELECT id, session_id, pid, started_at_utc, metadata_json " "FROM process_sessions WHERE ended_at_utc IS NULL ORDER BY started_at_utc, id" ).fetchall(): if cls._pid_is_running(row[2]): continue ended_at = cls._reconciliation_end( connection, "process_sessions", "started_at_utc", row[0], row[3], now ) connection.execute( "UPDATE process_sessions SET ended_at_utc=?, exit_reason=? WHERE session_id=? AND ended_at_utc IS NULL", (ended_at, "startup_reconciliation", row[1]), ) cls._upsert_live_session(connection, { "session_id": row[1], "kind": "process", "started_at_utc": row[3], "ended_at_utc": ended_at, "status": "interrupted", "source": "process_sessions", "payload_json": row[4] or "{}", }) stale_process_ids.append(str(row[1])) recovered["process_sessions"] += 1 live_process_exists = any( cls._pid_is_running(row[0]) for row in connection.execute( "SELECT pid FROM process_sessions WHERE ended_at_utc IS NULL" ).fetchall() ) if live_process_exists: connection.commit() return recovered for row in connection.execute( "SELECT id, connection_id, connected_at_utc, payload_json FROM bilibili_connections " "WHERE disconnected_at_utc IS NULL ORDER BY connected_at_utc, id" ).fetchall(): ended_at = cls._reconciliation_end( connection, "bilibili_connections", "connected_at_utc", row[0], row[2], now ) connection.execute( "UPDATE bilibili_connections SET disconnected_at_utc=?, status='failed', disconnect_reason=? " "WHERE connection_id=? AND disconnected_at_utc IS NULL", (ended_at, "startup_reconciliation", row[1]), ) cls._upsert_live_session(connection, { "session_id": row[1], "kind": "bilibili_connection", "started_at_utc": row[2], "ended_at_utc": ended_at, "status": "failed", "source": "bilibili_connections", "payload_json": row[3] or "{}", }) recovered["bilibili_connections"] += 1 for row in connection.execute( "SELECT id, group_run_id, started_at_utc, payload_json FROM group_runs " "WHERE ended_at_utc IS NULL AND status='running' ORDER BY started_at_utc, id" ).fetchall(): try: payload = json.loads(row[3] or "{}") except (TypeError, ValueError, json.JSONDecodeError): payload = {} payload["reason"] = "startup_reconciliation" if stale_process_ids: payload["stale_process_session_ids"] = stale_process_ids ended_at = cls._reconciliation_end( connection, "group_runs", "started_at_utc", row[0], row[2], now ) connection.execute( "UPDATE group_runs SET ended_at_utc=?, status='interrupted', payload_json=? " "WHERE group_run_id=? AND ended_at_utc IS NULL", (ended_at, _json(payload), row[1]), ) recovered["group_runs"] += 1 connection.commit() return recovered @staticmethod def make_event_id(*parts: Any) -> str: raw = _json(parts).encode("utf-8", "replace") return hashlib.sha256(raw).hexdigest() def record_process_start( self, session_id: str, *, pid: int | None = None, host: str | None = None, app_version: str | None = None, metadata: Mapping[str, Any] | None = None, started_at_utc: str | None = None, ) -> bool: self._process_session_id = session_id started = started_at_utc or utc_now() return self._enqueue_record("__process_start__", { "session_id": session_id, "started_at_utc": started, "business_date": business_date(started), "pid": pid, "host": host, "app_version": app_version, "metadata_json": _json(metadata or {}), }) def record_process_end(self, session_id: str | None = None, *, exit_reason: str | None = None) -> bool: target = session_id or self._process_session_id if not target or not self._started or self._closing: return False try: self.queue.put_nowait(("__process_end__", {"session_id": target, "ended_at_utc": utc_now(), "exit_reason": exit_reason})) if target == self._process_session_id: self._process_session_id = None return True except Exception: self._dropped += 1 return False def upsert_live_session( self, session_id: str, *, kind: str | None = None, started_at_utc: datetime | str | None = None, ended_at_utc: datetime | str | None = None, status: str | None = None, duration_ms: int | None = None, source: str | None = None, payload: Mapping[str, Any] | None = None, ) -> bool: if not session_id: return False fields: dict[str, Any] = {"session_id": session_id} if kind is not None: fields["kind"] = kind if started_at_utc is not None: fields["started_at_utc"] = started_at_utc if ended_at_utc is not None: fields["ended_at_utc"] = ended_at_utc if status is not None: fields["status"] = status if duration_ms is not None: fields["duration_ms"] = max(0, int(duration_ms)) if source is not None: fields["source"] = source if payload is not None: fields["payload_json"] = _json(payload) if not self._started or self._closing: self._dropped += 1 return False try: self.queue.put_nowait(("__live_session_upsert__", _clean_fields(fields))) return True except Exception as exc: self._dropped += 1 self._last_error = f"{type(exc).__name__}: {exc}" return False def begin_live_session( self, session_id: str, *, kind: str = "live", started_at_utc: datetime | str | None = None, status: str = "running", source: str | None = None, payload: Mapping[str, Any] | None = None, ) -> bool: return self.upsert_live_session( session_id, kind=kind, started_at_utc=started_at_utc or utc_now(), status=status, source=source, payload=payload, ) def end_live_session( self, session_id: str, *, ended_at_utc: datetime | str | None = None, status: str = "ended", duration_ms: int | None = None, payload: Mapping[str, Any] | None = None, ) -> bool: return self.upsert_live_session( session_id, ended_at_utc=ended_at_utc or utc_now(), status=status, duration_ms=duration_ms, payload=payload, ) def upsert_user_snapshot( self, platform: str, platform_user_id: str, *, display_name: str | None = None, avatar_url: str | None = None, user_level: int | None = None, is_admin: bool = False, snapshot: Mapping[str, Any] | None = None, seen_at_utc: str | None = None, ) -> bool: if not self._started or self._closing: return False seen = seen_at_utc or utc_now() fields = _clean_fields({ "platform": platform, "platform_user_id": platform_user_id, "display_name": display_name, "avatar_url": avatar_url, "user_level": user_level, "is_admin": is_admin, "first_seen_at_utc": seen, "last_seen_at_utc": seen, "snapshot_json": _json(snapshot or {}), }) try: self.queue.put_nowait(("__user_upsert__", fields)) return True except Exception: self._dropped += 1 return False def load_user_points(self, platform: str) -> dict[str, int]: """同步读取某平台的用户积分(权威余额)。仅可在事件循环线程调用。""" connection = self._connection if not self._started or connection is None: return {} rows = connection.execute( "SELECT platform_user_id, points FROM users WHERE platform=?", (platform,), ).fetchall() return {str(row[0]): int(row[1]) for row in rows} async def set_user_points(self, platform: str, platform_user_id: str, points: int) -> bool: """权威写入用户积分:串行执行并等待落盘后才返回。""" if not self._started or self._closing: return False try: await self.execute(self._set_user_points_impl, platform, platform_user_id, points) return True except Exception as exc: self._last_error = f"{type(exc).__name__}: {exc}" return False @staticmethod def _set_user_points_impl(connection, platform, platform_user_id, points): now = utc_now() connection.execute( """INSERT INTO users ( platform, platform_user_id, display_name, first_seen_at_utc, last_seen_at_utc, points ) VALUES (?, ?, NULL, ?, ?, ?) ON CONFLICT(platform, platform_user_id) DO UPDATE SET points=excluded.points, last_seen_at_utc=MAX(users.last_seen_at_utc, excluded.last_seen_at_utc)""", (platform, platform_user_id, now, now, int(points)), ) connection.commit() def record_signin(self, event_id: str, *, platform: str, platform_user_id: str, points_awarded: int = 0, **fields: Any) -> bool: return self._record_timed("signin_events", "signed_at_utc", event_id=event_id, platform=platform, platform_user_id=platform_user_id, points_awarded=points_awarded, **fields) def record_point_transaction(self, transaction_id: str, amount: int, **fields: Any) -> bool: return self._record_timed("point_transactions", "occurred_at_utc", transaction_id=transaction_id, amount=amount, **fields) def record_queue_event(self, queue_name: str, action: str, **fields: Any) -> bool: fields.setdefault("event_id", self.make_event_id(queue_name, action, utc_now(), fields)) return self._record_timed("queue_events", "occurred_at_utc", queue_name=queue_name, action=action, **fields) def record_login_session(self, login_session_id: str, **fields: Any) -> bool: return self._record_timed("login_sessions", "started_at_utc", login_session_id=login_session_id, **fields) def record_group_run(self, group_run_id: str, **fields: Any) -> bool: return self._record_timed("group_runs", "started_at_utc", group_run_id=group_run_id, **fields) def record_song_request(self, request_id: str, **fields: Any) -> bool: return self._record_timed("song_requests", "requested_at_utc", request_id=request_id, **fields) def record_playback_session(self, playback_id: str, **fields: Any) -> bool: return self._record_timed("playback_sessions", "started_at_utc", playback_id=playback_id, **fields) def record_gift(self, event_id: str, **fields: Any) -> bool: return self._record_timed("gift_events", "occurred_at_utc", event_id=event_id, **fields) def record_broadcast_request(self, request_id: str, **fields: Any) -> bool: return self._record_timed("broadcast_requests", "requested_at_utc", request_id=request_id, **fields) def record_tts_request(self, request_id: str, **fields: Any) -> bool: return self._record_timed("tts_requests", "requested_at_utc", request_id=request_id, **fields) def record_bilibili_connection(self, connection_id: str, **fields: Any) -> bool: return self._record_timed("__bilibili_connection__", "connected_at_utc", connection_id=connection_id, **fields) def record_danmu(self, event_id: str, **fields: Any) -> bool: return self._record_timed("danmu_events", "occurred_at_utc", event_id=event_id, **fields) def record_service_state(self, service_name: str, new_state: str, **fields: Any) -> bool: fields.setdefault("event_id", self.make_event_id(service_name, new_state, utc_now(), fields)) return self._record_timed("service_state_events", "occurred_at_utc", service_name=service_name, new_state=new_state, **fields) def record_admin_audit(self, event_id: str, action: str, **fields: Any) -> bool: if fields.get("remote_address"): fields["remote_address_hash"] = hashlib.sha256(str(fields.pop("remote_address")).encode()).hexdigest() return self._record_timed("admin_audit_events", "occurred_at_utc", event_id=event_id, action=action, **fields) def _record_timed(self, table: str, time_column: str, **fields: Any) -> bool: timestamp = _utc_iso(fields.get(time_column) or utc_now()) fields[time_column] = timestamp fields.setdefault("business_date", business_date(timestamp)) payload = fields.pop("payload", None) if payload is not None: fields["payload_json"] = _json(payload) return self._enqueue_record(table, fields) async def import_once( self, source_key: str, source_type: str, records: Sequence[tuple[str, Mapping[str, Any]]], *, fingerprint: str | None = None, checkpoint_key: str = "default", checkpoint_value: str | None = None, metadata: Mapping[str, Any] | None = None, ) -> bool: """幂等导入:相同 source_key+fingerprint 成功导入后不重复执行。""" if not self._started or self._closing: return False loop = asyncio.get_running_loop() completed = loop.create_future() try: await self.queue.put(("__import__", { "source_key": source_key, "source_type": source_type, "fingerprint": fingerprint, "records": [(table, _clean_fields(fields)) for table, fields in records], "checkpoint_key": checkpoint_key, "checkpoint_value": checkpoint_value, "metadata_json": _json(metadata or {}), "future": completed, })) return bool(await completed) except Exception as exc: self._last_error = f"{type(exc).__name__}: {exc}" return False async def query_daily_statistics(self, start_date: str | None = None, end_date: str | None = None) -> list[dict[str, Any]]: """查询每日统计视图;查询失败返回空列表。""" try: return await self.query_daily_statistics_strict(start_date, end_date) except Exception as exc: self._last_error = f"{type(exc).__name__}: {exc}" return [] async def query_daily_statistics_strict( self, start_date: str | None = None, end_date: str | None = None, ) -> list[dict[str, Any]]: """查询每日统计视图;查询失败时向调用方抛出异常。""" await self.flush() rows = await asyncio.to_thread(self._query_daily_sync, start_date, end_date) self._last_error = None return rows def runtime_snapshot(self) -> dict[str, Any]: """返回不执行磁盘完整性检查的轻量运行状态。""" return { "started": self._started, "closing": self._closing, "writer_alive": bool(self._writer_task and not self._writer_task.done()), "queue_size": self.queue.qsize(), "queue_capacity": self.queue.maxsize, "last_error": self._last_error, } def _query_daily_sync(self, start_date: str | None, end_date: str | None) -> list[dict[str, Any]]: clauses: list[str] = [] params: list[str] = [] if start_date: clauses.append("business_date >= ?") params.append(start_date) if end_date: clauses.append("business_date <= ?") params.append(end_date) where = f" WHERE {' AND '.join(clauses)}" if clauses else "" uri = f"file:{self.database_path.as_posix()}?mode=ro" connection = sqlite3.connect(uri, uri=True, timeout=5.0) try: connection.row_factory = sqlite3.Row connection.execute("PRAGMA busy_timeout=5000") cursor = connection.execute(f"SELECT * FROM daily_live_statistics{where} ORDER BY business_date", params) rows = [dict(row) for row in cursor] columns = [item[0] for item in cursor.description or ()] live_rows = connection.execute( "SELECT kind, started_at_utc, ended_at_utc FROM live_sessions WHERE ended_at_utc IS NOT NULL" ).fetchall() duration_maps = { "live_duration_ms": self._merged_daily_interval_durations([ (row[1], row[2]) for row in live_rows if row[0] in {"live", "stream", "streaming"} ]), "process_duration_ms": self._merged_daily_interval_durations([ (row[1], row[2]) for row in live_rows if row[0] == "process" ]), "bilibili_connection_duration_ms": self._merged_daily_interval_durations([ (row[1], row[2]) for row in live_rows if row[0] == "bilibili_connection" ]), "group_run_total_duration_ms": self._merged_daily_interval_durations([ (row[0], row[1]) for row in connection.execute( "SELECT started_at_utc, ended_at_utc FROM group_runs WHERE ended_at_utc IS NOT NULL" ).fetchall() ]), } by_day = {str(row["business_date"]): row for row in rows} duration_days = set().union(*(values.keys() for values in duration_maps.values())) for day in duration_days: if (start_date and day < start_date) or (end_date and day > end_date): continue if day not in by_day: by_day[day] = {column: 0 for column in columns} by_day[day]["business_date"] = day for day, row in by_day.items(): for column, values in duration_maps.items(): row[column] = values.get(day, 0) return [by_day[day] for day in sorted(by_day)] finally: connection.close() async def health_snapshot(self) -> dict[str, Any]: snapshot: dict[str, Any] = { "started": self._started, "closing": self._closing, "writer_alive": bool(self._writer_task and not self._writer_task.done()), "queue_size": self.queue.qsize(), "queue_capacity": self.queue.maxsize, "written": self._written, "dropped": self._dropped, "last_error": self._last_error, "database_path": str(self.database_path), "process_session_id": self._process_session_id, } if self._started: try: await self.flush() snapshot["database"] = await asyncio.to_thread(self._health_sync) except Exception as exc: snapshot["database"] = {"ok": False, "error": f"{type(exc).__name__}: {exc}"} else: snapshot["database"] = {"ok": False, "error": "not_started"} return snapshot def _health_sync(self) -> dict[str, Any]: uri = f"file:{self.database_path.as_posix()}?mode=ro" connection = sqlite3.connect(uri, uri=True, timeout=5.0) try: result = connection.execute("PRAGMA quick_check").fetchone() return { "ok": bool(result and result[0] == "ok"), "quick_check": result[0] if result else None, "size_bytes": self.database_path.stat().st_size if self.database_path.exists() else 0, "schema_version": connection.execute("PRAGMA user_version").fetchone()[0], } finally: connection.close()