"""统一超时配置(修永久卡死). 集中各阶段"单次外部调用"的总时长上限(用户拍板),并提供 httpx.Timeout 工厂。 要点:httpx 的 `timeout=标量` 只把 connect/read/write/pool 各设为 N,**没有"整次请求总时长"**; 服务端慢速吐字节时每次 read 都在 N 内返回一点 → read 永不触发 → 永久卡在 do_poll。 所以这里强制 **read 相设短**(停止吐数据即抛 ReadTimeout),总时长由 write 相 + 调用方护栏兜。 env 可覆盖各阶段总值,但被硬上限钳制,防再配出 3600 那种值。 """ from __future__ import annotations import os from typing import Mapping import httpx CONNECT_TIMEOUT_SECONDS = 10.0 # 各阶段总时长默认(秒)——用户拍板的"单次外部调用允许上限"。 _DEFAULTS: dict[str, float] = { "oss": 300.0, # OSS 上传/归档单次尝试 5min "video_download": 600.0, # 视频下载 10min "video_llm": 600.0, # qwen/gemini 单次判定 10min "crawapi": 180.0, # 平台搜索/作者/画像 3min "query_llm": 120.0, # query variant / Gate2 2min "pg": 30.0, # pattern PG Gate1 30s } # env 覆盖也不得超过(防误配)。 _HARD_CEILING: dict[str, float] = { "oss": 600.0, "video_download": 1200.0, "video_llm": 1200.0, "crawapi": 360.0, "query_llm": 300.0, "pg": 60.0, } # 单次 read(两次收到数据之间)上限。OSS 按业务要求放宽到 300s。 _READ: dict[str, float] = { "oss": 300.0, "video_download": 120.0, "video_llm": 120.0, "crawapi": 60.0, "query_llm": 60.0, "pg": 30.0, } _ENV_KEYS: dict[str, tuple[str, ...]] = { "oss": ("CONTENT_AGENT_OSS_TIMEOUT_SECONDS",), "video_download": ("CONTENT_AGENT_VIDEO_DOWNLOAD_TIMEOUT_SECONDS",), "video_llm": ("CONTENT_AGENT_VIDEO_LLM_TIMEOUT_SECONDS",), "crawapi": ("CONTENTFIND_API_CRAWAPI_TIMEOUT_SECONDS",), "query_llm": ("CONTENT_AGENT_QUERY_LLM_TIMEOUT_SECONDS",), "pg": ("OPEN_AIGC_PG_TIMEOUT_SECONDS",), } def total_timeout(stage: str, env: Mapping[str, str] | None = None) -> float: """阶段总时长(秒):env 覆盖 → 硬上限钳制 → 默认。""" src = os.environ if env is None else env value = _DEFAULTS[stage] for key in _ENV_KEYS[stage]: raw = src.get(key) if raw: try: value = float(raw) break except (TypeError, ValueError): pass return min(value, _HARD_CEILING[stage]) def read_timeout(stage: str) -> float: return _READ[stage] def as_httpx_timeout( total_seconds: float, *, read: float, connect: float = CONNECT_TIMEOUT_SECONDS, ) -> httpx.Timeout: """把一个总时长(秒)转成分段 httpx.Timeout:read 短、write=总、connect 短。""" total = max(float(total_seconds), 1.0) return httpx.Timeout( connect=min(connect, total), read=min(read, total), write=total, pool=min(connect, total), ) def httpx_timeout(stage: str, env: Mapping[str, str] | None = None) -> httpx.Timeout: """按阶段直接构造 httpx.Timeout(已含 env 覆盖 + read 短上限)。""" return as_httpx_timeout(total_timeout(stage, env=env), read=_READ[stage])