"""账号凭据与账号池(通用,支持错误分类换号与冷却)。 每个账号一份 ``account/.json``,启动时全量加载;请求时 round-robin 轮换可用子集。 失效按 :class:`FailReason` 分类:不可恢复(认证失败/封号)→ disabled 永久剔除; 可恢复(额度耗尽/人机验证)→ 冷却一段时间后自动恢复。 账号扩展字段(目标网站专属凭据:cookie/token/project_id 等)通过 ``extra=allow`` 容纳。 """ from __future__ import annotations import json import threading import time from collections.abc import Callable from enum import StrEnum from pathlib import Path from pydantic import BaseModel, ConfigDict from app._compat import flock_ex class FailReason(StrEnum): """账号失效原因分类(见 app/deps.py 的 classify_failure)。""" AUTH_FAILED = "auth_failed" # 认证失败(401/403)→ dead BANNED = "banned" # 账号被封 → dead QUOTA_EXHAUSTED = "quota_exhausted" # 额度耗尽 → cooling CF_CHALLENGE = "cf_challenge" # 人机验证 / CF 拦截 → cooling INSUFFICIENT_BALANCE = "insufficient_balance" # 402 永久无余额 → dead(lifetime 100 分花完不涨) # 不可恢复的原因:标记 disabled 永久剔除;其余进入冷却。 # AUTH_FAILED 不在此列:Privy token 过期属可恢复(重新登录刷新即可), # 且 auth.py 的过期预检在请求前即拦截已过期 token,不消耗上游调用。 _DEAD_REASONS: frozenset[FailReason] = frozenset({FailReason.BANNED, FailReason.INSUFFICIENT_BALANCE}) # 默认冷却时长(秒) COOLDOWN_SECONDS: float = 600.0 # 可恢复失效的冷却策略:action = "cooldown" / "disable" _COOLDOWN_POLICY: dict[str, str] = {"action": "cooldown"} # 按 FailReason 的冷却时长覆盖;未设置的原因回退到 COOLDOWN_SECONDS _COOLDOWN_SECONDS_MAP: dict[FailReason, float] = {} def set_cooldown_policy( action: str, seconds: float | None = None, seconds_map: dict[FailReason, float] | None = None, ) -> None: """设置非致命失效的处理策略与冷却时长。 - ``"cooldown"``:冷却后恢复;可用 ``seconds`` 设全局默认值,或用 ``seconds_map`` 按原因覆盖。 - ``"disable"``:直接标记 disabled,适合月额度或永不过期额度的上游。 """ if action not in ("cooldown", "disable"): raise ValueError(f"action must be 'cooldown' or 'disable', got {action!r}") _COOLDOWN_POLICY["action"] = action if seconds is not None: _COOLDOWN_SECONDS_MAP.setdefault(FailReason.QUOTA_EXHAUSTED, seconds) _COOLDOWN_SECONDS_MAP.setdefault(FailReason.CF_CHALLENGE, seconds) if seconds_map: _COOLDOWN_SECONDS_MAP.update(seconds_map) class Account(BaseModel): """单个账号凭据。目标网站专属字段通过 extra=allow 容纳(见 upstream/account_fields.py)。""" model_config = ConfigDict(extra="allow") name: str source_email: str = "" created_at: str = "" disabled: bool = False fail_reason: FailReason | None = None cooldown_until: float = 0.0 # Unix 时间戳;0=无冷却 conversation_id: str = "" # 生图会话 UUID(首次生图自动生成并持久化;普通对话不用) @classmethod def from_file(cls, path: Path) -> Account: data = json.loads(path.read_text(encoding="utf-8")) if "name" not in data: # account/.json 文件名即账号名,避免必填字段缺失导致启动崩溃 data["name"] = path.stem return cls(**data) class AccountPool: """账号池:round-robin 轮换可用子集(跳过 disabled 与未到期冷却),同步线程安全。""" def __init__(self, accounts: list[Account], account_dir: Path) -> None: # 按 name 排序保证 round-robin 游标确定性 self._all: list[Account] = sorted(accounts, key=lambda a: a.name) self._dir = account_dir self._idx = 0 self._lock = threading.Lock() self._on_changed: Callable[[], None] | None = None @property def account_dir(self) -> Path: """账号文件目录(与 :meth:`reload` 扫描的目录一致;上传/删除等写盘都用它)。""" return self._dir @classmethod def load(cls, account_dir: Path) -> AccountPool: """扫 account_dir 下的 *.json 构造账号池。 目录不存在或无 json **不抛错**(返回空池):HF 云端首启无账号, 面板上传后热加载即用;本地误配也只是空池,请求时报“无可用账号”。 """ if not account_dir.is_dir(): return cls([], account_dir) files = sorted(account_dir.glob("*.json")) accounts = [Account.from_file(f) for f in files] return cls(accounts, account_dir) def set_on_changed(self, callback: Callable[[], None] | None) -> None: """设置账号池变化回调(用于自动补账号任务即时唤醒)。""" self._on_changed = callback def _notify_changed(self) -> None: """轻量通知回调;在锁内调用,不得做耗时/阻塞操作。""" if self._on_changed is not None: try: self._on_changed() except Exception: # noqa: BLE001 pass def all(self) -> list[Account]: """返回全部账号(含 disabled / 冷却中)。""" return list(self._all) def available_count(self) -> int: """当前可用账号数(未 disabled 且未在冷却期)。""" return len(self._available()) def _is_available(self, a: Account) -> bool: if a.disabled: return False if a.cooldown_until and a.cooldown_until > time.time(): return False return True def _available(self) -> list[Account]: return [a for a in self._all if self._is_available(a)] def next(self) -> Account: """round-robin 返回下一个可用账号;无可用抛 RuntimeError。""" with self._lock: avail = self._available() if not avail: raise RuntimeError("无可用账号(全部 disabled 或冷却中)") if self._idx >= len(avail): self._idx = 0 acc = avail[self._idx] self._idx = (self._idx + 1) % len(avail) return acc def _save(self, acc: Account) -> None: """原子写 account/.json(内部调用,需已持有锁)。""" self._dir.mkdir(parents=True, exist_ok=True) target = self._dir / f"{acc.name}.json" tmp = target.with_suffix(".json.tmp") payload = acc.model_dump(mode="json") with tmp.open("w", encoding="utf-8") as f: json.dump(payload, f, ensure_ascii=False, indent=2) f.write("\n") flock_ex(f.fileno()) tmp.replace(target) def mark_failed(self, acc: Account, reason: FailReason) -> None: """按原因标记账号:不可恢复→disabled(剔除),可恢复→按策略冷却或禁用。原子写回 json。""" with self._lock: acc.fail_reason = reason if reason in _DEAD_REASONS: acc.disabled = True elif _COOLDOWN_POLICY.get("action") == "disable": acc.disabled = True else: seconds = _COOLDOWN_SECONDS_MAP.get(reason, COOLDOWN_SECONDS) acc.cooldown_until = time.time() + seconds self._save(acc) self._notify_changed() def cooldown(self, acc: Account, seconds: float) -> None: """主动冷却账号(不标记失败原因)。生图成功后调用:避免同一账号连击触发上游限速。""" with self._lock: acc.cooldown_until = time.time() + seconds self._save(acc) def mark_disabled(self, acc: Account) -> None: """兼容旧名:等价于 mark_failed(AUTH_FAILED)。""" self.mark_failed(acc, FailReason.AUTH_FAILED) def enable(self, acc: Account) -> None: """手动恢复账号:清 disabled / fail_reason / cooldown,重新纳入轮换。原子写回 json。""" with self._lock: acc.disabled = False acc.fail_reason = None acc.cooldown_until = 0.0 self._save(acc) self._notify_changed() def disable(self, acc: Account) -> None: """手动禁用账号:标记 disabled 并清冷却(永久剔除轮换,直到手动 enable)。原子写回 json。""" with self._lock: acc.disabled = True acc.cooldown_until = 0.0 self._save(acc) self._notify_changed() def add_or_update(self, acc: Account) -> None: """新增或替换账号,原子写盘并同步内存列表。""" with self._lock: self._save(acc) others = [a for a in self._all if a.name != acc.name] self._all = sorted([*others, acc], key=lambda a: a.name) avail_len = len(self._available()) if avail_len and self._idx >= avail_len: self._idx = 0 self._notify_changed() def remove(self, name: str) -> bool: """删除账号 json 并从内存池移除;存在返回 True。""" with self._lock: target = self._dir / f"{name}.json" existed = target.is_file() if existed: target.unlink() before = len(self._all) self._all = [a for a in self._all if a.name != name] avail_len = len(self._available()) if avail_len and self._idx >= avail_len: self._idx = 0 changed = existed or before != len(self._all) if changed: self._notify_changed() return changed def reload(self) -> None: """重新从磁盘加载全部账号,替换内存池。""" with self._lock: files = sorted(self._dir.glob("*.json")) self._all = [Account.from_file(f) for f in files] self._idx = 0