# -*- coding: utf-8 -*- # Author : Charley # Python : 3.12.10 # Date : 2026/09/08 """得卡 DECA · 账号池(框架版,2026/09/08 搭)。 把原来集中在单账号 token.json 的登录态,改为「多账号入池、按任务取号、异常自动切号、冷却复用」, 把带 token 请求压力平摊到多账号 + 各账号专属静态 IP,根治 9/7「单账号被风控击穿致整链停摆」。 方案与请求量测算见 HANDOFF.md 及 docs/账号池方案_得卡DECA_20260908.md。 数据模型:一张 deca_account_record 表(DDL 见同目录 deca_account_ddl.sql)+ 本模块。 每个账号自带独立 access/refresh/token_exp + 专属 proxy_url(1 号 1 IP 永久绑定)+ 健康状态。 线上生命周期原则(重要): · 线上**只做 refresh 续期,绝不自动跑密码登录**——密码登录会撞阿里云滑块、无人值守过不去。 续期失败即判 dead、摘池、企微通知,自动切下一个 healthy 号,采集不中断。 · 补货靠接码平台自动补货:healthy 数不足时 auto_replenish 经接码平台租号 + 短信登录(=注册)入库 + 绑定空闲专属 IP(见 register_one_via_isms / spiders/replenish_accounts.py),无需人工灌 token。 取号粒度:进程级独占——每个常驻任务启动时 acquire 一个号、独占整个运行期(保持账号↔IP 稳定降风控), 仅在该号挂了才 switch 换号。购买记录按商家分片时,每个分片进程各 acquire 一个号(带 owner_tag 区分)。 账号供给(2026/09/09 落地):得卡无独立注册接口——短信验证码登录「登录即注册」,故: - sms_send() / sms_login() 发验证码 + 短信登录入库(登录即注册;新号 INSERT、旧号补货 UPDATE) - register() 等价于 sms_login(两步:先 sms_send 收码,再带 code 调用),离线人工补货/开户主路径 - _login() 密码登录(撞阿里云滑块),仅极端兜底,线上默认不调;短信登录不撞滑块、优先用它 接入采集脚本(把 core.do_request 的 need_auth 取号改为走本池)属后续 step3,未在本框架内改 core。 """ import os import time from datetime import datetime, timedelta from loguru import logger # ==================== 配置 ==================== # 全部业务/敏感配置集中在 settings.py(接口路径 / 接码 Key / 代理账密 IP 池 / 阈值),改配置只动那处。 from settings import * # noqa: F401,F403 (T_ACCOUNT/LEASE_TTL_SEC/DEAD_CODES/ISMS_*/PROXY_*/*_PATH 等) class AccountPool: """账号池:从 deca_account_record 取号、续期、异常切号、状态维护。 单个常驻任务进程持有一个实例:启动 acquire 一个号独占,运行期用 ensure_access 保证 token 有效, 请求失败调 report_failure(达阈值 cooling/dead)→ switch 换号;结束 release 归还。 """ def __init__(self, pool, log=None, task_tag: str = ""): """初始化账号池。 Args: pool (MySQLConnectionPool): MySQL 连接池。 log (optional): 日志对象。Defaults to None(回落 loguru.logger)。 task_tag (str, optional): 本进程任务标签,写入 owner_tag 便于排障与独占区分 (如 "buy_record:881226408" / "team" / "onsale_alert")。Defaults to ""。 """ self.pool = pool self.log = log or logger self.task_tag = task_tag self.pid = os.getpid() self.account = None # 当前持有的账号行 dict;None 表示未持号 # ---------- 取号 / 归还 / 切号 ---------- def acquire(self) -> dict | None: """从池中租一个 healthy 且未被独占(或租约已过期)的账号,打上本进程租约并返回。 原子抢占:一条 UPDATE 按 last_used_at 升序挑一个可租的号打租约(owner_pid/lease_until), 再按本进程 pid + task_tag 回读刚租到的行。取号成功后写入 self.account。 Returns: dict | None: 账号行 dict(含 id/phone/access_token/refresh_token/token_exp/proxy_url 等); 池中无可用 healthy 号时返回 None(调用方应告警/退避)。 """ lease_until = datetime.now() + timedelta(seconds=LEASE_TTL_SEC) # 挑一个 healthy 且(未被占 或 租约已过期)的号,按最久未用优先做负载均衡,打上本进程租约 self.pool.update_one( f"UPDATE {T_ACCOUNT} SET owner_tag=%s, owner_pid=%s, lease_until=%s, last_used_at=NOW() " "WHERE status='healthy' AND (owner_pid IS NULL OR lease_until IS NULL OR lease_until < NOW()) " "ORDER BY last_used_at ASC LIMIT 1", # NULL 升序排最前:从未用过的号优先取,做负载均衡 (self.task_tag, self.pid, lease_until.strftime("%Y-%m-%d %H:%M:%S"))) # 回读本进程刚租到的号(owner_pid=自身 pid 唯一,别的进程不会用我的 pid) row = self.pool.select_one( f"SELECT id, phone, password, country_code, uid, access_token, refresh_token, token_exp, " f"proxy_url, status FROM {T_ACCOUNT} WHERE owner_pid=%s AND owner_tag=%s AND status='healthy' " f"ORDER BY lease_until DESC LIMIT 1", (self.pid, self.task_tag)) if not row: self.log.error(f"[账号池] 无可用 healthy 账号(task={self.task_tag})") self.account = None return None self.account = self._row_to_dict(row) self.log.info(f"[账号池] 取号成功 id={self.account['id']} phone={self.account['phone']} " f"proxy={self.account.get('proxy_url')} task={self.task_tag}") return self.account def acquire_random(self, lease_sec: int = LEASE_TTL_SEC) -> dict | None: """从 healthy 池随机挑一个未被占的账号,打上租约并返回(每批随机取号用)。 与 acquire(进程级独占、按最久未用负载均衡)的区别:本方法按 RAND() 随机挑号, 供「每批随机」——每 N 个请求换一个随机号,把带 token 请求量摊平到所有号、降单号风控面。 租约默认用 LEASE_TTL_SEC(较长,覆盖持号整个使用期):一个进程持号直到换号/退出, 期间租约不过期,别的进程不会抢走——避免「持号时长 > 租约 → 被抢 → 两进程同号续期 踩废一次性轮换的 refreshToken → 误判 10010 dead」(2026/09/09 修:原 300s 过短,慢/稀疏 任务如 team/onsale_alert 会在换号前租约过期)。进程崩溃后租约到期自动可被再租,不死占。 并发抢占:先 RAND() 选一个候选 id,再带「仍空闲」条件 UPDATE 认领;认领失败(被别的进程抢走) 重试几次换一个候选。认领成功后回读该行写入 self.account。 Args: lease_sec (int, optional): 租约时长(秒),覆盖持号整个使用期。Defaults to LEASE_TTL_SEC(6h)。 Returns: dict | None: 账号行 dict;池中无可用 healthy 号或多次抢占失败返回 None。 """ lease_until = (datetime.now() + timedelta(seconds=lease_sec)).strftime("%Y-%m-%d %H:%M:%S") for _ in range(5): # 抢占重试:候选被别的进程抢走时换一个 cand = self.pool.select_one( f"SELECT id FROM {T_ACCOUNT} WHERE status='healthy' " "AND (owner_pid IS NULL OR lease_until IS NULL OR lease_until < NOW()) " "ORDER BY RAND() LIMIT 1") if not cand: self.log.error(f"[账号池] 无可用 healthy 账号(随机取号,task={self.task_tag})") self.account = None return None cid = cand[0] # 带「仍空闲」条件认领:若被别的进程抢先占了,本 UPDATE 命中 0 行,回读会拿不到本 pid → 重试 self.pool.update_one( f"UPDATE {T_ACCOUNT} SET owner_tag=%s, owner_pid=%s, lease_until=%s, last_used_at=NOW() " "WHERE id=%s AND status='healthy' " "AND (owner_pid IS NULL OR lease_until IS NULL OR lease_until < NOW())", (self.task_tag, self.pid, lease_until, cid)) row = self.pool.select_one( f"SELECT id, phone, password, country_code, uid, access_token, refresh_token, " f"token_exp, proxy_url, status FROM {T_ACCOUNT} WHERE id=%s AND owner_pid=%s", (cid, self.pid)) if row: self.account = self._row_to_dict(row) self.log.info(f"[账号池] 随机取号 id={self.account['id']} phone={self.account['phone']} " f"proxy={self.account.get('proxy_url')} lease={lease_sec}s task={self.task_tag}") return self.account self.log.warning(f"[账号池] 随机取号多次抢占失败(并发激烈,task={self.task_tag})") self.account = None return None def release(self): """归还当前账号(清空租约 owner_pid/lease_until),进程正常退出时调用。""" if not self.account: return self.pool.update_one( f"UPDATE {T_ACCOUNT} SET owner_pid=NULL, lease_until=NULL WHERE id=%s", (self.account["id"],)) self.log.info(f"[账号池] 归还账号 id={self.account['id']}") self.account = None def switch(self) -> dict | None: """当前号不可用时换号:释放当前号的租约(不改其 status,status 由 report_* 决定)后重新 acquire。 Returns: dict | None: 新账号行 dict;池中无其他可用号时返回 None。 """ old = self.account if old: # 只解租约,不动 status——是否 cooling/dead 由 report_failure 已经写好;healthy 号只是让给别人 self.pool.update_one( f"UPDATE {T_ACCOUNT} SET owner_pid=NULL, lease_until=NULL WHERE id=%s AND owner_pid=%s", (old["id"], self.pid)) self.account = None new = self.acquire() if new: self.log.warning(f"[账号池] 切号:{old['phone'] if old else '-'} → {new['phone']}") else: self.log.error("[账号池] 切号失败:池中已无可用 healthy 号,请尽快离线补货") return new # ---------- token 续期 ---------- @staticmethod def _proxies(proxy_url) -> dict | None: """把该号 proxy_url 组装成 requests proxies(让续期/短信登录也走该号专属 IP)。 Args: proxy_url: 该号绑定的代理 URL;None/空表示直连。 Returns: dict | None: {"http":.., "https":..};proxy_url 为空返回 None。 """ return {"http": proxy_url, "https": proxy_url} if proxy_url else None def ensure_access(self, log=None) -> str | None: """确保当前账号 access 有效:未临期直接返回;临期用该号 refresh 续期并回写 DB。 线上主路径:只 refresh、不自动登录。续期失败即判该号 dead 并切号(见 report_failure)。 Args: log (optional): 日志对象。Defaults to None(用 self.log)。 Returns: str | None: 有效 access token;无持号/续期失败返回 None(调用方据此 switch 换号)。 """ log = log or self.log acc = self.account if not acc: log.error("[账号池] ensure_access 时未持号") return None # 未临期:直接用缓存 access if acc.get("access_token") and acc.get("token_exp") and time.time() < acc["token_exp"] - ACCESS_SKEW_SEC: return acc["access_token"] # 临期/缺失:用该号 refresh 续期(need_auth=False,不依赖全局 token)。走该号专属 proxy_url, # 让续期这类登录态操作也从该号绑定的 IP 发出(不暴露本机 IP,防 20 号在登录态层被同源关联)。 if not acc.get("refresh_token"): log.warning(f"[账号池] id={acc['id']} 无 refresh_token,无法续期") self.report_failure("no_refresh_token", dead=True) return None try: import deca_sold_core as core # 惰性 import:避免 import 本模块即触发 core 副作用 resp = core.do_request(log, REFRESH_PATH, {"refreshToken": acc["refresh_token"]}, need_auth=False, use_proxy=False, proxy_override=self._proxies(acc.get("proxy_url"))) code = (resp or {}).get("code") msg = (resp or {}).get("msg") data = (resp or {}).get("data") or {} access = data.get("accessToken") if not access: if code in DEAD_CODES: # 实测 10010=refresh 失效(被轮换/过期/长期未登录)→ 确定判 dead 摘池(线上不自动登录) log.warning(f"[账号池] id={acc['id']} 续期判死 code={code} msg={msg}") self.report_failure(f"refresh_dead:code={code}", dead=True) else: # 其他未知失败(限流/网络/临时):软失败累计,达 MAX_FAIL 才 cooling,不急着判死 log.warning(f"[账号池] id={acc['id']} 续期失败(非判死码) code={code} msg={msg}") self.report_failure(f"refresh_fail:code={code}", dead=False) return None expires_in = int(data.get("expiresIn") or 900) new_refresh = data.get("refreshToken") or acc["refresh_token"] # 一次性轮换:有新的就换 new_exp = int(time.time() + expires_in) # 回写 DB 该号最新 token self.pool.update_one( f"UPDATE {T_ACCOUNT} SET access_token=%s, refresh_token=%s, token_exp=%s, " "last_success_at=NOW(), fail_count=0 WHERE id=%s", (access, new_refresh, new_exp, acc["id"])) acc["access_token"], acc["refresh_token"], acc["token_exp"] = access, new_refresh, new_exp log.info(f"[账号池] id={acc['id']} 续期成功,有效 {expires_in}s") return access except Exception as e: # 网络/超时等异常(非明确判死码)→ 软失败累计,不永久判死(网络问题 ≠ 账号死) log.warning(f"[账号池] id={acc['id']} 续期异常(按软失败处理): {e}") self.report_failure(f"refresh_exc:{e}", dead=False) return None # ---------- 状态机 ---------- def report_success(self): """标记当前账号一次成功(清零 fail_count、刷新 last_success_at)。""" if not self.account: return self.pool.update_one( f"UPDATE {T_ACCOUNT} SET fail_count=0, last_success_at=NOW(), last_error=NULL WHERE id=%s", (self.account["id"],)) def report_failure(self, error: str = "", dead: bool = False): """记录当前账号一次失败:dead=True 直接判死摘池;否则累计失败,达 MAX_FAIL 转 cooling。 Args: error (str, optional): 错误信息,写入 last_error 供排障。Defaults to ""。 dead (bool, optional): True=登录态死了(续期失败等)直接置 dead 并企微通知人工。 False=软失败(限流/偶发),累计到 MAX_FAIL 才转 cooling。Defaults to False。 """ if not self.account: return aid = self.account["id"] err = (error or "")[:255] if dead: self.pool.update_one( f"UPDATE {T_ACCOUNT} SET status='dead', owner_pid=NULL, lease_until=NULL, " "last_error=%s WHERE id=%s", (err, aid)) self.log.error(f"[账号池] id={aid} 判 dead 摘池({err})→ 需离线补货") # TODO(step3): 企微通知人工「账号 {phone} dead,请离线补货」,复用 auto_send_wx_msg。 self.account = None return # 软失败累计 self.pool.update_one( f"UPDATE {T_ACCOUNT} SET fail_count=fail_count+1, last_error=%s WHERE id=%s", (err, aid)) row = self.pool.select_one(f"SELECT fail_count FROM {T_ACCOUNT} WHERE id=%s", (aid,)) fc = int(row[0]) if row and row[0] is not None else 0 if fc >= MAX_FAIL: cd = (datetime.now() + timedelta(seconds=COOLDOWN_SEC)).strftime("%Y-%m-%d %H:%M:%S") self.pool.update_one( f"UPDATE {T_ACCOUNT} SET status='cooling', cooldown_until=%s, owner_pid=NULL, " "lease_until=NULL WHERE id=%s", (cd, aid)) self.log.warning(f"[账号池] id={aid} 连续失败 {fc} 次转 cooling,冷却到 {cd}") self.account = None def revive_cooling(self) -> int: """巡检:把冷却到期的 cooling 账号恢复为 healthy 重新入池(dead 不自动恢复,需人工)。 由一个后台巡检循环定期调用(如每几分钟一次)。 Returns: int: 本次恢复的账号数。 """ self.pool.update_one( f"UPDATE {T_ACCOUNT} SET status='healthy', fail_count=0, cooldown_until=NULL " "WHERE status='cooling' AND cooldown_until IS NOT NULL AND cooldown_until < NOW()", ()) # 恢复条数取决于底层驱动是否回报 affected rows;此处仅执行,数量另查(保持接口简单) return 0 def keepalive_idle(self, idle_hours: int = 6, log=None) -> tuple[int, int]: """token 保活巡检:对闲置超过 idle_hours 的 healthy 号主动刷新一次,避免 refresh_token 闲到过期被判死。 背景(2026/09/15):access token 仅 15 分钟、refresh_token 一次性轮换;号若长时间没被业务请求抽到, refresh_token 会自然过期,下次一用即 10010 判死(实测闲置 30~77h 死)。故后台每小时扫一遍: 把「healthy 且 last_success_at 超过 idle_hours 没成功 且 当前未被占用」的号逐个强制刷新—— 刷成功则 token 轮换保温热;刷回 10010 则本就已死、判 dead 交自动补货补。 并发安全:每个号先带「仍空闲」条件抢租约(claim)、抢到才刷、刷完释放,避免与业务进程并发刷同一号 踩废其一次性轮换的 refreshToken。 Args: idle_hours (int, optional): 闲置阈值(小时),last_success_at 早于此才刷。Defaults to 6。 log (optional): 日志对象。Defaults to None(用 self.log)。 Returns: tuple[int, int]: (温热成功数, 判死数)。 """ log = log or self.log rows = self.pool.select_all( f"SELECT id FROM {T_ACCOUNT} WHERE status='healthy' " "AND (last_success_at IS NULL OR last_success_at < NOW() - INTERVAL %s HOUR) " "AND (owner_pid IS NULL OR lease_until IS NULL OR lease_until < NOW())", (idle_hours,)) or [] warmed, died = 0, 0 for (aid,) in rows: lease_until = (datetime.now() + timedelta(seconds=120)).strftime("%Y-%m-%d %H:%M:%S") # 短租,够刷新用 # 带「仍空闲」条件抢占:抢到才刷(防与业务进程并发刷同一号) self.pool.update_one( f"UPDATE {T_ACCOUNT} SET owner_tag=%s, owner_pid=%s, lease_until=%s " "WHERE id=%s AND status='healthy' " "AND (owner_pid IS NULL OR lease_until IS NULL OR lease_until < NOW())", (self.task_tag, self.pid, lease_until, aid)) row = self.pool.select_one( f"SELECT id, phone, password, country_code, uid, access_token, refresh_token, " f"token_exp, proxy_url, status FROM {T_ACCOUNT} WHERE id=%s AND owner_pid=%s AND status='healthy'", (aid, self.pid)) if not row: continue # 抢占失败(被业务进程抢走)或状态已变,跳过 self.account = self._row_to_dict(row) self.account["token_exp"] = 0 # 置 0 强制走刷新分支,确保轮换 refresh_token(而非命中未临期缓存) access = self.ensure_access(log) # 内部:刷成功回写 token;10010 则 report_failure(dead) 已摘池 if access: warmed += 1 # 仍 healthy:解租约归还(ensure_access 判死时已清 self.account+解租约,不会走到这) if self.account is not None: self.pool.update_one( f"UPDATE {T_ACCOUNT} SET owner_pid=NULL, lease_until=NULL WHERE id=%s AND owner_pid=%s", (self.account["id"], self.pid)) else: died += 1 self.account = None if warmed or died: log.info(f"[保活] 闲置>{idle_hours}h 刷新:温热 {warmed} 个,判死 {died} 个") return warmed, died # ---------- 账号供给(依赖明天抓的接口,占位)---------- def _login(self, log=None) -> bool: """密码登录换取全新一套 token(照 deca_sold_core.login 的接口写实)。 ⚠️ 线上默认不调用——密码登录会撞阿里云滑块,无人值守过不去(9/7 就死在这)。仅供 「离线补货 / 手动兜底」场景显式调用。线上续期失败一律走 report_failure(dead=True)。 Args: log (optional): 日志对象。Defaults to None。 Returns: bool: 登录成功并回写 token 返回 True;失败返回 False。 """ log = log or self.log acc = self.account if not acc: return False body = {"countryCode": acc.get("country_code") or "86", "password": acc.get("password"), "phone": acc.get("phone")} try: import deca_sold_core as core resp = core.do_request(log, LOGIN_PATH, body, need_auth=False, use_proxy=False) data = (resp or {}).get("data") or {} access, refresh = data.get("accessToken"), data.get("refreshToken") if not access or not refresh: log.warning(f"[账号池] id={acc['id']} 密码登录未返回 token: {resp.get('msg') if resp else None}") return False new_exp = int(time.time() + int(data.get("expiresIn") or 900)) self.pool.update_one( f"UPDATE {T_ACCOUNT} SET access_token=%s, refresh_token=%s, token_exp=%s, " "status='healthy', fail_count=0, last_success_at=NOW() WHERE id=%s", (access, refresh, new_exp, acc["id"])) acc.update({"access_token": access, "refresh_token": refresh, "token_exp": new_exp}) log.info(f"[账号池] id={acc['id']} 密码登录成功(离线补货)") return True except Exception as e: log.warning(f"[账号池] id={acc['id']} 密码登录异常: {e}") return False def sms_send(self, phone: str, country_code: str = "86", proxy_url: str = None, log=None) -> bool: """向指定手机号发送登录短信验证码(登录/注册第一步)。 need_auth=False,不依赖任何登录态。用于「补货 / 新号开户」流程第一步:调用后手机收到验证码, 再调 sms_login。走该号将绑的专属 proxy_url,保证发码也从该号 IP 发出(与后续登录同源)。 Args: phone (str): 目标手机号。 country_code (str, optional): 区号。Defaults to "86"。 proxy_url (str, optional): 该号专属 IP(走它发码,与登录同源)。Defaults to None(直连)。 log (optional): 日志对象。Defaults to None(用 self.log)。 Returns: bool: code==0 视为发送成功返回 True;否则 False。 """ log = log or self.log body = {"countryCode": country_code, "phone": phone} try: import deca_sold_core as core # 惰性 import:避免 import 本模块即触发 core 副作用 resp = core.do_request(log, SMS_SEND_PATH, body, need_auth=False, use_proxy=False, proxy_override=self._proxies(proxy_url)) if resp and resp.get("code") == 0: log.info(f"[账号池] 验证码已发送 phone={phone}") return True log.warning(f"[账号池] 发验证码失败 phone={phone}: {resp.get('msg') if resp else None}") return False except Exception as e: log.warning(f"[账号池] 发验证码异常 phone={phone}: {e}") return False def sms_login(self, phone: str, code: str, country_code: str = "86", proxy_url: str = None, password: str = None, log=None) -> dict | None: """短信验证码登录(登录即注册)并入库为 healthy 账号(新号 INSERT / 已存在号 UPDATE token)。 得卡无独立注册接口:新手机号首次短信登录即自动开户,故本方法同时承担「注册新号」与 「离线补货已死账号」两职。成功后把 access/refresh/token_exp/uid 落库 deca_account_record, status 置 healthy、清空 fail_count 与租约,让该号重新可被 acquire。 Args: phone (str): 手机号。 code (str): 收到的短信验证码(人工填入)。 country_code (str, optional): 区号。Defaults to "86"。 proxy_url (str, optional): 该号专属静态出口 IP(仅新号入库时写入;补货已存在号不覆盖)。 Defaults to None。 password (str, optional): 若该号另设了密码,一并存库供密码登录兜底。Defaults to None。 log (optional): 日志对象。Defaults to None(用 self.log)。 Returns: dict | None: 入库后的账号字段 dict(id/phone/uid/access_token/refresh_token/token_exp/ proxy_url/status);登录失败返回 None。 """ log = log or self.log body = {"countryCode": country_code, "phone": phone, "code": code} try: import deca_sold_core as core resp = core.do_request(log, SMS_LOGIN_PATH, body, need_auth=False, use_proxy=False, proxy_override=self._proxies(proxy_url)) data = (resp or {}).get("data") or {} access, refresh = data.get("accessToken"), data.get("refreshToken") if not access or not refresh: log.warning(f"[账号池] 短信登录未返回 token phone={phone}: {resp.get('msg') if resp else None}") return None uid = data.get("userId") new_exp = int(time.time() + int(data.get("expiresIn") or 900)) # phone 唯一键:新号 INSERT、已存在号(补货)UPDATE 回写最新 token 并复活为 healthy。 # proxy_url 用 COALESCE(旧值, 新值):补货时不覆盖已绑定的专属 IP;新号才写入。 self.pool.update_one( f"INSERT INTO {T_ACCOUNT} " "(phone, password, country_code, uid, access_token, refresh_token, token_exp, " " proxy_url, status, fail_count, last_success_at) " "VALUES (%s,%s,%s,%s,%s,%s,%s,%s,'healthy',0,NOW()) " "ON DUPLICATE KEY UPDATE " " uid=VALUES(uid), access_token=VALUES(access_token), " " refresh_token=VALUES(refresh_token), token_exp=VALUES(token_exp), " " proxy_url=COALESCE(proxy_url, VALUES(proxy_url)), " " password=COALESCE(VALUES(password), password), " " status='healthy', fail_count=0, cooldown_until=NULL, " " owner_pid=NULL, lease_until=NULL, last_success_at=NOW(), last_error=NULL", (phone, password, country_code, uid, access, refresh, new_exp, proxy_url)) row = self.pool.select_one( f"SELECT id, phone, password, country_code, uid, access_token, refresh_token, " f"token_exp, proxy_url, status FROM {T_ACCOUNT} WHERE phone=%s", (phone,)) acc = self._row_to_dict(row) if row else None log.info(f"[账号池] 短信登录入库成功 phone={phone} uid={uid} " f"id={acc['id'] if acc else '?'}(access 有效 {data.get('expiresIn')}s)") return acc except Exception as e: log.warning(f"[账号池] 短信登录异常 phone={phone}: {e}") return None def register(self, phone: str, code: str, country_code: str = "86", proxy_url: str = None, password: str = None, log=None) -> dict | None: """注册/开户一个账号并入池——得卡「登录即注册」,本方法等价于 sms_login。 得卡无独立注册接口:新手机号短信验证码登录即自动开户。故注册流程为两步(人工收验证码): 1) 先调 self.sms_send(phone) → 手机收到验证码; 2) 人工把验证码填进来,调 self.register(phone, code, proxy_url=...) 完成开户并入库。 风控提示:得卡风控强,批量开户需错开时间/设备/IP,避免一批新号被批量识别连坐(见方案文档)。 Args: phone (str): 手机号。 code (str): 短信验证码(人工填入)。 country_code (str, optional): 区号。Defaults to "86"。 proxy_url (str, optional): 该号专属静态出口 IP,入库绑定。Defaults to None。 password (str, optional): 如另设密码则一并存库(供兜底)。Defaults to None。 log (optional): 日志对象。Defaults to None。 Returns: dict | None: 入库后的账号字段 dict;失败返回 None。 """ return self.sms_login(phone, code, country_code=country_code, proxy_url=proxy_url, password=password, log=log) # ---------- 自动补货(接码平台)---------- def count_healthy(self) -> int: """统计池中当前 healthy(可用)账号数,供自动补货判断是否需要补。 Returns: int: healthy 账号数。 """ row = self.pool.select_one(f"SELECT COUNT(*) FROM {T_ACCOUNT} WHERE status='healthy'") return int(row[0]) if row and row[0] is not None else 0 def _free_proxy_urls(self) -> list: """返回 PROXY_URLS 中当前 DB 里还没被任何账号占用的空闲专属 IP(供补货分配给新号)。 Returns: list: 空闲 proxy_url 列表(PROXY_URLS 减去 deca_account_record 里已占用的 proxy_url)。 """ rows = self.pool.select_all(f"SELECT proxy_url FROM {T_ACCOUNT} WHERE proxy_url IS NOT NULL") or [] used = {r[0] for r in rows} return [u for u in PROXY_URLS if u not in used] def _release_dead_proxies(self) -> int: """把 dead 账号占用的专属 IP 释放(proxy_url 置空),供补货回收给新号(不复活、仅回收 IP)。 背景(2026/09/14):1 号 1 IP、IP 是号池规模硬上限;dead 号死了却一直占着 proxy_url, 当 PROXY_URLS 被 healthy+dead 占满时 _free_proxy_urls 返回空 → auto_replenish 撞「空闲 IP 已用尽」 永远补不进。故补货前先释放 dead 号的 IP:保留 dead 行做审计(status 不变),仅清空其 proxy_url, 让这些 IP 回到空闲池、由新号顶替。只动 status='dead' 的行,绝不碰 healthy/cooling 占用的 IP。 Returns: int: 释放(清空 proxy_url)的 dead 账号数。 """ row = self.pool.select_one( f"SELECT COUNT(*) FROM {T_ACCOUNT} WHERE status='dead' AND proxy_url IS NOT NULL") n = int(row[0]) if row and row[0] is not None else 0 if n: self.pool.update_one( f"UPDATE {T_ACCOUNT} SET proxy_url=NULL WHERE status='dead' AND proxy_url IS NOT NULL", ()) return n def register_one_via_isms(self, isms, project: dict, proxy_url: str = None, ascription: int = ISMS_ASCRIPTION, log=None) -> dict | None: """经接码平台自动注册一个得卡账号并入库(取号→得卡发码→收码→短信登录→释放号)。 全自动单号注册:接码平台租一个能收得卡短信的号 → 得卡 sms_send 发验证码 → 接码平台 poll_sms 收码 → 得卡 sms_login(登录即注册)入库 → 释放接码平台号。 Args: isms (ISmsClient): 已初始化的接码平台客户端。 project (dict): search_projects 选中的项目(含 project_id/name/token)。 proxy_url (str, optional): 分配给新号的专属静态 IP(入库绑定)。Defaults to None。 ascription (int, optional): 卡类型 2=实体卡/1=虚拟卡。Defaults to ISMS_ASCRIPTION。 log (optional): 日志对象。Defaults to None(用 self.log)。 Returns: dict | None: 入库后的账号 dict;任一步失败返回 None。 """ log = log or self.log order_id = None try: # 1) 接码平台取号 nr = isms.get_number(project_id=project.get("project_id"), project_name=project.get("name"), project_token=project.get("token"), quantity=1, ascription=ascription) data = (nr.get("data") or []) if nr.get("success") else [] if not data: log.warning(f"[补货] 接码取号失败: {nr.get('message') or nr.get('msg')}") return None phone = str(data[0].get("number")) order_id = data[0].get("orderId") or data[0].get("order_id") log.info(f"[补货] 接码取到号 {phone} order={order_id}") # 2) 得卡发验证码(走该号将绑的专属 IP,与登录同源) if not self.sms_send(phone, proxy_url=proxy_url, log=log): return None # 3) 接码平台轮询收验证码 code = isms.poll_sms(order_id, log=log) if not code: return None # 4) 得卡短信登录(=注册)入库 acc = self.sms_login(phone, code, proxy_url=proxy_url, log=log) return acc except Exception as e: log.warning(f"[补货] 单号注册异常: {e}") return None finally: # 5) 释放接码平台号(无论成败,及时释放降占用) if order_id is not None: try: isms.release_number(order_id=order_id) except Exception: pass def auto_replenish(self, target_min: int, target_max: int, free_proxy_urls: list = None, keyword: str = ISMS_PROJECT_KEYWORD, ascription: int = ISMS_ASCRIPTION, api_key: str = None, log=None) -> int: """巡检:healthy 号 < target_min 时,经接码平台自动注册补到 target_max。 供后台巡检循环定期调用。补货受两处上限约束:free_proxy_urls(空闲 IP,1 号 1 IP)与接码平台余额。 每补一个号随机停顿(REPLENISH_STAGGER_SEC)错峰,降一批新号被批量连坐识别的风险。 ⚠️ 依赖外部资源(有 API Key + IP + 得卡项目存在后才能真跑): - api_key:接码平台 Key(传参或环境变量 ISMS_API_KEY); - keyword:得卡项目关键词需 search_projects 实测确认; - free_proxy_urls:可分配的空闲专属静态 IP 列表(1 号 1 IP)。 Args: target_min (int): healthy 低于该值才触发补货。 target_max (int): 补货目标(补到 healthy 达此值)。 free_proxy_urls (list, optional): 可分配给新号的空闲专属 IP 列表。Defaults to None。 keyword (str, optional): 接码平台得卡项目搜索关键词。Defaults to ISMS_PROJECT_KEYWORD。 ascription (int, optional): 卡类型 2=实体卡/1=虚拟卡。Defaults to ISMS_ASCRIPTION。 api_key (str, optional): 接码平台 Key;None 时用配置常量 ISMS_API_KEY。Defaults to None。 log (optional): 日志对象。Defaults to None(用 self.log)。 Returns: int: 本轮成功补充入库的账号数。 """ import random log = log or self.log cur = self.count_healthy() if cur >= target_min: return 0 log.info(f"[补货] healthy={cur} < {target_min},开始自动补货至 {target_max}") try: from isms_client import ISmsClient isms = ISmsClient(api_key=api_key or ISMS_API_KEY) # 默认用配置里的 Key,传参可覆盖 except Exception as e: log.error(f"[补货] 接码平台客户端初始化失败(检查 ISMS_API_KEY): {e}") return 0 # 搜得卡项目 pr = isms.search_projects(keyword) projects = (pr.get("data") or []) if pr.get("success") else [] if not projects: log.error(f"[补货] 接码平台未搜到项目 '{keyword}'——请确认关键词或该平台是否有得卡项目") return 0 project = projects[0] # 补货前先释放 dead 号占用的专属 IP(不复活、仅回收 IP),否则 IP 被 healthy+dead 占满时永远补不进 released = self._release_dead_proxies() if released: log.info(f"[补货] 已释放 {released} 个 dead 账号占用的专属 IP 供新号回收") if free_proxy_urls is None: free_proxy_urls = self._free_proxy_urls() # 默认用配置 PROXY_URLS 里 DB 未占用的空闲 IP free = list(free_proxy_urls) added = 0 fail_streak = 0 # 连续失败计数:防接码没号/异常时死循环空转烧钱 while self.count_healthy() < target_max: proxy = free.pop(0) if free else None # 1 号 1 IP:无空闲 IP 则停(IP 是硬上限) if proxy is None: log.warning("[补货] 空闲 IP 已用尽,停止补货(IP 是号池规模硬上限)") break acc = self.register_one_via_isms(isms, project, proxy_url=proxy, ascription=ascription, log=log) if acc: added += 1 fail_streak = 0 log.info(f"[补货] 已补 {added} 个(最新 phone={acc['phone']} proxy={proxy})") else: if proxy is not None: free.insert(0, proxy) # 本号失败,IP 归还继续给下一个用 fail_streak += 1 if fail_streak >= 5: # 连续 5 次失败:接码可能没号/余额不足/异常,停止避免空转 log.error("[补货] 连续 5 次注册失败,停止补货(检查接码余额/项目号源/网络)") break log.warning(f"[补货] 本次注册失败(连续 {fail_streak} 次),稍后重试") time.sleep(random.randint(*REPLENISH_STAGGER_SEC)) # 错峰,避免一批号被连坐 log.info(f"[补货] 本轮补充完成 {added} 个,当前 healthy={self.count_healthy()}") return added # ---------- 内部 ---------- def _row_to_dict(self, row) -> dict: """把 acquire 回读的行元组转成账号 dict(列顺序与 acquire 的 SELECT 对齐)。 Args: row (tuple): SELECT 出的一行。 Returns: dict: 账号字段 dict。 """ keys = ["id", "phone", "password", "country_code", "uid", "access_token", "refresh_token", "token_exp", "proxy_url", "status"] return dict(zip(keys, row))