| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696 |
- # -*- 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))
|