account_pool.py 40 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696
  1. # -*- coding: utf-8 -*-
  2. # Author : Charley
  3. # Python : 3.12.10
  4. # Date : 2026/09/08
  5. """得卡 DECA · 账号池(框架版,2026/09/08 搭)。
  6. 把原来集中在单账号 token.json 的登录态,改为「多账号入池、按任务取号、异常自动切号、冷却复用」,
  7. 把带 token 请求压力平摊到多账号 + 各账号专属静态 IP,根治 9/7「单账号被风控击穿致整链停摆」。
  8. 方案与请求量测算见 HANDOFF.md 及 docs/账号池方案_得卡DECA_20260908.md。
  9. 数据模型:一张 deca_account_record 表(DDL 见同目录 deca_account_ddl.sql)+ 本模块。
  10. 每个账号自带独立 access/refresh/token_exp + 专属 proxy_url(1 号 1 IP 永久绑定)+ 健康状态。
  11. 线上生命周期原则(重要):
  12. · 线上**只做 refresh 续期,绝不自动跑密码登录**——密码登录会撞阿里云滑块、无人值守过不去。
  13. 续期失败即判 dead、摘池、企微通知,自动切下一个 healthy 号,采集不中断。
  14. · 补货靠接码平台自动补货:healthy 数不足时 auto_replenish 经接码平台租号 + 短信登录(=注册)入库 +
  15. 绑定空闲专属 IP(见 register_one_via_isms / spiders/replenish_accounts.py),无需人工灌 token。
  16. 取号粒度:进程级独占——每个常驻任务启动时 acquire 一个号、独占整个运行期(保持账号↔IP 稳定降风控),
  17. 仅在该号挂了才 switch 换号。购买记录按商家分片时,每个分片进程各 acquire 一个号(带 owner_tag 区分)。
  18. 账号供给(2026/09/09 落地):得卡无独立注册接口——短信验证码登录「登录即注册」,故:
  19. - sms_send() / sms_login() 发验证码 + 短信登录入库(登录即注册;新号 INSERT、旧号补货 UPDATE)
  20. - register() 等价于 sms_login(两步:先 sms_send 收码,再带 code 调用),离线人工补货/开户主路径
  21. - _login() 密码登录(撞阿里云滑块),仅极端兜底,线上默认不调;短信登录不撞滑块、优先用它
  22. 接入采集脚本(把 core.do_request 的 need_auth 取号改为走本池)属后续 step3,未在本框架内改 core。
  23. """
  24. import os
  25. import time
  26. from datetime import datetime, timedelta
  27. from loguru import logger
  28. # ==================== 配置 ====================
  29. # 全部业务/敏感配置集中在 settings.py(接口路径 / 接码 Key / 代理账密 IP 池 / 阈值),改配置只动那处。
  30. from settings import * # noqa: F401,F403 (T_ACCOUNT/LEASE_TTL_SEC/DEAD_CODES/ISMS_*/PROXY_*/*_PATH 等)
  31. class AccountPool:
  32. """账号池:从 deca_account_record 取号、续期、异常切号、状态维护。
  33. 单个常驻任务进程持有一个实例:启动 acquire 一个号独占,运行期用 ensure_access 保证 token 有效,
  34. 请求失败调 report_failure(达阈值 cooling/dead)→ switch 换号;结束 release 归还。
  35. """
  36. def __init__(self, pool, log=None, task_tag: str = ""):
  37. """初始化账号池。
  38. Args:
  39. pool (MySQLConnectionPool): MySQL 连接池。
  40. log (optional): 日志对象。Defaults to None(回落 loguru.logger)。
  41. task_tag (str, optional): 本进程任务标签,写入 owner_tag 便于排障与独占区分
  42. (如 "buy_record:881226408" / "team" / "onsale_alert")。Defaults to ""。
  43. """
  44. self.pool = pool
  45. self.log = log or logger
  46. self.task_tag = task_tag
  47. self.pid = os.getpid()
  48. self.account = None # 当前持有的账号行 dict;None 表示未持号
  49. # ---------- 取号 / 归还 / 切号 ----------
  50. def acquire(self) -> dict | None:
  51. """从池中租一个 healthy 且未被独占(或租约已过期)的账号,打上本进程租约并返回。
  52. 原子抢占:一条 UPDATE 按 last_used_at 升序挑一个可租的号打租约(owner_pid/lease_until),
  53. 再按本进程 pid + task_tag 回读刚租到的行。取号成功后写入 self.account。
  54. Returns:
  55. dict | None: 账号行 dict(含 id/phone/access_token/refresh_token/token_exp/proxy_url 等);
  56. 池中无可用 healthy 号时返回 None(调用方应告警/退避)。
  57. """
  58. lease_until = datetime.now() + timedelta(seconds=LEASE_TTL_SEC)
  59. # 挑一个 healthy 且(未被占 或 租约已过期)的号,按最久未用优先做负载均衡,打上本进程租约
  60. self.pool.update_one(
  61. f"UPDATE {T_ACCOUNT} SET owner_tag=%s, owner_pid=%s, lease_until=%s, last_used_at=NOW() "
  62. "WHERE status='healthy' AND (owner_pid IS NULL OR lease_until IS NULL OR lease_until < NOW()) "
  63. "ORDER BY last_used_at ASC LIMIT 1", # NULL 升序排最前:从未用过的号优先取,做负载均衡
  64. (self.task_tag, self.pid, lease_until.strftime("%Y-%m-%d %H:%M:%S")))
  65. # 回读本进程刚租到的号(owner_pid=自身 pid 唯一,别的进程不会用我的 pid)
  66. row = self.pool.select_one(
  67. f"SELECT id, phone, password, country_code, uid, access_token, refresh_token, token_exp, "
  68. f"proxy_url, status FROM {T_ACCOUNT} WHERE owner_pid=%s AND owner_tag=%s AND status='healthy' "
  69. f"ORDER BY lease_until DESC LIMIT 1",
  70. (self.pid, self.task_tag))
  71. if not row:
  72. self.log.error(f"[账号池] 无可用 healthy 账号(task={self.task_tag})")
  73. self.account = None
  74. return None
  75. self.account = self._row_to_dict(row)
  76. self.log.info(f"[账号池] 取号成功 id={self.account['id']} phone={self.account['phone']} "
  77. f"proxy={self.account.get('proxy_url')} task={self.task_tag}")
  78. return self.account
  79. def acquire_random(self, lease_sec: int = LEASE_TTL_SEC) -> dict | None:
  80. """从 healthy 池随机挑一个未被占的账号,打上租约并返回(每批随机取号用)。
  81. 与 acquire(进程级独占、按最久未用负载均衡)的区别:本方法按 RAND() 随机挑号,
  82. 供「每批随机」——每 N 个请求换一个随机号,把带 token 请求量摊平到所有号、降单号风控面。
  83. 租约默认用 LEASE_TTL_SEC(较长,覆盖持号整个使用期):一个进程持号直到换号/退出,
  84. 期间租约不过期,别的进程不会抢走——避免「持号时长 > 租约 → 被抢 → 两进程同号续期
  85. 踩废一次性轮换的 refreshToken → 误判 10010 dead」(2026/09/09 修:原 300s 过短,慢/稀疏
  86. 任务如 team/onsale_alert 会在换号前租约过期)。进程崩溃后租约到期自动可被再租,不死占。
  87. 并发抢占:先 RAND() 选一个候选 id,再带「仍空闲」条件 UPDATE 认领;认领失败(被别的进程抢走)
  88. 重试几次换一个候选。认领成功后回读该行写入 self.account。
  89. Args:
  90. lease_sec (int, optional): 租约时长(秒),覆盖持号整个使用期。Defaults to LEASE_TTL_SEC(6h)。
  91. Returns:
  92. dict | None: 账号行 dict;池中无可用 healthy 号或多次抢占失败返回 None。
  93. """
  94. lease_until = (datetime.now() + timedelta(seconds=lease_sec)).strftime("%Y-%m-%d %H:%M:%S")
  95. for _ in range(5): # 抢占重试:候选被别的进程抢走时换一个
  96. cand = self.pool.select_one(
  97. f"SELECT id FROM {T_ACCOUNT} WHERE status='healthy' "
  98. "AND (owner_pid IS NULL OR lease_until IS NULL OR lease_until < NOW()) "
  99. "ORDER BY RAND() LIMIT 1")
  100. if not cand:
  101. self.log.error(f"[账号池] 无可用 healthy 账号(随机取号,task={self.task_tag})")
  102. self.account = None
  103. return None
  104. cid = cand[0]
  105. # 带「仍空闲」条件认领:若被别的进程抢先占了,本 UPDATE 命中 0 行,回读会拿不到本 pid → 重试
  106. self.pool.update_one(
  107. f"UPDATE {T_ACCOUNT} SET owner_tag=%s, owner_pid=%s, lease_until=%s, last_used_at=NOW() "
  108. "WHERE id=%s AND status='healthy' "
  109. "AND (owner_pid IS NULL OR lease_until IS NULL OR lease_until < NOW())",
  110. (self.task_tag, self.pid, lease_until, cid))
  111. row = self.pool.select_one(
  112. f"SELECT id, phone, password, country_code, uid, access_token, refresh_token, "
  113. f"token_exp, proxy_url, status FROM {T_ACCOUNT} WHERE id=%s AND owner_pid=%s",
  114. (cid, self.pid))
  115. if row:
  116. self.account = self._row_to_dict(row)
  117. self.log.info(f"[账号池] 随机取号 id={self.account['id']} phone={self.account['phone']} "
  118. f"proxy={self.account.get('proxy_url')} lease={lease_sec}s task={self.task_tag}")
  119. return self.account
  120. self.log.warning(f"[账号池] 随机取号多次抢占失败(并发激烈,task={self.task_tag})")
  121. self.account = None
  122. return None
  123. def release(self):
  124. """归还当前账号(清空租约 owner_pid/lease_until),进程正常退出时调用。"""
  125. if not self.account:
  126. return
  127. self.pool.update_one(
  128. f"UPDATE {T_ACCOUNT} SET owner_pid=NULL, lease_until=NULL WHERE id=%s",
  129. (self.account["id"],))
  130. self.log.info(f"[账号池] 归还账号 id={self.account['id']}")
  131. self.account = None
  132. def switch(self) -> dict | None:
  133. """当前号不可用时换号:释放当前号的租约(不改其 status,status 由 report_* 决定)后重新 acquire。
  134. Returns:
  135. dict | None: 新账号行 dict;池中无其他可用号时返回 None。
  136. """
  137. old = self.account
  138. if old:
  139. # 只解租约,不动 status——是否 cooling/dead 由 report_failure 已经写好;healthy 号只是让给别人
  140. self.pool.update_one(
  141. f"UPDATE {T_ACCOUNT} SET owner_pid=NULL, lease_until=NULL WHERE id=%s AND owner_pid=%s",
  142. (old["id"], self.pid))
  143. self.account = None
  144. new = self.acquire()
  145. if new:
  146. self.log.warning(f"[账号池] 切号:{old['phone'] if old else '-'} → {new['phone']}")
  147. else:
  148. self.log.error("[账号池] 切号失败:池中已无可用 healthy 号,请尽快离线补货")
  149. return new
  150. # ---------- token 续期 ----------
  151. @staticmethod
  152. def _proxies(proxy_url) -> dict | None:
  153. """把该号 proxy_url 组装成 requests proxies(让续期/短信登录也走该号专属 IP)。
  154. Args:
  155. proxy_url: 该号绑定的代理 URL;None/空表示直连。
  156. Returns:
  157. dict | None: {"http":.., "https":..};proxy_url 为空返回 None。
  158. """
  159. return {"http": proxy_url, "https": proxy_url} if proxy_url else None
  160. def ensure_access(self, log=None) -> str | None:
  161. """确保当前账号 access 有效:未临期直接返回;临期用该号 refresh 续期并回写 DB。
  162. 线上主路径:只 refresh、不自动登录。续期失败即判该号 dead 并切号(见 report_failure)。
  163. Args:
  164. log (optional): 日志对象。Defaults to None(用 self.log)。
  165. Returns:
  166. str | None: 有效 access token;无持号/续期失败返回 None(调用方据此 switch 换号)。
  167. """
  168. log = log or self.log
  169. acc = self.account
  170. if not acc:
  171. log.error("[账号池] ensure_access 时未持号")
  172. return None
  173. # 未临期:直接用缓存 access
  174. if acc.get("access_token") and acc.get("token_exp") and time.time() < acc["token_exp"] - ACCESS_SKEW_SEC:
  175. return acc["access_token"]
  176. # 临期/缺失:用该号 refresh 续期(need_auth=False,不依赖全局 token)。走该号专属 proxy_url,
  177. # 让续期这类登录态操作也从该号绑定的 IP 发出(不暴露本机 IP,防 20 号在登录态层被同源关联)。
  178. if not acc.get("refresh_token"):
  179. log.warning(f"[账号池] id={acc['id']} 无 refresh_token,无法续期")
  180. self.report_failure("no_refresh_token", dead=True)
  181. return None
  182. try:
  183. import deca_sold_core as core # 惰性 import:避免 import 本模块即触发 core 副作用
  184. resp = core.do_request(log, REFRESH_PATH, {"refreshToken": acc["refresh_token"]},
  185. need_auth=False, use_proxy=False,
  186. proxy_override=self._proxies(acc.get("proxy_url")))
  187. code = (resp or {}).get("code")
  188. msg = (resp or {}).get("msg")
  189. data = (resp or {}).get("data") or {}
  190. access = data.get("accessToken")
  191. if not access:
  192. if code in DEAD_CODES:
  193. # 实测 10010=refresh 失效(被轮换/过期/长期未登录)→ 确定判 dead 摘池(线上不自动登录)
  194. log.warning(f"[账号池] id={acc['id']} 续期判死 code={code} msg={msg}")
  195. self.report_failure(f"refresh_dead:code={code}", dead=True)
  196. else:
  197. # 其他未知失败(限流/网络/临时):软失败累计,达 MAX_FAIL 才 cooling,不急着判死
  198. log.warning(f"[账号池] id={acc['id']} 续期失败(非判死码) code={code} msg={msg}")
  199. self.report_failure(f"refresh_fail:code={code}", dead=False)
  200. return None
  201. expires_in = int(data.get("expiresIn") or 900)
  202. new_refresh = data.get("refreshToken") or acc["refresh_token"] # 一次性轮换:有新的就换
  203. new_exp = int(time.time() + expires_in)
  204. # 回写 DB 该号最新 token
  205. self.pool.update_one(
  206. f"UPDATE {T_ACCOUNT} SET access_token=%s, refresh_token=%s, token_exp=%s, "
  207. "last_success_at=NOW(), fail_count=0 WHERE id=%s",
  208. (access, new_refresh, new_exp, acc["id"]))
  209. acc["access_token"], acc["refresh_token"], acc["token_exp"] = access, new_refresh, new_exp
  210. log.info(f"[账号池] id={acc['id']} 续期成功,有效 {expires_in}s")
  211. return access
  212. except Exception as e:
  213. # 网络/超时等异常(非明确判死码)→ 软失败累计,不永久判死(网络问题 ≠ 账号死)
  214. log.warning(f"[账号池] id={acc['id']} 续期异常(按软失败处理): {e}")
  215. self.report_failure(f"refresh_exc:{e}", dead=False)
  216. return None
  217. # ---------- 状态机 ----------
  218. def report_success(self):
  219. """标记当前账号一次成功(清零 fail_count、刷新 last_success_at)。"""
  220. if not self.account:
  221. return
  222. self.pool.update_one(
  223. f"UPDATE {T_ACCOUNT} SET fail_count=0, last_success_at=NOW(), last_error=NULL WHERE id=%s",
  224. (self.account["id"],))
  225. def report_failure(self, error: str = "", dead: bool = False):
  226. """记录当前账号一次失败:dead=True 直接判死摘池;否则累计失败,达 MAX_FAIL 转 cooling。
  227. Args:
  228. error (str, optional): 错误信息,写入 last_error 供排障。Defaults to ""。
  229. dead (bool, optional): True=登录态死了(续期失败等)直接置 dead 并企微通知人工。
  230. False=软失败(限流/偶发),累计到 MAX_FAIL 才转 cooling。Defaults to False。
  231. """
  232. if not self.account:
  233. return
  234. aid = self.account["id"]
  235. err = (error or "")[:255]
  236. if dead:
  237. self.pool.update_one(
  238. f"UPDATE {T_ACCOUNT} SET status='dead', owner_pid=NULL, lease_until=NULL, "
  239. "last_error=%s WHERE id=%s", (err, aid))
  240. self.log.error(f"[账号池] id={aid} 判 dead 摘池({err})→ 需离线补货")
  241. # TODO(step3): 企微通知人工「账号 {phone} dead,请离线补货」,复用 auto_send_wx_msg。
  242. self.account = None
  243. return
  244. # 软失败累计
  245. self.pool.update_one(
  246. f"UPDATE {T_ACCOUNT} SET fail_count=fail_count+1, last_error=%s WHERE id=%s", (err, aid))
  247. row = self.pool.select_one(f"SELECT fail_count FROM {T_ACCOUNT} WHERE id=%s", (aid,))
  248. fc = int(row[0]) if row and row[0] is not None else 0
  249. if fc >= MAX_FAIL:
  250. cd = (datetime.now() + timedelta(seconds=COOLDOWN_SEC)).strftime("%Y-%m-%d %H:%M:%S")
  251. self.pool.update_one(
  252. f"UPDATE {T_ACCOUNT} SET status='cooling', cooldown_until=%s, owner_pid=NULL, "
  253. "lease_until=NULL WHERE id=%s", (cd, aid))
  254. self.log.warning(f"[账号池] id={aid} 连续失败 {fc} 次转 cooling,冷却到 {cd}")
  255. self.account = None
  256. def revive_cooling(self) -> int:
  257. """巡检:把冷却到期的 cooling 账号恢复为 healthy 重新入池(dead 不自动恢复,需人工)。
  258. 由一个后台巡检循环定期调用(如每几分钟一次)。
  259. Returns:
  260. int: 本次恢复的账号数。
  261. """
  262. self.pool.update_one(
  263. f"UPDATE {T_ACCOUNT} SET status='healthy', fail_count=0, cooldown_until=NULL "
  264. "WHERE status='cooling' AND cooldown_until IS NOT NULL AND cooldown_until < NOW()",
  265. ())
  266. # 恢复条数取决于底层驱动是否回报 affected rows;此处仅执行,数量另查(保持接口简单)
  267. return 0
  268. def keepalive_idle(self, idle_hours: int = 6, log=None) -> tuple[int, int]:
  269. """token 保活巡检:对闲置超过 idle_hours 的 healthy 号主动刷新一次,避免 refresh_token 闲到过期被判死。
  270. 背景(2026/09/15):access token 仅 15 分钟、refresh_token 一次性轮换;号若长时间没被业务请求抽到,
  271. refresh_token 会自然过期,下次一用即 10010 判死(实测闲置 30~77h 死)。故后台每小时扫一遍:
  272. 把「healthy 且 last_success_at 超过 idle_hours 没成功 且 当前未被占用」的号逐个强制刷新——
  273. 刷成功则 token 轮换保温热;刷回 10010 则本就已死、判 dead 交自动补货补。
  274. 并发安全:每个号先带「仍空闲」条件抢租约(claim)、抢到才刷、刷完释放,避免与业务进程并发刷同一号
  275. 踩废其一次性轮换的 refreshToken。
  276. Args:
  277. idle_hours (int, optional): 闲置阈值(小时),last_success_at 早于此才刷。Defaults to 6。
  278. log (optional): 日志对象。Defaults to None(用 self.log)。
  279. Returns:
  280. tuple[int, int]: (温热成功数, 判死数)。
  281. """
  282. log = log or self.log
  283. rows = self.pool.select_all(
  284. f"SELECT id FROM {T_ACCOUNT} WHERE status='healthy' "
  285. "AND (last_success_at IS NULL OR last_success_at < NOW() - INTERVAL %s HOUR) "
  286. "AND (owner_pid IS NULL OR lease_until IS NULL OR lease_until < NOW())",
  287. (idle_hours,)) or []
  288. warmed, died = 0, 0
  289. for (aid,) in rows:
  290. lease_until = (datetime.now() + timedelta(seconds=120)).strftime("%Y-%m-%d %H:%M:%S") # 短租,够刷新用
  291. # 带「仍空闲」条件抢占:抢到才刷(防与业务进程并发刷同一号)
  292. self.pool.update_one(
  293. f"UPDATE {T_ACCOUNT} SET owner_tag=%s, owner_pid=%s, lease_until=%s "
  294. "WHERE id=%s AND status='healthy' "
  295. "AND (owner_pid IS NULL OR lease_until IS NULL OR lease_until < NOW())",
  296. (self.task_tag, self.pid, lease_until, aid))
  297. row = self.pool.select_one(
  298. f"SELECT id, phone, password, country_code, uid, access_token, refresh_token, "
  299. f"token_exp, proxy_url, status FROM {T_ACCOUNT} WHERE id=%s AND owner_pid=%s AND status='healthy'",
  300. (aid, self.pid))
  301. if not row:
  302. continue # 抢占失败(被业务进程抢走)或状态已变,跳过
  303. self.account = self._row_to_dict(row)
  304. self.account["token_exp"] = 0 # 置 0 强制走刷新分支,确保轮换 refresh_token(而非命中未临期缓存)
  305. access = self.ensure_access(log) # 内部:刷成功回写 token;10010 则 report_failure(dead) 已摘池
  306. if access:
  307. warmed += 1
  308. # 仍 healthy:解租约归还(ensure_access 判死时已清 self.account+解租约,不会走到这)
  309. if self.account is not None:
  310. self.pool.update_one(
  311. f"UPDATE {T_ACCOUNT} SET owner_pid=NULL, lease_until=NULL WHERE id=%s AND owner_pid=%s",
  312. (self.account["id"], self.pid))
  313. else:
  314. died += 1
  315. self.account = None
  316. if warmed or died:
  317. log.info(f"[保活] 闲置>{idle_hours}h 刷新:温热 {warmed} 个,判死 {died} 个")
  318. return warmed, died
  319. # ---------- 账号供给(依赖明天抓的接口,占位)----------
  320. def _login(self, log=None) -> bool:
  321. """密码登录换取全新一套 token(照 deca_sold_core.login 的接口写实)。
  322. ⚠️ 线上默认不调用——密码登录会撞阿里云滑块,无人值守过不去(9/7 就死在这)。仅供
  323. 「离线补货 / 手动兜底」场景显式调用。线上续期失败一律走 report_failure(dead=True)。
  324. Args:
  325. log (optional): 日志对象。Defaults to None。
  326. Returns:
  327. bool: 登录成功并回写 token 返回 True;失败返回 False。
  328. """
  329. log = log or self.log
  330. acc = self.account
  331. if not acc:
  332. return False
  333. body = {"countryCode": acc.get("country_code") or "86",
  334. "password": acc.get("password"), "phone": acc.get("phone")}
  335. try:
  336. import deca_sold_core as core
  337. resp = core.do_request(log, LOGIN_PATH, body, need_auth=False, use_proxy=False)
  338. data = (resp or {}).get("data") or {}
  339. access, refresh = data.get("accessToken"), data.get("refreshToken")
  340. if not access or not refresh:
  341. log.warning(f"[账号池] id={acc['id']} 密码登录未返回 token: {resp.get('msg') if resp else None}")
  342. return False
  343. new_exp = int(time.time() + int(data.get("expiresIn") or 900))
  344. self.pool.update_one(
  345. f"UPDATE {T_ACCOUNT} SET access_token=%s, refresh_token=%s, token_exp=%s, "
  346. "status='healthy', fail_count=0, last_success_at=NOW() WHERE id=%s",
  347. (access, refresh, new_exp, acc["id"]))
  348. acc.update({"access_token": access, "refresh_token": refresh, "token_exp": new_exp})
  349. log.info(f"[账号池] id={acc['id']} 密码登录成功(离线补货)")
  350. return True
  351. except Exception as e:
  352. log.warning(f"[账号池] id={acc['id']} 密码登录异常: {e}")
  353. return False
  354. def sms_send(self, phone: str, country_code: str = "86", proxy_url: str = None, log=None) -> bool:
  355. """向指定手机号发送登录短信验证码(登录/注册第一步)。
  356. need_auth=False,不依赖任何登录态。用于「补货 / 新号开户」流程第一步:调用后手机收到验证码,
  357. 再调 sms_login。走该号将绑的专属 proxy_url,保证发码也从该号 IP 发出(与后续登录同源)。
  358. Args:
  359. phone (str): 目标手机号。
  360. country_code (str, optional): 区号。Defaults to "86"。
  361. proxy_url (str, optional): 该号专属 IP(走它发码,与登录同源)。Defaults to None(直连)。
  362. log (optional): 日志对象。Defaults to None(用 self.log)。
  363. Returns:
  364. bool: code==0 视为发送成功返回 True;否则 False。
  365. """
  366. log = log or self.log
  367. body = {"countryCode": country_code, "phone": phone}
  368. try:
  369. import deca_sold_core as core # 惰性 import:避免 import 本模块即触发 core 副作用
  370. resp = core.do_request(log, SMS_SEND_PATH, body, need_auth=False, use_proxy=False,
  371. proxy_override=self._proxies(proxy_url))
  372. if resp and resp.get("code") == 0:
  373. log.info(f"[账号池] 验证码已发送 phone={phone}")
  374. return True
  375. log.warning(f"[账号池] 发验证码失败 phone={phone}: {resp.get('msg') if resp else None}")
  376. return False
  377. except Exception as e:
  378. log.warning(f"[账号池] 发验证码异常 phone={phone}: {e}")
  379. return False
  380. def sms_login(self, phone: str, code: str, country_code: str = "86",
  381. proxy_url: str = None, password: str = None, log=None) -> dict | None:
  382. """短信验证码登录(登录即注册)并入库为 healthy 账号(新号 INSERT / 已存在号 UPDATE token)。
  383. 得卡无独立注册接口:新手机号首次短信登录即自动开户,故本方法同时承担「注册新号」与
  384. 「离线补货已死账号」两职。成功后把 access/refresh/token_exp/uid 落库 deca_account_record,
  385. status 置 healthy、清空 fail_count 与租约,让该号重新可被 acquire。
  386. Args:
  387. phone (str): 手机号。
  388. code (str): 收到的短信验证码(人工填入)。
  389. country_code (str, optional): 区号。Defaults to "86"。
  390. proxy_url (str, optional): 该号专属静态出口 IP(仅新号入库时写入;补货已存在号不覆盖)。
  391. Defaults to None。
  392. password (str, optional): 若该号另设了密码,一并存库供密码登录兜底。Defaults to None。
  393. log (optional): 日志对象。Defaults to None(用 self.log)。
  394. Returns:
  395. dict | None: 入库后的账号字段 dict(id/phone/uid/access_token/refresh_token/token_exp/
  396. proxy_url/status);登录失败返回 None。
  397. """
  398. log = log or self.log
  399. body = {"countryCode": country_code, "phone": phone, "code": code}
  400. try:
  401. import deca_sold_core as core
  402. resp = core.do_request(log, SMS_LOGIN_PATH, body, need_auth=False, use_proxy=False,
  403. proxy_override=self._proxies(proxy_url))
  404. data = (resp or {}).get("data") or {}
  405. access, refresh = data.get("accessToken"), data.get("refreshToken")
  406. if not access or not refresh:
  407. log.warning(f"[账号池] 短信登录未返回 token phone={phone}: {resp.get('msg') if resp else None}")
  408. return None
  409. uid = data.get("userId")
  410. new_exp = int(time.time() + int(data.get("expiresIn") or 900))
  411. # phone 唯一键:新号 INSERT、已存在号(补货)UPDATE 回写最新 token 并复活为 healthy。
  412. # proxy_url 用 COALESCE(旧值, 新值):补货时不覆盖已绑定的专属 IP;新号才写入。
  413. self.pool.update_one(
  414. f"INSERT INTO {T_ACCOUNT} "
  415. "(phone, password, country_code, uid, access_token, refresh_token, token_exp, "
  416. " proxy_url, status, fail_count, last_success_at) "
  417. "VALUES (%s,%s,%s,%s,%s,%s,%s,%s,'healthy',0,NOW()) "
  418. "ON DUPLICATE KEY UPDATE "
  419. " uid=VALUES(uid), access_token=VALUES(access_token), "
  420. " refresh_token=VALUES(refresh_token), token_exp=VALUES(token_exp), "
  421. " proxy_url=COALESCE(proxy_url, VALUES(proxy_url)), "
  422. " password=COALESCE(VALUES(password), password), "
  423. " status='healthy', fail_count=0, cooldown_until=NULL, "
  424. " owner_pid=NULL, lease_until=NULL, last_success_at=NOW(), last_error=NULL",
  425. (phone, password, country_code, uid, access, refresh, new_exp, proxy_url))
  426. row = self.pool.select_one(
  427. f"SELECT id, phone, password, country_code, uid, access_token, refresh_token, "
  428. f"token_exp, proxy_url, status FROM {T_ACCOUNT} WHERE phone=%s", (phone,))
  429. acc = self._row_to_dict(row) if row else None
  430. log.info(f"[账号池] 短信登录入库成功 phone={phone} uid={uid} "
  431. f"id={acc['id'] if acc else '?'}(access 有效 {data.get('expiresIn')}s)")
  432. return acc
  433. except Exception as e:
  434. log.warning(f"[账号池] 短信登录异常 phone={phone}: {e}")
  435. return None
  436. def register(self, phone: str, code: str, country_code: str = "86",
  437. proxy_url: str = None, password: str = None, log=None) -> dict | None:
  438. """注册/开户一个账号并入池——得卡「登录即注册」,本方法等价于 sms_login。
  439. 得卡无独立注册接口:新手机号短信验证码登录即自动开户。故注册流程为两步(人工收验证码):
  440. 1) 先调 self.sms_send(phone) → 手机收到验证码;
  441. 2) 人工把验证码填进来,调 self.register(phone, code, proxy_url=...) 完成开户并入库。
  442. 风控提示:得卡风控强,批量开户需错开时间/设备/IP,避免一批新号被批量识别连坐(见方案文档)。
  443. Args:
  444. phone (str): 手机号。
  445. code (str): 短信验证码(人工填入)。
  446. country_code (str, optional): 区号。Defaults to "86"。
  447. proxy_url (str, optional): 该号专属静态出口 IP,入库绑定。Defaults to None。
  448. password (str, optional): 如另设密码则一并存库(供兜底)。Defaults to None。
  449. log (optional): 日志对象。Defaults to None。
  450. Returns:
  451. dict | None: 入库后的账号字段 dict;失败返回 None。
  452. """
  453. return self.sms_login(phone, code, country_code=country_code,
  454. proxy_url=proxy_url, password=password, log=log)
  455. # ---------- 自动补货(接码平台)----------
  456. def count_healthy(self) -> int:
  457. """统计池中当前 healthy(可用)账号数,供自动补货判断是否需要补。
  458. Returns:
  459. int: healthy 账号数。
  460. """
  461. row = self.pool.select_one(f"SELECT COUNT(*) FROM {T_ACCOUNT} WHERE status='healthy'")
  462. return int(row[0]) if row and row[0] is not None else 0
  463. def _free_proxy_urls(self) -> list:
  464. """返回 PROXY_URLS 中当前 DB 里还没被任何账号占用的空闲专属 IP(供补货分配给新号)。
  465. Returns:
  466. list: 空闲 proxy_url 列表(PROXY_URLS 减去 deca_account_record 里已占用的 proxy_url)。
  467. """
  468. rows = self.pool.select_all(f"SELECT proxy_url FROM {T_ACCOUNT} WHERE proxy_url IS NOT NULL") or []
  469. used = {r[0] for r in rows}
  470. return [u for u in PROXY_URLS if u not in used]
  471. def _release_dead_proxies(self) -> int:
  472. """把 dead 账号占用的专属 IP 释放(proxy_url 置空),供补货回收给新号(不复活、仅回收 IP)。
  473. 背景(2026/09/14):1 号 1 IP、IP 是号池规模硬上限;dead 号死了却一直占着 proxy_url,
  474. 当 PROXY_URLS 被 healthy+dead 占满时 _free_proxy_urls 返回空 → auto_replenish 撞「空闲 IP 已用尽」
  475. 永远补不进。故补货前先释放 dead 号的 IP:保留 dead 行做审计(status 不变),仅清空其 proxy_url,
  476. 让这些 IP 回到空闲池、由新号顶替。只动 status='dead' 的行,绝不碰 healthy/cooling 占用的 IP。
  477. Returns:
  478. int: 释放(清空 proxy_url)的 dead 账号数。
  479. """
  480. row = self.pool.select_one(
  481. f"SELECT COUNT(*) FROM {T_ACCOUNT} WHERE status='dead' AND proxy_url IS NOT NULL")
  482. n = int(row[0]) if row and row[0] is not None else 0
  483. if n:
  484. self.pool.update_one(
  485. f"UPDATE {T_ACCOUNT} SET proxy_url=NULL WHERE status='dead' AND proxy_url IS NOT NULL", ())
  486. return n
  487. def register_one_via_isms(self, isms, project: dict, proxy_url: str = None,
  488. ascription: int = ISMS_ASCRIPTION, log=None) -> dict | None:
  489. """经接码平台自动注册一个得卡账号并入库(取号→得卡发码→收码→短信登录→释放号)。
  490. 全自动单号注册:接码平台租一个能收得卡短信的号 → 得卡 sms_send 发验证码 →
  491. 接码平台 poll_sms 收码 → 得卡 sms_login(登录即注册)入库 → 释放接码平台号。
  492. Args:
  493. isms (ISmsClient): 已初始化的接码平台客户端。
  494. project (dict): search_projects 选中的项目(含 project_id/name/token)。
  495. proxy_url (str, optional): 分配给新号的专属静态 IP(入库绑定)。Defaults to None。
  496. ascription (int, optional): 卡类型 2=实体卡/1=虚拟卡。Defaults to ISMS_ASCRIPTION。
  497. log (optional): 日志对象。Defaults to None(用 self.log)。
  498. Returns:
  499. dict | None: 入库后的账号 dict;任一步失败返回 None。
  500. """
  501. log = log or self.log
  502. order_id = None
  503. try:
  504. # 1) 接码平台取号
  505. nr = isms.get_number(project_id=project.get("project_id"), project_name=project.get("name"),
  506. project_token=project.get("token"), quantity=1, ascription=ascription)
  507. data = (nr.get("data") or []) if nr.get("success") else []
  508. if not data:
  509. log.warning(f"[补货] 接码取号失败: {nr.get('message') or nr.get('msg')}")
  510. return None
  511. phone = str(data[0].get("number"))
  512. order_id = data[0].get("orderId") or data[0].get("order_id")
  513. log.info(f"[补货] 接码取到号 {phone} order={order_id}")
  514. # 2) 得卡发验证码(走该号将绑的专属 IP,与登录同源)
  515. if not self.sms_send(phone, proxy_url=proxy_url, log=log):
  516. return None
  517. # 3) 接码平台轮询收验证码
  518. code = isms.poll_sms(order_id, log=log)
  519. if not code:
  520. return None
  521. # 4) 得卡短信登录(=注册)入库
  522. acc = self.sms_login(phone, code, proxy_url=proxy_url, log=log)
  523. return acc
  524. except Exception as e:
  525. log.warning(f"[补货] 单号注册异常: {e}")
  526. return None
  527. finally:
  528. # 5) 释放接码平台号(无论成败,及时释放降占用)
  529. if order_id is not None:
  530. try:
  531. isms.release_number(order_id=order_id)
  532. except Exception:
  533. pass
  534. def auto_replenish(self, target_min: int, target_max: int, free_proxy_urls: list = None,
  535. keyword: str = ISMS_PROJECT_KEYWORD, ascription: int = ISMS_ASCRIPTION,
  536. api_key: str = None, log=None) -> int:
  537. """巡检:healthy 号 < target_min 时,经接码平台自动注册补到 target_max。
  538. 供后台巡检循环定期调用。补货受两处上限约束:free_proxy_urls(空闲 IP,1 号 1 IP)与接码平台余额。
  539. 每补一个号随机停顿(REPLENISH_STAGGER_SEC)错峰,降一批新号被批量连坐识别的风险。
  540. ⚠️ 依赖外部资源(有 API Key + IP + 得卡项目存在后才能真跑):
  541. - api_key:接码平台 Key(传参或环境变量 ISMS_API_KEY);
  542. - keyword:得卡项目关键词需 search_projects 实测确认;
  543. - free_proxy_urls:可分配的空闲专属静态 IP 列表(1 号 1 IP)。
  544. Args:
  545. target_min (int): healthy 低于该值才触发补货。
  546. target_max (int): 补货目标(补到 healthy 达此值)。
  547. free_proxy_urls (list, optional): 可分配给新号的空闲专属 IP 列表。Defaults to None。
  548. keyword (str, optional): 接码平台得卡项目搜索关键词。Defaults to ISMS_PROJECT_KEYWORD。
  549. ascription (int, optional): 卡类型 2=实体卡/1=虚拟卡。Defaults to ISMS_ASCRIPTION。
  550. api_key (str, optional): 接码平台 Key;None 时用配置常量 ISMS_API_KEY。Defaults to None。
  551. log (optional): 日志对象。Defaults to None(用 self.log)。
  552. Returns:
  553. int: 本轮成功补充入库的账号数。
  554. """
  555. import random
  556. log = log or self.log
  557. cur = self.count_healthy()
  558. if cur >= target_min:
  559. return 0
  560. log.info(f"[补货] healthy={cur} < {target_min},开始自动补货至 {target_max}")
  561. try:
  562. from isms_client import ISmsClient
  563. isms = ISmsClient(api_key=api_key or ISMS_API_KEY) # 默认用配置里的 Key,传参可覆盖
  564. except Exception as e:
  565. log.error(f"[补货] 接码平台客户端初始化失败(检查 ISMS_API_KEY): {e}")
  566. return 0
  567. # 搜得卡项目
  568. pr = isms.search_projects(keyword)
  569. projects = (pr.get("data") or []) if pr.get("success") else []
  570. if not projects:
  571. log.error(f"[补货] 接码平台未搜到项目 '{keyword}'——请确认关键词或该平台是否有得卡项目")
  572. return 0
  573. project = projects[0]
  574. # 补货前先释放 dead 号占用的专属 IP(不复活、仅回收 IP),否则 IP 被 healthy+dead 占满时永远补不进
  575. released = self._release_dead_proxies()
  576. if released:
  577. log.info(f"[补货] 已释放 {released} 个 dead 账号占用的专属 IP 供新号回收")
  578. if free_proxy_urls is None:
  579. free_proxy_urls = self._free_proxy_urls() # 默认用配置 PROXY_URLS 里 DB 未占用的空闲 IP
  580. free = list(free_proxy_urls)
  581. added = 0
  582. fail_streak = 0 # 连续失败计数:防接码没号/异常时死循环空转烧钱
  583. while self.count_healthy() < target_max:
  584. proxy = free.pop(0) if free else None # 1 号 1 IP:无空闲 IP 则停(IP 是硬上限)
  585. if proxy is None:
  586. log.warning("[补货] 空闲 IP 已用尽,停止补货(IP 是号池规模硬上限)")
  587. break
  588. acc = self.register_one_via_isms(isms, project, proxy_url=proxy, ascription=ascription, log=log)
  589. if acc:
  590. added += 1
  591. fail_streak = 0
  592. log.info(f"[补货] 已补 {added} 个(最新 phone={acc['phone']} proxy={proxy})")
  593. else:
  594. if proxy is not None:
  595. free.insert(0, proxy) # 本号失败,IP 归还继续给下一个用
  596. fail_streak += 1
  597. if fail_streak >= 5: # 连续 5 次失败:接码可能没号/余额不足/异常,停止避免空转
  598. log.error("[补货] 连续 5 次注册失败,停止补货(检查接码余额/项目号源/网络)")
  599. break
  600. log.warning(f"[补货] 本次注册失败(连续 {fail_streak} 次),稍后重试")
  601. time.sleep(random.randint(*REPLENISH_STAGGER_SEC)) # 错峰,避免一批号被连坐
  602. log.info(f"[补货] 本轮补充完成 {added} 个,当前 healthy={self.count_healthy()}")
  603. return added
  604. # ---------- 内部 ----------
  605. def _row_to_dict(self, row) -> dict:
  606. """把 acquire 回读的行元组转成账号 dict(列顺序与 acquire 的 SELECT 对齐)。
  607. Args:
  608. row (tuple): SELECT 出的一行。
  609. Returns:
  610. dict: 账号字段 dict。
  611. """
  612. keys = ["id", "phone", "password", "country_code", "uid", "access_token",
  613. "refresh_token", "token_exp", "proxy_url", "status"]
  614. return dict(zip(keys, row))