responses_survival.py 33 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855
  1. """Fixed cohort survival tracking using the codex responses API."""
  2. from __future__ import annotations
  3. import argparse
  4. from datetime import datetime
  5. import json
  6. import time
  7. from pathlib import Path
  8. from typing import Any
  9. from core.proxy_pool import ProxyLease, ProxyPool
  10. from core.settings import AppSettings
  11. from ops.common import get_management_key
  12. from platforms.chatgpt.fingerprint import OPENAI_FINGERPRINT_PROFILE
  13. from platforms.chatgpt.constants import OPENAI_USER_AGENT
  14. from platforms.chatgpt.pool import load_token_record, update_token_record
  15. from .scan import (
  16. ScanResult,
  17. _extract_credentials,
  18. _load_token_payload,
  19. _probe_responses_path,
  20. )
  21. def now_iso() -> str:
  22. return datetime.now().astimezone().isoformat(timespec="seconds")
  23. def _parse_iso(value: str) -> datetime | None:
  24. raw = str(value or "").strip()
  25. if not raw:
  26. return None
  27. try:
  28. return datetime.fromisoformat(raw)
  29. except Exception:
  30. return None
  31. def _duration_seconds(started_at: str, ended_at: str) -> int | None:
  32. start_dt = _parse_iso(started_at)
  33. end_dt = _parse_iso(ended_at)
  34. if start_dt is None or end_dt is None:
  35. return None
  36. return max(0, int((end_dt - start_dt).total_seconds()))
  37. def _member_age_seconds(member: dict[str, Any], probed_at: str) -> int | None:
  38. return _duration_seconds(str(member.get("created_at") or ""), probed_at)
  39. def _compact_text(value: str, limit: int = 320) -> str:
  40. return " ".join(str(value or "").split())[:limit]
  41. def _extract_error_facts(detail: str) -> tuple[str, str]:
  42. raw = str(detail or "").strip()
  43. if not raw:
  44. return "", ""
  45. try:
  46. payload = json.loads(raw)
  47. except Exception:
  48. return "", raw[:160]
  49. error = payload.get("error")
  50. if not isinstance(error, dict):
  51. return "", raw[:160]
  52. return str(error.get("code") or "").strip(), str(error.get("message") or "").strip()[:160]
  53. def _state_template(
  54. *,
  55. pool_dir: Path,
  56. cohort_size: int,
  57. proxy: str | None,
  58. timeout_seconds: int,
  59. seed_source: str = "latest_generated_pool_files",
  60. ) -> dict[str, Any]:
  61. return {
  62. "probe_mode": "responses",
  63. "probe_target": "codex_responses",
  64. "updated_at": "",
  65. "seeded_at": "",
  66. "seed_source": seed_source,
  67. "pool_dir": str(pool_dir),
  68. "cohort_size": max(1, int(cohort_size)),
  69. "proxy": str(proxy or "").strip() or None,
  70. "probe_fingerprint_profile": OPENAI_FINGERPRINT_PROFILE,
  71. "probe_user_agent": OPENAI_USER_AGENT,
  72. "timeout_seconds": max(5, int(timeout_seconds)),
  73. "members": [],
  74. "summary": {
  75. "tracked": 0,
  76. "alive": 0,
  77. "invalid": 0,
  78. "missing": 0,
  79. "removed_after_invalid": 0,
  80. "transport_error": 0,
  81. "suspicious": 0,
  82. "never_probed": 0,
  83. "first_invalid_count": 0,
  84. },
  85. "promotion_stats": {
  86. "promoted_success_total": 0,
  87. "promoted_failure_total": 0,
  88. "last_promoted_at": "",
  89. },
  90. "changes": [],
  91. "round_count": 0,
  92. }
  93. def load_responses_survival_state(path: Path) -> dict[str, Any]:
  94. if not path.is_file():
  95. return {}
  96. try:
  97. payload = json.loads(path.read_text(encoding="utf-8"))
  98. except Exception:
  99. return {}
  100. return payload if isinstance(payload, dict) else {}
  101. def _persist_state(path: Path, payload: dict[str, Any]) -> None:
  102. path.parent.mkdir(parents=True, exist_ok=True)
  103. tmp_path = path.with_name(f"{path.name}.tmp")
  104. tmp_path.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
  105. tmp_path.replace(path)
  106. def _seed_member(path: Path) -> dict[str, Any] | None:
  107. try:
  108. payload = load_token_record(path)
  109. except Exception:
  110. return None
  111. email = str(payload.get("email") or "").strip()
  112. access_token = str(payload.get("access_token") or "").strip()
  113. account_id = str(payload.get("account_id") or "").strip()
  114. if not email or not access_token or not account_id:
  115. return None
  116. created_at = str(payload.get("created_at") or "").strip()
  117. if not created_at:
  118. created_at = datetime.fromtimestamp(path.stat().st_mtime).astimezone().isoformat(timespec="seconds")
  119. return {
  120. "email": email,
  121. "file_name": path.name,
  122. "path": str(path),
  123. "created_at": created_at,
  124. "selected_at": now_iso(),
  125. "first_probe_at": "",
  126. "last_probe_at": "",
  127. "probe_count": 0,
  128. "last_probe_status_code": None,
  129. "last_probe_category": "",
  130. "last_probe_detail": "",
  131. "transport_error_count": 0,
  132. "suspicious_count": 0,
  133. "missing_at": "",
  134. "removed_after_invalid_at": "",
  135. "last_missing_detail": "",
  136. "first_invalid_at": "",
  137. "first_invalid_error_code": "",
  138. "first_invalid_error_message": "",
  139. "first_use_at": "",
  140. "first_use_age_seconds": None,
  141. "first_use_fingerprint_profile": "",
  142. "fingerprint_consistent": None,
  143. "registration_fingerprint_profile": str(payload.get("registration_fingerprint_profile") or "").strip(),
  144. "registration_proxy_key": str(payload.get("registration_proxy_key") or "").strip(),
  145. "registration_proxy_region": str(payload.get("registration_proxy_region") or "").strip(),
  146. "registration_post_create_gate": str(payload.get("registration_post_create_gate") or "").strip(),
  147. "warmup_required": bool(payload.get("warmup_required")) or str(payload.get("registration_post_create_gate") or "").strip().lower() == "add_phone",
  148. "warmup_state": str(payload.get("warmup_state") or "").strip()
  149. or ("pending" if str(payload.get("registration_post_create_gate") or "").strip().lower() == "add_phone" else "not_required"),
  150. "warmup_passed": bool(payload.get("warmup_passed")) or str(payload.get("warmup_state") or "").strip() == "passed",
  151. "warmup_completed_at": str(payload.get("warmup_completed_at") or "").strip(),
  152. "successful_probe_count": int(payload.get("successful_probe_count") or 0),
  153. "warmup_promotion_recorded": bool(payload.get("warmup_promotion_recorded")),
  154. "warmup_promotion_result": str(payload.get("warmup_promotion_result") or "").strip(),
  155. "first_use_proxy_key": str(payload.get("first_use_proxy_key") or "").strip(),
  156. "first_use_proxy_region": str(payload.get("first_use_proxy_region") or "").strip(),
  157. "first_invalid_proxy_key": str(payload.get("first_invalid_proxy_key") or "").strip(),
  158. "first_invalid_proxy_region": str(payload.get("first_invalid_proxy_region") or "").strip(),
  159. "survival_seconds": None,
  160. "state": "tracking",
  161. }
  162. def _collect_seed_members(
  163. pool_dir: Path,
  164. *,
  165. require_provenance: bool = False,
  166. recent_window_seconds: int = 0,
  167. ) -> list[dict[str, Any]]:
  168. recent_with_provenance: list[tuple[float, dict[str, Any]]] = []
  169. recent_without_provenance: list[tuple[float, dict[str, Any]]] = []
  170. cutoff_seconds = max(0, int(recent_window_seconds or 0))
  171. cutoff_ts = time.time() - cutoff_seconds if cutoff_seconds > 0 else 0.0
  172. for path in pool_dir.glob("*.json"):
  173. if not path.is_file():
  174. continue
  175. member = _seed_member(path)
  176. if member is None:
  177. continue
  178. created_at = _parse_iso(str(member.get("created_at") or ""))
  179. sort_ts = created_at.timestamp() if created_at is not None else path.stat().st_mtime
  180. if cutoff_ts and sort_ts < cutoff_ts:
  181. continue
  182. has_provenance = bool(str(member.get("registration_fingerprint_profile") or "").strip())
  183. if has_provenance:
  184. recent_with_provenance.append((sort_ts, member))
  185. elif not require_provenance:
  186. recent_without_provenance.append((sort_ts, member))
  187. recent_with_provenance.sort(key=lambda item: item[0], reverse=True)
  188. recent_without_provenance.sort(key=lambda item: item[0], reverse=True)
  189. ordered = [member for _ts, member in recent_with_provenance]
  190. ordered.extend(member for _ts, member in recent_without_provenance)
  191. return ordered
  192. def _seed_members(
  193. pool_dir: Path,
  194. cohort_size: int,
  195. *,
  196. require_provenance: bool = False,
  197. recent_window_seconds: int = 0,
  198. ) -> list[dict[str, Any]]:
  199. ordered = _collect_seed_members(
  200. pool_dir,
  201. require_provenance=require_provenance,
  202. recent_window_seconds=recent_window_seconds,
  203. )
  204. return ordered[: max(1, int(cohort_size))]
  205. def _member_identity(member: dict[str, Any]) -> str:
  206. return str(member.get("path") or member.get("file_name") or member.get("email") or "").strip()
  207. def _member_created_ts(member: dict[str, Any]) -> float:
  208. created = _parse_iso(str(member.get("created_at") or "").strip())
  209. if created is not None:
  210. return created.timestamp()
  211. path = Path(str(member.get("path") or "")).expanduser()
  212. if path.is_file():
  213. return path.stat().st_mtime
  214. return 0.0
  215. def _is_pending_warmup_member(member: dict[str, Any]) -> bool:
  216. return bool(member.get("warmup_required")) and not bool(member.get("warmup_passed")) and not _member_has_invalid_history(member)
  217. def _merge_member_state(existing: dict[str, Any], fresh: dict[str, Any]) -> dict[str, Any]:
  218. merged = dict(fresh)
  219. merged.update(existing)
  220. merged["selected_at"] = str(existing.get("selected_at") or fresh.get("selected_at") or now_iso())
  221. return merged
  222. def _member_priority(member: dict[str, Any]) -> tuple[float, float]:
  223. created_ts = _member_created_ts(member)
  224. if _is_pending_warmup_member(member):
  225. return (0.0, created_ts)
  226. last_probe = _parse_iso(str(member.get("last_probe_at") or "").strip())
  227. last_probe_ts = last_probe.timestamp() if last_probe is not None else 0.0
  228. if _member_has_invalid_history(member):
  229. return (1.0, -last_probe_ts)
  230. if str(member.get("last_probe_at") or "").strip():
  231. return (2.0, -last_probe_ts)
  232. return (3.0, -created_ts)
  233. def _refresh_active_members(
  234. pool_dir: Path,
  235. members: list[dict[str, Any]],
  236. *,
  237. cohort_size: int,
  238. require_provenance: bool = False,
  239. recent_window_seconds: int = 0,
  240. ) -> list[dict[str, Any]]:
  241. candidates = _collect_seed_members(
  242. pool_dir,
  243. require_provenance=require_provenance,
  244. recent_window_seconds=recent_window_seconds,
  245. )
  246. candidate_by_key = {_member_identity(member): member for member in candidates if _member_identity(member)}
  247. existing_keys: set[str] = set()
  248. combined: list[dict[str, Any]] = []
  249. for raw_member in members:
  250. if not isinstance(raw_member, dict):
  251. continue
  252. key = _member_identity(raw_member)
  253. if key:
  254. existing_keys.add(key)
  255. fresh = candidate_by_key.get(key)
  256. combined.append(_merge_member_state(raw_member, fresh) if isinstance(fresh, dict) else raw_member)
  257. for candidate in candidates:
  258. key = _member_identity(candidate)
  259. if not key or key in existing_keys:
  260. continue
  261. combined.append(candidate)
  262. deduped: list[dict[str, Any]] = []
  263. seen: set[str] = set()
  264. for member in sorted(combined, key=_member_priority):
  265. key = _member_identity(member)
  266. if key and key in seen:
  267. continue
  268. if key:
  269. seen.add(key)
  270. deduped.append(member)
  271. if len(deduped) >= max(1, int(cohort_size)):
  272. break
  273. return deduped
  274. def probe_responses_token_file(path: Path, proxy: str | None, timeout_seconds: int) -> ScanResult:
  275. payload, load_error = _load_token_payload(path)
  276. if load_error is not None:
  277. return load_error
  278. assert payload is not None
  279. credentials = _extract_credentials(path, payload)
  280. if isinstance(credentials, ScanResult):
  281. return credentials
  282. access_token, account_id = credentials
  283. return _probe_responses_path(path, access_token, account_id, proxy, timeout_seconds)
  284. def _member_outcome(member: dict[str, Any]) -> str:
  285. state = str(member.get("state") or "").strip()
  286. if state == "invalid_removed":
  287. return "invalid_removed"
  288. category = str(member.get("last_probe_category") or "").strip()
  289. return category or "never_probed"
  290. def _member_has_invalid_history(member: dict[str, Any]) -> bool:
  291. state = str(member.get("state") or "").strip()
  292. category = str(member.get("last_probe_category") or "").strip()
  293. return bool(str(member.get("first_invalid_at") or "").strip()) or state in {"invalid", "invalid_removed"} or category == "invalid"
  294. def _preserve_terminal_invalid(member: dict[str, Any]) -> None:
  295. if str(member.get("last_probe_category") or "").strip() != "invalid":
  296. member["last_probe_category"] = "invalid"
  297. if member.get("last_probe_status_code") in {None, ""}:
  298. member["last_probe_status_code"] = 401
  299. detail = str(member.get("last_probe_detail") or "").strip()
  300. if not detail or detail.startswith("missing_file:"):
  301. member["last_probe_detail"] = "invalid_before_pool_removal"
  302. def _persist_member_fields(member: dict[str, Any]) -> None:
  303. path = Path(str(member.get("path") or "")).expanduser()
  304. if not path.is_file():
  305. return
  306. try:
  307. update_token_record(
  308. path,
  309. warmup_required=bool(member.get("warmup_required")),
  310. warmup_state=str(member.get("warmup_state") or "").strip(),
  311. warmup_passed=bool(member.get("warmup_passed")),
  312. warmup_completed_at=str(member.get("warmup_completed_at") or "").strip(),
  313. successful_probe_count=int(member.get("successful_probe_count") or 0),
  314. warmup_promotion_recorded=bool(member.get("warmup_promotion_recorded")),
  315. warmup_promotion_result=str(member.get("warmup_promotion_result") or "").strip(),
  316. first_use_at=str(member.get("first_use_at") or "").strip(),
  317. first_use_age_seconds=member.get("first_use_age_seconds"),
  318. first_use_proxy_key=str(member.get("first_use_proxy_key") or "").strip(),
  319. first_use_proxy_region=str(member.get("first_use_proxy_region") or "").strip(),
  320. first_invalid_at=str(member.get("first_invalid_at") or "").strip(),
  321. first_invalid_error_code=str(member.get("first_invalid_error_code") or "").strip(),
  322. first_invalid_error_message=str(member.get("first_invalid_error_message") or "").strip(),
  323. first_invalid_proxy_key=str(member.get("first_invalid_proxy_key") or "").strip(),
  324. first_invalid_proxy_region=str(member.get("first_invalid_proxy_region") or "").strip(),
  325. survival_seconds=member.get("survival_seconds"),
  326. )
  327. except Exception:
  328. return
  329. def _maybe_promote_warmup_member(
  330. member: dict[str, Any],
  331. *,
  332. settings: AppSettings | None,
  333. probed_at: str,
  334. ) -> str | None:
  335. if settings is None or str(settings.backend or "").strip().lower() != "cpa":
  336. return None
  337. if not bool(member.get("warmup_required")) or not bool(member.get("warmup_passed")):
  338. return None
  339. path = Path(str(member.get("path") or "")).expanduser()
  340. if not path.is_file():
  341. return None
  342. try:
  343. payload = load_token_record(path)
  344. except Exception:
  345. return None
  346. if not isinstance(payload, dict):
  347. return None
  348. if str(payload.get("cpa_sync_status") or "").strip().lower() == "synced":
  349. return None
  350. key = str(get_management_key() or "").strip()
  351. if not key:
  352. update_token_record(
  353. path,
  354. cpa_sync_status="failed",
  355. last_cpa_sync_at=probed_at,
  356. last_cpa_sync_error="CPA management key unavailable",
  357. )
  358. member["cpa_sync_status"] = "failed"
  359. return "failed"
  360. api_url = str(settings.cpa_management_base_url or "").strip()
  361. if api_url.endswith("/v0/management"):
  362. api_url = api_url[: -len("/v0/management")]
  363. from platforms.chatgpt.cpa_upload import upload_to_cpa
  364. ok, message = upload_to_cpa(payload, api_url=api_url.rstrip("/"), api_key=key, proxy=None)
  365. if ok:
  366. update_token_record(
  367. path,
  368. health_status="good",
  369. cpa_sync_status="synced",
  370. last_cpa_sync_at=probed_at,
  371. last_cpa_sync_error="",
  372. )
  373. member["cpa_sync_status"] = "synced"
  374. return "success"
  375. else:
  376. update_token_record(
  377. path,
  378. cpa_sync_status="failed",
  379. last_cpa_sync_at=probed_at,
  380. last_cpa_sync_error=message,
  381. )
  382. member["cpa_sync_status"] = "failed"
  383. return "failed"
  384. def _update_member(
  385. member: dict[str, Any],
  386. result: ScanResult,
  387. probed_at: str,
  388. *,
  389. probe_proxy_key: str = "",
  390. probe_proxy_region: str = "",
  391. warmup_min_age_seconds: int = 600,
  392. warmup_min_successful_probes: int = 2,
  393. ) -> dict[str, Any]:
  394. previous_outcome = _member_outcome(member)
  395. member["last_probe_at"] = probed_at
  396. if not str(member.get("first_probe_at") or "").strip():
  397. member["first_probe_at"] = probed_at
  398. if not str(member.get("first_use_at") or "").strip():
  399. member["first_use_at"] = probed_at
  400. member["first_use_age_seconds"] = _duration_seconds(
  401. str(member.get("created_at") or "").strip(),
  402. probed_at,
  403. )
  404. member["first_use_fingerprint_profile"] = OPENAI_FINGERPRINT_PROFILE
  405. member["first_use_proxy_key"] = str(probe_proxy_key or "").strip()
  406. member["first_use_proxy_region"] = str(probe_proxy_region or "").strip()
  407. registration_profile = str(member.get("registration_fingerprint_profile") or "").strip()
  408. if registration_profile:
  409. member["fingerprint_consistent"] = registration_profile == OPENAI_FINGERPRINT_PROFILE
  410. member["probe_count"] = int(member.get("probe_count") or 0) + 1
  411. if result.category == "missing" and _member_has_invalid_history(member):
  412. if not str(member.get("missing_at") or "").strip():
  413. member["missing_at"] = probed_at
  414. if not str(member.get("removed_after_invalid_at") or "").strip():
  415. member["removed_after_invalid_at"] = probed_at
  416. member["last_missing_detail"] = _compact_text(result.detail or "")
  417. _preserve_terminal_invalid(member)
  418. member["state"] = "invalid_removed"
  419. next_outcome = _member_outcome(member)
  420. detail = _compact_text(
  421. f"removed_after_invalid | {member.get('last_missing_detail') or ''}"
  422. )
  423. return {
  424. "email": str(member.get("email") or "").strip(),
  425. "from": previous_outcome,
  426. "to": next_outcome,
  427. "probed_at": probed_at,
  428. "survival_seconds": member.get("survival_seconds"),
  429. "detail": detail,
  430. }
  431. if result.category != "invalid" and _member_has_invalid_history(member):
  432. member["post_invalid_probe_at"] = probed_at
  433. member["post_invalid_probe_category"] = result.category
  434. member["post_invalid_probe_detail"] = _compact_text(result.detail or "")
  435. _preserve_terminal_invalid(member)
  436. member["state"] = "invalid"
  437. next_outcome = _member_outcome(member)
  438. return {
  439. "email": str(member.get("email") or "").strip(),
  440. "from": previous_outcome,
  441. "to": next_outcome,
  442. "probed_at": probed_at,
  443. "survival_seconds": member.get("survival_seconds"),
  444. "detail": member["post_invalid_probe_detail"],
  445. }
  446. member["last_probe_status_code"] = result.status_code
  447. member["last_probe_category"] = result.category
  448. member["last_probe_detail"] = _compact_text(result.detail or "")
  449. if result.category == "transport_error":
  450. member["transport_error_count"] = int(member.get("transport_error_count") or 0) + 1
  451. elif result.category not in {"normal", "invalid", "missing"}:
  452. member["suspicious_count"] = int(member.get("suspicious_count") or 0) + 1
  453. elif result.category == "missing" and not str(member.get("missing_at") or "").strip():
  454. member["missing_at"] = probed_at
  455. if result.category == "invalid":
  456. if not str(member.get("first_invalid_at") or "").strip():
  457. member["first_invalid_at"] = probed_at
  458. member["survival_seconds"] = _duration_seconds(
  459. str(member.get("created_at") or "").strip() or str(member.get("first_probe_at") or "").strip(),
  460. probed_at,
  461. )
  462. error_code, error_message = _extract_error_facts(result.detail or "")
  463. member["first_invalid_error_code"] = error_code
  464. member["first_invalid_error_message"] = error_message
  465. member["first_invalid_proxy_key"] = str(probe_proxy_key or "").strip()
  466. member["first_invalid_proxy_region"] = str(probe_proxy_region or "").strip()
  467. if bool(member.get("warmup_required")):
  468. member["warmup_state"] = "failed"
  469. member["warmup_passed"] = False
  470. member["state"] = "invalid"
  471. elif result.category == "missing":
  472. member["state"] = "missing"
  473. else:
  474. if result.category == "normal":
  475. member["successful_probe_count"] = int(member.get("successful_probe_count") or 0) + 1
  476. if bool(member.get("warmup_required")) and not bool(member.get("warmup_passed")):
  477. current_age = _member_age_seconds(member, probed_at)
  478. if (
  479. int(member.get("successful_probe_count") or 0) >= max(1, int(warmup_min_successful_probes))
  480. and (current_age is not None and int(current_age) >= max(0, int(warmup_min_age_seconds)))
  481. ):
  482. member["warmup_passed"] = True
  483. member["warmup_state"] = "passed"
  484. member["warmup_completed_at"] = probed_at
  485. member["state"] = "tracking"
  486. _persist_member_fields(member)
  487. next_outcome = _member_outcome(member)
  488. return {
  489. "email": str(member.get("email") or "").strip(),
  490. "from": previous_outcome,
  491. "to": next_outcome,
  492. "probed_at": probed_at,
  493. "survival_seconds": member.get("survival_seconds"),
  494. "detail": member["last_probe_detail"],
  495. }
  496. def _probe_member_proxy(
  497. member: dict[str, Any],
  498. *,
  499. default_proxy: str | None,
  500. proxy_pool: ProxyPool | None,
  501. ) -> tuple[str | None, str, str, ProxyLease | None]:
  502. preferred_name = str(member.get("registration_proxy_key") or "").strip()
  503. preferred_region = str(member.get("registration_proxy_region") or "").strip().lower()
  504. if proxy_pool is None:
  505. return default_proxy, preferred_name, preferred_region, None
  506. try:
  507. lease = proxy_pool.acquire(
  508. timeout=5.0,
  509. preferred_name=preferred_name or None,
  510. preferred_regions=(preferred_region,) if preferred_region else (),
  511. )
  512. except Exception:
  513. return default_proxy, preferred_name, preferred_region, None
  514. return lease.proxy_url, lease.name, preferred_region or "", lease
  515. def _build_summary(members: list[dict[str, Any]]) -> dict[str, int]:
  516. summary = {
  517. "tracked": len(members),
  518. "alive": 0,
  519. "invalid": 0,
  520. "missing": 0,
  521. "removed_after_invalid": 0,
  522. "transport_error": 0,
  523. "suspicious": 0,
  524. "never_probed": 0,
  525. "first_invalid_count": 0,
  526. }
  527. for member in members:
  528. outcome = _member_outcome(member)
  529. if outcome == "never_probed":
  530. summary["never_probed"] += 1
  531. elif outcome == "normal":
  532. summary["alive"] += 1
  533. elif outcome == "invalid":
  534. summary["invalid"] += 1
  535. elif outcome == "invalid_removed":
  536. summary["invalid"] += 1
  537. summary["removed_after_invalid"] += 1
  538. elif outcome == "missing":
  539. summary["missing"] += 1
  540. elif outcome == "transport_error":
  541. summary["transport_error"] += 1
  542. else:
  543. summary["suspicious"] += 1
  544. if str(member.get("first_invalid_at") or "").strip():
  545. summary["first_invalid_count"] += 1
  546. return summary
  547. def responses_survival_once(
  548. *,
  549. pool_dir: Path,
  550. state_file: Path,
  551. cohort_size: int,
  552. proxy: str | None,
  553. timeout_seconds: int,
  554. reseed: bool = False,
  555. proxy_pool: ProxyPool | None = None,
  556. require_provenance: bool = False,
  557. recent_window_seconds: int = 0,
  558. warmup_min_age_seconds: int = 600,
  559. warmup_min_successful_probes: int = 2,
  560. settings: AppSettings | None = None,
  561. ) -> dict[str, Any]:
  562. state = load_responses_survival_state(state_file)
  563. seeded = False
  564. reseeded = False
  565. if not state or reseed:
  566. state = _state_template(
  567. pool_dir=pool_dir,
  568. cohort_size=cohort_size,
  569. proxy=proxy,
  570. timeout_seconds=timeout_seconds,
  571. )
  572. state["members"] = _seed_members(
  573. pool_dir,
  574. int(state.get("cohort_size") or cohort_size),
  575. require_provenance=require_provenance,
  576. recent_window_seconds=recent_window_seconds,
  577. )
  578. state["seeded_at"] = now_iso()
  579. seeded = True
  580. reseeded = reseed
  581. else:
  582. state.setdefault("probe_mode", "responses")
  583. state.setdefault("probe_target", "codex_responses")
  584. state.setdefault("pool_dir", str(pool_dir))
  585. state.setdefault("cohort_size", max(1, int(cohort_size)))
  586. state.setdefault("proxy", str(proxy or "").strip() or None)
  587. state.setdefault("probe_fingerprint_profile", OPENAI_FINGERPRINT_PROFILE)
  588. state.setdefault("probe_user_agent", OPENAI_USER_AGENT)
  589. state.setdefault("timeout_seconds", max(5, int(timeout_seconds)))
  590. state.setdefault("members", [])
  591. state.setdefault("summary", {})
  592. state.setdefault("changes", [])
  593. state.setdefault(
  594. "promotion_stats",
  595. {
  596. "promoted_success_total": 0,
  597. "promoted_failure_total": 0,
  598. "last_promoted_at": "",
  599. },
  600. )
  601. state.setdefault("seed_source", "latest_generated_pool_files")
  602. state.setdefault("round_count", 0)
  603. if not isinstance(state.get("members"), list):
  604. state["members"] = []
  605. if not isinstance(state.get("promotion_stats"), dict):
  606. state["promotion_stats"] = {
  607. "promoted_success_total": 0,
  608. "promoted_failure_total": 0,
  609. "last_promoted_at": "",
  610. }
  611. if not state["members"]:
  612. state["members"] = _seed_members(
  613. pool_dir,
  614. int(state.get("cohort_size") or cohort_size),
  615. require_provenance=require_provenance,
  616. recent_window_seconds=recent_window_seconds,
  617. )
  618. state["seeded_at"] = now_iso()
  619. seeded = True
  620. else:
  621. state["members"] = _refresh_active_members(
  622. pool_dir,
  623. [member for member in state["members"] if isinstance(member, dict)],
  624. cohort_size=int(state.get("cohort_size") or cohort_size),
  625. require_provenance=require_provenance,
  626. recent_window_seconds=recent_window_seconds,
  627. )
  628. changes: list[dict[str, Any]] = []
  629. for raw_member in state["members"]:
  630. if not isinstance(raw_member, dict):
  631. continue
  632. member = raw_member
  633. probed_at = now_iso()
  634. selected_proxy, probe_proxy_key, probe_proxy_region, lease = _probe_member_proxy(
  635. member,
  636. default_proxy=str(state.get("proxy") or "").strip() or None,
  637. proxy_pool=proxy_pool,
  638. )
  639. try:
  640. result = probe_responses_token_file(
  641. Path(str(member.get("path") or "")),
  642. selected_proxy,
  643. max(5, int(state.get("timeout_seconds") or timeout_seconds)),
  644. )
  645. finally:
  646. if lease is not None and proxy_pool is not None:
  647. try:
  648. proxy_pool.release(lease, success=result.category == "normal" if 'result' in locals() else None, stage=result.category if 'result' in locals() else None)
  649. except Exception:
  650. pass
  651. change = _update_member(
  652. member,
  653. result,
  654. probed_at,
  655. probe_proxy_key=probe_proxy_key,
  656. probe_proxy_region=probe_proxy_region,
  657. warmup_min_age_seconds=warmup_min_age_seconds,
  658. warmup_min_successful_probes=warmup_min_successful_probes,
  659. )
  660. promotion_outcome = _maybe_promote_warmup_member(
  661. member,
  662. settings=settings,
  663. probed_at=probed_at,
  664. )
  665. if promotion_outcome in {"success", "failed"} and not bool(member.get("warmup_promotion_recorded")):
  666. stats = state["promotion_stats"]
  667. member["warmup_promotion_recorded"] = True
  668. member["warmup_promotion_result"] = promotion_outcome
  669. if promotion_outcome == "success":
  670. stats["promoted_success_total"] = int(stats.get("promoted_success_total") or 0) + 1
  671. else:
  672. stats["promoted_failure_total"] = int(stats.get("promoted_failure_total") or 0) + 1
  673. stats["last_promoted_at"] = probed_at
  674. _persist_member_fields(member)
  675. if change["from"] != change["to"]:
  676. changes.append(change)
  677. state["updated_at"] = now_iso()
  678. state["summary"] = _build_summary([member for member in state["members"] if isinstance(member, dict)])
  679. state["changes"] = changes
  680. state["seeded"] = seeded
  681. state["reseeded"] = reseeded
  682. state["state_file"] = str(state_file)
  683. state["round_count"] = int(state.get("round_count") or 0) + 1
  684. _persist_state(state_file, state)
  685. return state
  686. def print_responses_survival_summary(result: dict[str, Any]) -> None:
  687. summary = result.get("summary") if isinstance(result.get("summary"), dict) else {}
  688. print(
  689. "[responses-survival] summary | "
  690. f"tracked={int(summary.get('tracked') or 0)} | alive={int(summary.get('alive') or 0)} "
  691. f"| invalid={int(summary.get('invalid') or 0)} | missing={int(summary.get('missing') or 0)} "
  692. f"| removed_after_invalid={int(summary.get('removed_after_invalid') or 0)} "
  693. f"| transport_error={int(summary.get('transport_error') or 0)} "
  694. f"| suspicious={int(summary.get('suspicious') or 0)}"
  695. )
  696. for change in result.get("changes") or []:
  697. if not isinstance(change, dict):
  698. continue
  699. survival_seconds = change.get("survival_seconds")
  700. survival_text = f" | survival={survival_seconds}s" if survival_seconds is not None else ""
  701. print(
  702. f"[responses-survival] state change | {change.get('email') or '?'} | "
  703. f"{change.get('from') or 'never_probed'} -> {change.get('to') or '?'}{survival_text}"
  704. )
  705. print(f"[responses-survival] state={result.get('state_file') or '-'}")
  706. def run_responses_survival_loop(
  707. *,
  708. pool_dir: Path,
  709. state_file: Path,
  710. cohort_size: int,
  711. proxy: str | None,
  712. timeout_seconds: int,
  713. interval_seconds: int,
  714. reseed: bool = False,
  715. max_rounds: int = 0,
  716. settings: AppSettings | None = None,
  717. ) -> dict[str, Any]:
  718. round_index = 0
  719. last_result: dict[str, Any] = {}
  720. proxy_pool = ProxyPool.from_settings(settings) if settings is not None else None
  721. if proxy_pool is not None:
  722. proxy_pool.start()
  723. while True:
  724. round_index += 1
  725. last_result = responses_survival_once(
  726. pool_dir=pool_dir,
  727. state_file=state_file,
  728. cohort_size=cohort_size,
  729. proxy=proxy,
  730. timeout_seconds=timeout_seconds,
  731. reseed=reseed and round_index == 1,
  732. proxy_pool=proxy_pool,
  733. require_provenance=bool(settings.responses_survival_require_provenance) if settings is not None else False,
  734. recent_window_seconds=int(settings.responses_survival_recent_window_seconds) if settings is not None else 0,
  735. warmup_min_age_seconds=int(settings.warmup_min_age_seconds) if settings is not None else 600,
  736. warmup_min_successful_probes=int(settings.warmup_min_successful_probes) if settings is not None else 2,
  737. settings=settings,
  738. )
  739. print_responses_survival_summary(last_result)
  740. if max_rounds > 0 and round_index >= max_rounds:
  741. if proxy_pool is not None:
  742. proxy_pool.close()
  743. return last_result
  744. time.sleep(max(5, int(interval_seconds)))
  745. def build_arg_parser() -> argparse.ArgumentParser:
  746. parser = argparse.ArgumentParser(description="Run responses survival tracking for recently created accounts.")
  747. parser.add_argument("--pool-dir", required=True)
  748. parser.add_argument("--state-file", required=True)
  749. parser.add_argument("--cohort-size", type=int, default=8)
  750. parser.add_argument("--proxy", default="")
  751. parser.add_argument("--timeout-seconds", type=int, default=30)
  752. parser.add_argument("--interval-seconds", type=int, default=60)
  753. parser.add_argument("--max-rounds", type=int, default=0, help="0 means run forever")
  754. parser.add_argument("--reseed", action="store_true")
  755. return parser
  756. def main(argv: list[str] | None = None) -> int:
  757. parser = build_arg_parser()
  758. args = parser.parse_args(argv)
  759. run_responses_survival_loop(
  760. pool_dir=Path(args.pool_dir).expanduser().resolve(),
  761. state_file=Path(args.state_file).expanduser().resolve(),
  762. cohort_size=max(1, int(args.cohort_size)),
  763. proxy=str(args.proxy or "").strip() or None,
  764. timeout_seconds=max(5, int(args.timeout_seconds)),
  765. interval_seconds=max(5, int(args.interval_seconds)),
  766. reseed=bool(args.reseed),
  767. max_rounds=max(0, int(args.max_rounds)),
  768. )
  769. return 0
  770. if __name__ == "__main__":
  771. raise SystemExit(main())