feat(cheap-worker): compare_node 加有界并发(端口池+线程,可选)

并发的真正约束=九门 play/smoke 固定端口+Chrome 实例(gameId 前缀已隔离产物)。
查清整条链端口可参数化:gen.mjs --port/--cdp(done 门 smoke 用)、play.mjs
每次独立 spawn server+Chrome+userDataDir。据此:
- _port_pool:每并发槽一对独立 (server_port,cdp_port),段不重叠。
- run_pair_sync:同步整对跑(线程内 run_studio 自带 loop)。
- run_multi 加 conc:端口池(size=conc 天然限并发)+ asyncio.to_thread 每对一线程
  + gather。conc=1 串行;conc>1 有界并发。前台进程内有界并发,非后台子代理 task。
- _gen_node 传 --port/--cdp 给 gen.mjs(并发防 smoke 撞端口)。
CLI 加 --conc。test_compare_multi 14/14(+端口池不撞)。
This commit is contained in:
lili 2026-06-26 11:00:00 -07:00
parent 3a117c1ac4
commit 04dbd0608b
2 changed files with 77 additions and 24 deletions

View File

@ -195,28 +195,37 @@ def load_brief(genre: dict) -> str:
return ""
def _gen_node(brief: str, game_id: str, timeout: float = 900.0) -> dict:
"""Node 旧路 generation-only:gen.mjs --mode saa(scaffold+stage+自动 spec,SAA 不跑九门)。不在此 play。"""
def _gen_node(brief: str, game_id: str, port: int = 4320, cdp_port: int = 9222,
timeout: float = 900.0) -> dict:
"""Node 旧路 generation-only:gen.mjs --mode saa(scaffold+stage+done 门 smoke,SAA 不跑九门)。不在此 play。
传 --port/--cdp:gen.mjs done 门内部 smoke-boot 用这俩端口(行 142-143/249),并发时每槽独立端口防撞。
"""
r = subprocess.run(
["node", str(cheap_run._AMODEL_GEN / "gen.mjs"), "--mode", "saa", "--game-id", game_id, "--brief", brief],
["node", str(cheap_run._AMODEL_GEN / "gen.mjs"), "--mode", "saa", "--game-id", game_id,
"--brief", brief, "--port", str(port), "--cdp", str(cdp_port)],
cwd=str(cheap_run._GAME_RUNTIME), capture_output=True, text=True, timeout=timeout,
env=cheap_run._shell_env(), check=False,
)
return {"rc": r.returncode, "raw": (r.stdout + r.stderr)[-1200:]}
async def run_pair(genre: dict, brief: str, py_id: str, node_id: str,
def run_pair_sync(genre: dict, brief: str, py_id: str, node_id: str,
port: int = 4320, cdp_port: int = 9222) -> dict:
"""一局对照:两路各 generation-only → 注入同一份金标 spec → 各自单独 play。返回逐门 + 形态。"""
"""一局对照(同步、线程内整对跑):两路各 generation-only → 注入同一份金标 spec → 各自单独 play。
设计为同步,由 run_multi 经 asyncio.to_thread 并发调度(每对一个线程、独立端口);线程内 run_studio 用自带
event loop(asyncio.run)。返回逐门 + 形态。
"""
golden = golden_spec_path(genre)
# Python 路:generation-only(run_gates=False,内部不 play)→ 注入金标 → play。
await cheap_studio.run_studio(py_id, brief, run_gates=False, port=port, cdp_port=cdp_port)
asyncio.run(cheap_studio.run_studio(py_id, brief, run_gates=False, port=port, cdp_port=cdp_port))
inject_golden(py_id, golden)
py_play = cheap_run.play(py_id, port=port, cdp_port=cdp_port)
# Node 路:gen.mjs --mode saa(不 play)→ 注入金标(覆写自动 spec)→ play。
node_gen = _gen_node(brief, node_id)
# Node 路:gen.mjs --mode saa(不 play,内部 smoke 用同槽端口)→ 注入金标(覆写自动 spec)→ play。
node_gen = _gen_node(brief, node_id, port=port, cdp_port=cdp_port)
inject_golden(node_id, golden)
node_play = cheap_run.play(node_id, port=port, cdp_port=cdp_port)
@ -264,6 +273,15 @@ def _pair_ids(genre_key: str, k: int) -> tuple:
return f"cheap-{genre_key}-{k}", f"node-{genre_key}-{k}"
def _port_pool(conc: int, base_port: int = 4320, base_cdp: int = 9222) -> list:
"""每并发槽一对独立 (server_port, cdp_port):两两不撞、server 段(43xx)与 cdp 段(92xx)不重叠。
并发的真正约束 = 九门 play/smoke 的固定端口 + Chrome 实例;每槽独立端口才能并发不撞(gameId 前缀已隔离产物)。
"""
conc = max(1, conc)
return [(base_port + i * 2, base_cdp + i * 2) for i in range(conc)]
def _next_index(stems: list) -> int:
"""从已有 compare-multi-<N> 文件名算下一个序号(不覆盖历史)。"""
nums = [int(s.rsplit("-", 1)[-1]) for s in stems if s.rsplit("-", 1)[-1].isdigit()]
@ -291,11 +309,22 @@ def print_multi_report(report: dict) -> None:
print(f">>> 注:{s['note']}\n", file=sys.stderr)
async def run_multi(genre_keys: list, n: int, prefix: str = "cmp",
port: int = 4320, cdp_port: int = 9222) -> dict:
"""三品类 × n 串行对照:逐品类逐局生成→注入金标→play,按品类聚合。前台串行、前缀隔离。"""
async def run_multi(genre_keys: list, n: int, conc: int = 1,
base_port: int = 4320, base_cdp: int = 9222) -> dict:
"""三品类 × n 对照:逐局生成→注入金标→play,按品类聚合。
conc=1 串行;conc>1 有界并发 —— 端口池(每槽独立 port/cdp)+ 信号量限并发(≤conc 同时跑),每对经
asyncio.to_thread 在独立线程跑(线程内 run_studio 自带 loop)。前台进程内有界并发,非后台子代理。
"""
n = _clamp_n(n)
conc = max(1, conc)
pool = asyncio.Queue() # 端口池:size=conc,get 空则阻塞 → 天然限并发到 ≤conc
for pp in _port_pool(conc, base_port, base_cdp):
pool.put_nowait(pp)
# 收集每品类的 runs(线程并发写不同 key 的 list,key 互斥、无共享写)。
genre_runs = {}
valid = []
for key in genre_keys:
genre = _GENRE_BY_KEY.get(key)
if genre is None:
@ -305,16 +334,27 @@ async def run_multi(genre_keys: list, n: int, prefix: str = "cmp",
if not brief:
print(f"[skip] {key} 无 brief(base {genre['base']} run-summary 缺失)", file=sys.stderr)
continue
runs = []
for k in range(n):
valid.append((key, genre, brief))
genre_runs[key] = []
async def one(key, genre, brief, k):
py_id, node_id = _pair_ids(key, k)
print(f"[{key} {k + 1}/{n}] py={py_id} node={node_id} 生成→注入金标→play …", file=sys.stderr)
port, cdp = await pool.get() # 取一对独立端口(无空闲则等)
try:
runs.append(await run_pair(genre, brief, py_id, node_id, port=port, cdp_port=cdp_port))
print(f"[{key} {k + 1}/{n}] py={py_id} node={node_id} port={port}/{cdp} 生成→注入金标→play …",
file=sys.stderr)
r = await asyncio.to_thread(run_pair_sync, genre, brief, py_id, node_id, port, cdp)
genre_runs[key].append(r)
except Exception as e: # noqa: BLE001 单局失败不拖垮整批,记录后继续
print(f"[{key} {k + 1}/{n}] 异常:{type(e).__name__}: {e}", file=sys.stderr)
if runs:
genre_runs[key] = runs
finally:
pool.put_nowait((port, cdp)) # 归还端口
tasks = [one(key, genre, brief, k) for (key, genre, brief) in valid for k in range(n)]
print(f"[run] {len(valid)} 品类 × n={n} = {len(tasks)} 对,conc={conc}(端口段 {base_port}.. / {base_cdp}..)",
file=sys.stderr)
await asyncio.gather(*tasks)
genre_runs = {k: v for k, v in genre_runs.items() if v} # 去掉全失败的品类
report = aggregate(genre_runs)
out = _next_report_path()
@ -368,9 +408,9 @@ if __name__ == "__main__":
ap.add_argument("--genres", default="click-score,whack-mole,shop-serve",
help="逗号分隔品类键(默认三 tap-targets 代表品类)")
ap.add_argument("--n", type=int, default=3, help="每品类重复次数(n>5 小批,上限夹 5)")
ap.add_argument("--prefix", default="cmp", help="gameId 前缀基")
ap.add_argument("--conc", type=int, default=1, help="并发对数(1=串行;>1 端口池+线程有界并发,本机 Chrome 建议 ≤4)")
a = ap.parse_args()
keys = [k.strip() for k in a.genres.split(",") if k.strip()]
rep = asyncio.run(run_multi(keys, a.n, prefix=a.prefix))
rep = asyncio.run(run_multi(keys, a.n, conc=a.conc))
# 退出码:有任一品类坐实等价即 0(机制验证,非达标判定)。
sys.exit(0 if rep["summary"]["equivalentGenres"] else 1)

View File

@ -114,6 +114,19 @@ def test_pair_ids_isolation():
assert len(seen) == 3 * 3 * 2
def test_port_pool_distinct_non_overlapping():
"""并发端口池:每槽一对,所有 server/cdp 端口两两不撞、server 段与 cdp 段不重叠。"""
pool = C._port_pool(4)
assert len(pool) == 4
all_ports = [p for pair in pool for p in pair]
assert len(set(all_ports)) == len(all_ports), "端口有重复(并发会撞)"
server_ports = {pair[0] for pair in pool}
cdp_ports = {pair[1] for pair in pool}
assert not (server_ports & cdp_ports), "server 段与 cdp 段重叠"
assert C._port_pool(1) == [(4320, 9222)] # 串行=单槽默认端口
assert C._port_pool(0) == [(4320, 9222)] # conc<1 兜底至少一槽
def test_next_index_no_overwrite():
assert C._next_index([]) == 1
assert C._next_index(["compare-multi-1", "compare-multi-2"]) == 3