两段式第二段的 regenerate-module 类:改玩法 = 据 intent 有界重写 game-logic.js,不动其余。
复用机制(参数化、不复制 resume/熔断/收口):
- cheap_studio.run_studio 加 4 个默认 None 可选参(system_prompt/initial_kick/write_whitelist/prepare)
切到 modify 态,create 路零行为变化(params=None 等价原逻辑)
- cheap_toolkit.build_toolkit 加 write_whitelist(非白名单 basename 写直接拒)
- cheap_roles.build_modify_system_prompt(复用 create 红线契约块)
cheap_modify.execute_regenerate_modify:
- 取 intent(空→failed)→ materialize base 源 → 重写前后对非目标文件算 hash
- 有界重写(写边界收窄到只 game-logic.js)→ 九门
- status=succeeded ⟺ 九门过 ∧ 非目标稳;manifest{file,kind:behavior,intent,changed,untouchedStable}(断言②地基)
worker_service._process_regenerate_job(regen_fn):把 M3 的 regenerate-module 显式 failed 占位换成真执行;
deterministic/create/生成路一字未动。
测试:test_a11_m4_regenerate 10/10(全注入桩零 LLM)+ 全回归独立复跑 73 测试全绿。
一次真 LLM smoke(单跑):intent=连击递增 → status=succeeded、untouchedStable=True、九门 pass、
attempts=1、¥0.31、137s;独立 diff 佐证 5 非目标文件字节相同、仅 game-logic.js 变。
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
528 lines
26 KiB
Python
528 lines
26 KiB
Python
"""worker_service.py — 把 cheap-worker 的 run_studio 包成 §6.1 HTTP worker(M3a U1)。
|
|
|
|
创始人 D1:Python 新路走 HTTP worker 复用后端默认派发通道(WorkerDispatchClient)——
|
|
接 §6.1 job-in POST → HTTP 2xx 投递握手立返 → 内部有界串行队列跑 run_studio(生成+九门)→
|
|
result_out 组 §6.1 result-out(三字段语义分离)→ 带 HMAC POST 回 job 携带的 callback URL
|
|
(/admin-api/aigc/dify/callback-internal)。镜像 wg1/gen-worker/worker/service.py 骨架,只换生成核心。
|
|
|
|
KTD4(治串行 worker 忙时 409 误判):job-in 一律入有界队列(有容量 202、队满 503),单 worker 线程
|
|
串行消费(Chrome serve 4320 / CDP 9222 不可并发);执行器 2xx=ok 契约不动。队满 503 仍非 2xx →
|
|
执行器置 LLM_ERROR 终态(罕见背压、可接受;配套 U4 失败任务允许同 idempotencyKey 重提)。
|
|
|
|
可测性:核心逻辑(parse_job / compute_signature / try_enqueue / process_job)纯函数化 + 依赖
|
|
(run_fn / send_fn / profile_fn / opener)可注入——单测零真实网络 / 零 LLM(对 wg1_stub_worker 可测范式)。
|
|
"""
|
|
|
|
import argparse
|
|
import hashlib
|
|
import hmac
|
|
import json
|
|
import queue
|
|
import threading
|
|
import urllib.error
|
|
import urllib.request
|
|
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
|
from pathlib import Path
|
|
|
|
import cheap_classify
|
|
import dedup
|
|
import result_out
|
|
|
|
# 默认监听端口 / 队列容量(KTD4 有界;并发上限另由创始人容量画像 conc≤15 + 远期多实例控)。
|
|
DEFAULT_PORT = 9501
|
|
DEFAULT_QUEUE_MAXSIZE = 8
|
|
# 内网免鉴权回调子路由(后端 AdminAigcTaskController.difyCallbackInternal,@PermitAll + HMAC 验签)。
|
|
DEFAULT_CALLBACK = "http://localhost:48080/admin-api/aigc/dify/callback-internal"
|
|
|
|
# 内网回调 opener:禁代理。回调目标恒为内网(localhost / Tailscale 100.64.x),经系统代理会被
|
|
# fake-ip(198.18.x)拦成 502(本机实测踩坑);故回调始终走直连、绕过代理(对 _bootstrap proxy-bypass 同理)。
|
|
_NO_PROXY_OPENER = urllib.request.build_opener(urllib.request.ProxyHandler({}))
|
|
|
|
|
|
def log(msg: str) -> None:
|
|
"""带前缀日志(stdout,便于 nohup/setsid 落盘排障)。"""
|
|
print(f"[cheap-worker-service] {msg}", flush=True)
|
|
|
|
|
|
# ---------- §6.1 job-in 解析 ----------
|
|
|
|
def parse_job(raw: bytes) -> dict | None:
|
|
"""解析 §6.1 job-in 原始字节为 dict;坏 JSON 返 None(不抛,由 handler 回 400)。"""
|
|
if not raw:
|
|
return {}
|
|
try:
|
|
obj = json.loads(raw.decode("utf-8"))
|
|
return obj if isinstance(obj, dict) else None
|
|
except Exception: # noqa: BLE001 — 任意解析异常都归 None,真因由 handler 日志记
|
|
return None
|
|
|
|
|
|
def parse_modify(job: dict) -> dict | None:
|
|
"""从 §6.1 job 提取 modify 区(A11 调整路);无 modifyMode = create 路、返 None(现行 create 零副作用)。
|
|
|
|
后端 putModifyFieldsIntoJob 放的键:modifyMode / baseVersionId / modifyPatch / sourceProject(base 源工程
|
|
JSON 串,据 baseVersionId 反查注入)。A11 执行段(M3 确定性类 / M4 模块重生成)据此取 base 源 + 改动意图;
|
|
判意图段(M2)产出的结构化改动经 modifyPatch 带入。本函数只提取不执行——执行落 M3/M4。
|
|
"""
|
|
mode = job.get("modifyMode")
|
|
if not mode:
|
|
return None
|
|
return {
|
|
"mode": mode,
|
|
"baseVersionId": job.get("baseVersionId"),
|
|
"modifyPatch": job.get("modifyPatch"),
|
|
"sourceProject": job.get("sourceProject"), # base 源工程 JSON 串(后端反查注入;worker 只消费不反查)
|
|
}
|
|
|
|
|
|
# ---------- HMAC 回调签名(与 Java CallbackSignatureVerifier 对账)----------
|
|
|
|
def compute_signature(secret: str, data: bytes) -> str | None:
|
|
"""HMAC-SHA256(secret, raw bytes) → hex 小写;空 secret 返 None(验签关闭,与 service.py 一致)。"""
|
|
if not secret:
|
|
return None
|
|
return hmac.new(secret.encode("utf-8"), data, hashlib.sha256).hexdigest()
|
|
|
|
|
|
def post_callback(url: str, payload: dict, secret: str, *, opener=None) -> tuple[int, str]:
|
|
"""把 result-out POST 回后端内网回调子路由(签名字节 == 发送字节;HMAC 头与 Java 对账)。
|
|
|
|
Args:
|
|
url: 回调 URL(job.callback.target)。
|
|
payload: result-out dict。
|
|
secret: callback-secret(空=不签)。
|
|
opener: urlopen 注入位(单测桩;默认 urllib.request.urlopen)。
|
|
|
|
Returns:
|
|
(http_status, body_text);网络异常返 (-1, 异常文本),不抛(日志记真因)。
|
|
"""
|
|
data = json.dumps(payload, ensure_ascii=False).encode("utf-8") # 签名字节 == 发送字节(铁律)
|
|
headers = {"Content-Type": "application/json; charset=utf-8", "tenant-id": "0"}
|
|
sig = compute_signature(secret, data)
|
|
if sig:
|
|
headers["X-Callback-Signature"] = sig
|
|
req = urllib.request.Request(url, data=data, method="POST", headers=headers)
|
|
_open = opener or _NO_PROXY_OPENER.open # 默认禁代理(内网回调,经代理会被 fake-ip 拦 502)
|
|
try:
|
|
with _open(req, timeout=30) as resp:
|
|
return resp.status, resp.read().decode("utf-8", "replace")
|
|
except urllib.error.HTTPError as e:
|
|
return e.code, (e.read().decode("utf-8", "replace") if e.fp else "")
|
|
except Exception as e: # noqa: BLE001 — 网络异常返 -1,不抛
|
|
return -1, f"{type(e).__name__}: {e}"
|
|
|
|
|
|
# ---------- best-effort profile 派生(sourceProject 用;派生不出返 None=省略)----------
|
|
|
|
def derive_profile(job: dict, game_dir) -> dict | None:
|
|
"""派生源工程 2.0 profile{tickModel,inputModel,progressModel}(便宜档基线)。
|
|
|
|
A11(切片三)起:sourceProject 由非承重转为 modify 路的硬前置——源不回传落库则 base 版本反查为空、
|
|
无源可改。便宜档当前只产 tap-targets 玩法(点击得分/打地鼠/经营点客等),都是实时循环 + 离散点击 +
|
|
指标进度,故返基线 profile{realtime, discrete-choice, metric}(非伪造:tap-targets 本就这三维)。
|
|
未来便宜档扩到回合制(如 2048)/叙事类时,从 play-spec/产物派生而非固定基线。
|
|
"""
|
|
return {"tickModel": "realtime", "inputModel": "discrete-choice", "progressModel": "metric"}
|
|
|
|
|
|
# ---------- 默认生成核心(lazy import,避免单测触发 tier2/key 装配)----------
|
|
|
|
def _default_run_fn(job: dict):
|
|
"""默认 run_fn:跑 cheap_studio.run_studio(异步)→ 返 (run-summary, game_dir)。
|
|
|
|
懒导入 cheap_studio/cheap_run/_bootstrap——只在真运行时装配 tier2 框架 + key,单测注入 run_fn 不触发。
|
|
"""
|
|
import asyncio
|
|
|
|
import cheap_run
|
|
import cheap_studio
|
|
|
|
game_id = job.get("gameId") or job.get("job_id")
|
|
brief = job.get("brief") or ""
|
|
summary = asyncio.run(cheap_studio.run_studio(game_id, brief))
|
|
return summary, cheap_run.game_dir(game_id)
|
|
|
|
|
|
# ---------- worker 状态 + 有界队列 + 串行消费 ----------
|
|
|
|
class WorkerState:
|
|
"""进程级状态:配置 + 有界队列 + 可注入依赖 + job_id 去重。"""
|
|
|
|
def __init__(self, *, callback_secret: str = "", queue_maxsize: int = DEFAULT_QUEUE_MAXSIZE,
|
|
run_fn=None, send_fn=None, profile_fn=None, classify_fn=None, modify_fn=None, regen_fn=None):
|
|
self.callback_secret = callback_secret
|
|
self.queue: queue.Queue = queue.Queue(maxsize=queue_maxsize)
|
|
self.run_fn = run_fn or _default_run_fn # job -> (summary, game_dir)
|
|
self.send_fn = send_fn or post_callback # (url, payload, secret) -> (status, body)
|
|
self.profile_fn = profile_fn or derive_profile # (job, game_dir) -> profile | None
|
|
self.classify_fn = classify_fn or cheap_classify.classify # (rawText, sourceProject) -> 建议改动(A11 M2 判意图)
|
|
# A11 切片三 M3:确定性类执行器(game_id, sourceProject, modifyPatch) -> 执行结果。
|
|
# None=懒绑 cheap_modify.execute_deterministic_modify(单测注入桩免真 esbuild/Chrome)。
|
|
self.modify_fn = modify_fn
|
|
# A11 切片三 M4:模块重生成执行器(game_id, sourceProject, modifyPatch) -> 执行结果。
|
|
# None=懒绑 cheap_modify.execute_regenerate_modify(单测注入桩免真 LLM/esbuild/Chrome)。
|
|
self.regen_fn = regen_fn
|
|
self._seen: set = set() # job_id 去重(防同 job 重投重跑)
|
|
self._lock = threading.Lock()
|
|
|
|
|
|
def try_enqueue(state: WorkerState, job: dict) -> tuple[bool, int]:
|
|
"""把 job 入有界队列。有容量 → (True,202);队满 → (False,503);同 job_id 重投 → (True,202) 幂等不重入。"""
|
|
jid = job.get("job_id") or job.get("traceId")
|
|
with state._lock:
|
|
if jid is not None and jid in state._seen:
|
|
return True, 202 # 幂等:已受理过,不重复入队
|
|
try:
|
|
state.queue.put_nowait(job)
|
|
except queue.Full:
|
|
return False, 503 # 队满:非 2xx → 执行器 LLM_ERROR(诚实残留,见 KTD4/m1)
|
|
if jid is not None:
|
|
state._seen.add(jid)
|
|
return True, 202
|
|
|
|
|
|
def process_job(state: WorkerState, job: dict) -> dict:
|
|
"""处理一个 job:run_fn 生成 → 组 §6.1 result-out → 向 callback.target 发(带 HMAC)。返回 payload。
|
|
|
|
A11 切片三:有 modifyMode 的 job 走调整路(_process_modify_job),**create/生成路一字不动**——
|
|
modify 路在最前分流,下方生成路对无 modifyMode 的 job 零副作用。
|
|
"""
|
|
trace_id = job.get("traceId") or job.get("job_id")
|
|
|
|
# ── A11 调整路分流(modify):有 modifyMode → deterministic 确定性执行 / regenerate-module 留 M4。──
|
|
modify = parse_modify(job)
|
|
if modify is not None:
|
|
return _process_modify_job(state, job, modify)
|
|
|
|
# ↓↓↓ 以下 create/生成路保持原样、零改动 ↓↓↓
|
|
summary, game_dir = state.run_fn(job)
|
|
|
|
# ── M3b U3:D9 反同质化(vendored dedup)。撞重只告警不阻断(status 仍按九门判);
|
|
# check_similarity 内部已对 FS 异常返 dupHit=None,另包一层 try 兜底——D9 整体失败不得让回调崩。──
|
|
try:
|
|
title = result_out._derive_title(job.get("brief") or "")
|
|
sim = dedup.check_similarity(title=title, theme="generic", trace_id=trace_id)
|
|
if isinstance(summary, dict):
|
|
summary["similarity"] = sim # 随 trace.similarity 落库(result_out._build_trace 读 summary.similarity)
|
|
except Exception as e: # noqa: BLE001 — D9 非承重,失败只记日志、不阻断主链
|
|
log(f"D9 查重失败(省略 similarity,不阻断) trace_id={trace_id}: {type(e).__name__}: {e}")
|
|
|
|
# sourceProject best-effort:profile 可派生才组,否则省略(不伪造、不阻断成功)。
|
|
source_project = None
|
|
try:
|
|
profile = state.profile_fn(job, game_dir)
|
|
if profile:
|
|
source_project = result_out.build_source_project(game_dir, profile)
|
|
except Exception as e: # noqa: BLE001 — sourceProject 非承重,派生失败只记日志、不影响主链
|
|
log(f"sourceProject 派生失败(省略,不阻断)trace_id={trace_id}: {type(e).__name__}: {e}")
|
|
|
|
payload = result_out.build_result_out(job, summary, game_dir, source_project=source_project)
|
|
callback_url = (job.get("callback") or {}).get("target")
|
|
if callback_url:
|
|
status, body = state.send_fn(callback_url, payload, state.callback_secret)
|
|
log(f"回调完成 trace_id={trace_id}, status={payload['status']}, callbackHttp={status}, body={str(body)[:200]}")
|
|
else:
|
|
log(f"job 无 callback.target,跳过回调 trace_id={trace_id}, status={payload['status']}")
|
|
return payload
|
|
|
|
|
|
# ---------- A11 切片三:调整路(modify)处理 ----------
|
|
|
|
def _modify_summary(result: dict, game_id: str) -> dict:
|
|
"""把 execute_deterministic_modify 的执行结果转成 build_result_out 吃的 run-summary 形态。
|
|
|
|
确定性改一次成型(无 resume 轮)→ attempts=1;verdict.pass 取执行结果 status;verdictFull 透传九门 verdict
|
|
供 trace.sevenGateVerdict(后端 D11)。manifest/error 随 summary 带出,由 _process_modify_job 附进 trace。
|
|
"""
|
|
pass_ = result.get("status") == "succeeded"
|
|
verdict_full = result.get("verdict") if isinstance(result.get("verdict"), dict) else None
|
|
failed = []
|
|
if verdict_full:
|
|
guards = verdict_full.get("guards") or {}
|
|
failed = [k for k, g in guards.items() if isinstance(g, dict) and g.get("pass") is False]
|
|
return {
|
|
"ok": pass_, "gameId": game_id, "stage": result.get("stage", "play"), "attempts": 1,
|
|
"verdict": {"pass": pass_, "failedGates": failed},
|
|
"verdictFull": verdict_full,
|
|
}
|
|
|
|
|
|
def _process_modify_job(state: WorkerState, job: dict, modify: dict) -> dict:
|
|
"""A11 调整路:deterministic 走确定性执行(改→重建→九门);regenerate-module 走模块重生成(M4);未知 mode → failed。
|
|
|
|
deterministic:execute_deterministic_modify(注入 modify_fn 可桩)→ 组 result-out(manifest 附进 trace、
|
|
best-effort 带新版 sourceProject 供链式改)→ 回调。regenerate-module:_process_regenerate_job(M4 有界重写
|
|
game-logic.js)。未知 mode:显式 failed + 日志(不静默吞)。
|
|
"""
|
|
trace_id = job.get("traceId") or job.get("job_id")
|
|
mode = modify.get("mode")
|
|
game_id = job.get("gameId") or job.get("job_id")
|
|
|
|
# A11 M4:模块重生成(改玩法)分流——有界单文件重写 game-logic.js,复用 cheap_studio resume + 三层校验。
|
|
if mode == "regenerate-module":
|
|
return _process_regenerate_job(state, job, modify)
|
|
|
|
if mode != "deterministic":
|
|
# 未知 mode 显式 failed,绝不静默吞(错误路红线)。
|
|
why = f"未知 modifyMode={mode}"
|
|
log(f"调整路不支持的 mode → 发兜底 failed trace_id={trace_id}: {why}")
|
|
_send_failed(state, job, why)
|
|
return {"status": "failed", "mode": mode, "error": why}
|
|
|
|
# deterministic:懒绑执行器(单测注入桩免真 esbuild/Chrome)。
|
|
import cheap_run
|
|
|
|
modify_fn = state.modify_fn
|
|
if modify_fn is None:
|
|
import cheap_modify
|
|
modify_fn = cheap_modify.execute_deterministic_modify
|
|
|
|
result = modify_fn(game_id, modify.get("sourceProject"), modify.get("modifyPatch"))
|
|
summary = _modify_summary(result, game_id)
|
|
game_dir = cheap_run.game_dir(game_id)
|
|
|
|
# best-effort:重建成功则把新版源回传(供下一次链式改用 base);派生失败只记日志、不阻断。
|
|
source_project = None
|
|
if result.get("status") == "succeeded":
|
|
try:
|
|
profile = state.profile_fn(job, game_dir)
|
|
if profile:
|
|
source_project = result_out.build_source_project(game_dir, profile)
|
|
except Exception as e: # noqa: BLE001 — sourceProject 非承重,派生失败不阻断回调
|
|
log(f"调整路 sourceProject 派生失败(省略,不阻断)trace_id={trace_id}: {type(e).__name__}: {e}")
|
|
|
|
payload = result_out.build_result_out(job, summary, game_dir, source_project=source_project)
|
|
# 改动清单 manifest 附进 trace(后端 trace 是开放 Map、忽略未知键;此键供前端展示「改了哪几处」+ 排障)。
|
|
manifest = result.get("manifest")
|
|
if manifest is not None:
|
|
payload.setdefault("trace", {})["modifyManifest"] = manifest
|
|
if result.get("error"):
|
|
payload.setdefault("trace", {})["modifyError"] = result["error"]
|
|
|
|
callback_url = (job.get("callback") or {}).get("target")
|
|
if callback_url:
|
|
http_status, body = state.send_fn(callback_url, payload, state.callback_secret)
|
|
log(f"调整路回调完成 trace_id={trace_id}, status={payload['status']}, "
|
|
f"callbackHttp={http_status}, manifest={len(manifest) if manifest else 0} 处, body={str(body)[:160]}")
|
|
else:
|
|
log(f"调整路 job 无 callback.target,跳过回调 trace_id={trace_id}, status={payload['status']}")
|
|
return payload
|
|
|
|
|
|
def _regen_summary(result: dict, game_id: str) -> dict:
|
|
"""把 execute_regenerate_modify 结果转成 build_result_out 吃的 run-summary。
|
|
|
|
复用重写产出的富 run-summary(cost/models/attempts/wallSec/stage/driverType,供 trace 完整);**但
|
|
verdict.pass 以执行器最终判定为准**(= 九门过 ∧ 断言②非目标稳定)——否则「九门过但误伤了非目标」会被
|
|
build_result_out 读 run-summary 原始 verdict.pass=True 误判成功,与执行器 status=failed 矛盾。
|
|
早失败(无重写、result.summary=None)→ 用最小 failed summary 兜底。
|
|
"""
|
|
pass_ = result.get("status") == "succeeded"
|
|
run_summary = result.get("summary")
|
|
verdict_full = result.get("verdict") if isinstance(result.get("verdict"), dict) else None
|
|
failed = []
|
|
if verdict_full:
|
|
guards = verdict_full.get("guards") or {}
|
|
failed = [k for k, g in guards.items() if isinstance(g, dict) and g.get("pass") is False]
|
|
if isinstance(run_summary, dict):
|
|
s = dict(run_summary) # 浅拷贝富字段(cost/models/attempts/wallSec/stage/gameId/driverType)
|
|
else:
|
|
s = {"gameId": game_id, "attempts": 1, "stage": result.get("stage", "code")}
|
|
# verdict 以执行器最终判定为准(含断言②);verdictFull 透传供 trace.sevenGateVerdict(后端 D11)。
|
|
s["verdict"] = {"pass": pass_, "failedGates": failed}
|
|
s["verdictFull"] = verdict_full
|
|
s["gameId"] = game_id
|
|
return s
|
|
|
|
|
|
def _process_regenerate_job(state: WorkerState, job: dict, modify: dict) -> dict:
|
|
"""A11 M4 模块重生成(改玩法):execute_regenerate_modify 有界重写 game-logic.js → 九门 → 改动清单 +
|
|
断言②非目标稳定。manifest 进 trace、best-effort 带新版 sourceProject(同 deterministic 口径)。
|
|
|
|
与 deterministic 路两点差异:① 执行器是 regen_fn(LLM 重写、有 resume 轮),summary 取重写富 run-summary
|
|
(cost/models/attempts 完整)而非 attempts=1 的 _modify_summary;② manifest 是单条 dict(非 list)、含
|
|
changed/untouchedStable。
|
|
"""
|
|
trace_id = job.get("traceId") or job.get("job_id")
|
|
game_id = job.get("gameId") or job.get("job_id")
|
|
|
|
import cheap_run
|
|
|
|
regen_fn = state.regen_fn
|
|
if regen_fn is None:
|
|
import cheap_modify
|
|
regen_fn = cheap_modify.execute_regenerate_modify
|
|
|
|
result = regen_fn(game_id, modify.get("sourceProject"), modify.get("modifyPatch"))
|
|
summary = _regen_summary(result, game_id)
|
|
game_dir = cheap_run.game_dir(game_id)
|
|
|
|
# best-effort:重写成功则把新版源回传(供下一次链式改用 base);派生失败只记日志、不阻断。
|
|
source_project = None
|
|
if result.get("status") == "succeeded":
|
|
try:
|
|
profile = state.profile_fn(job, game_dir)
|
|
if profile:
|
|
source_project = result_out.build_source_project(game_dir, profile)
|
|
except Exception as e: # noqa: BLE001 — sourceProject 非承重,派生失败不阻断回调
|
|
log(f"模块重生成 sourceProject 派生失败(省略,不阻断)trace_id={trace_id}: {type(e).__name__}: {e}")
|
|
|
|
payload = result_out.build_result_out(job, summary, game_dir, source_project=source_project)
|
|
# 改动清单 manifest 附进 trace(后端 trace 是开放 Map、忽略未知键;供前端展示「改了玩法、非目标稳没稳」+ 排障)。
|
|
manifest = result.get("manifest")
|
|
if manifest is not None:
|
|
payload.setdefault("trace", {})["modifyManifest"] = manifest
|
|
if result.get("error"):
|
|
payload.setdefault("trace", {})["modifyError"] = result["error"]
|
|
|
|
callback_url = (job.get("callback") or {}).get("target")
|
|
if callback_url:
|
|
http_status, body = state.send_fn(callback_url, payload, state.callback_secret)
|
|
log(f"模块重生成回调完成 trace_id={trace_id}, status={payload['status']}, callbackHttp={http_status}, "
|
|
f"changed={(manifest or {}).get('changed')}, untouchedStable={(manifest or {}).get('untouchedStable')}, "
|
|
f"body={str(body)[:160]}")
|
|
else:
|
|
log(f"模块重生成 job 无 callback.target,跳过回调 trace_id={trace_id}, status={payload['status']}")
|
|
return payload
|
|
|
|
|
|
def _send_failed(state: WorkerState, job: dict, reason: str) -> None:
|
|
"""生成异常兜底:发 failed result-out,避免任务挂 RUNNING(对执行器 callbackFailed 语义)。"""
|
|
trace_id = job.get("traceId") or job.get("job_id")
|
|
# 用不存在的 game_dir 让 build_result_out 走失败路(bundle 读不到 → status=failed)。
|
|
# failure_reason 传七值枚举 catch-all `llm_error`(worker 异常兜底=生成执行链异常的最近桶):
|
|
# M3a 前误传非枚举值 `generation_failed`,虽经 _map_failure_reason 映射成 llm_error、回调行为正确,
|
|
# 但传真枚举值语义更清、可追溯(错误路红线);具体异常 reason 另由下方日志落痕。
|
|
payload = result_out.build_result_out(job, {}, Path("/nonexistent-amgen"), failure_reason="llm_error")
|
|
url = (job.get("callback") or {}).get("target")
|
|
if url:
|
|
try:
|
|
state.send_fn(url, payload, state.callback_secret)
|
|
except Exception as e: # noqa: BLE001
|
|
log(f"兜底失败回调发送也失败 trace_id={trace_id}: {type(e).__name__}: {e}")
|
|
log(f"已发兜底 failed 回调 trace_id={trace_id}, reason={reason[:200]}")
|
|
|
|
|
|
def _worker_loop(state: WorkerState) -> None:
|
|
"""单 worker 线程:串行消费队列(Chrome/CDP 端口不可并发);异常发兜底 failed 回调、不漏 task_done。"""
|
|
while True:
|
|
job = state.queue.get()
|
|
if job is None: # 停机哨兵
|
|
state.queue.task_done()
|
|
break
|
|
try:
|
|
process_job(state, job)
|
|
except Exception as e: # noqa: BLE001 — 生成/组装异常:发兜底 failed,任务不挂 RUNNING
|
|
log(f"process_job 异常 → 发兜底 failed: {type(e).__name__}: {e}")
|
|
try:
|
|
_send_failed(state, job, f"{type(e).__name__}: {e}")
|
|
except Exception as e2: # noqa: BLE001
|
|
log(f"兜底也异常: {e2}")
|
|
finally:
|
|
state.queue.task_done()
|
|
|
|
|
|
def start_worker(state: WorkerState) -> threading.Thread:
|
|
"""起单 worker daemon 线程(串行消费)。"""
|
|
t = threading.Thread(target=_worker_loop, args=(state,), daemon=True, name="cheap-worker-loop")
|
|
t.start()
|
|
return t
|
|
|
|
|
|
def make_handler(state: WorkerState):
|
|
"""工厂:绑定 state 的 HTTP handler(GET /health 探针;POST /generate 收 §6.1 job-in)。"""
|
|
|
|
class Handler(BaseHTTPRequestHandler):
|
|
def log_message(self, fmt, *args): # 静默默认访问日志(用自定义 log)
|
|
return
|
|
|
|
def _send_json(self, code: int, obj: dict):
|
|
data = json.dumps(obj, ensure_ascii=False).encode("utf-8")
|
|
self.send_response(code)
|
|
self.send_header("Content-Type", "application/json")
|
|
self.send_header("Content-Length", str(len(data)))
|
|
self.end_headers()
|
|
self.wfile.write(data)
|
|
|
|
def do_GET(self):
|
|
if self.path.rstrip("/") in ("/health", "/generate/health"):
|
|
self._send_json(200, {"ok": True, "queued": state.queue.qsize()})
|
|
else:
|
|
self._send_json(404, {"error": "not found"})
|
|
|
|
def _handle_classify(self, raw: bytes):
|
|
"""A11 M2 判意图(同步):读 {rawText, sourceProject} → classify_fn → 返建议改动。
|
|
|
|
不入生成队列、不起 Chrome——判意图是一次轻 LLM 调用(ThreadingHTTPServer 天然并发安全)。
|
|
异常一律兜底为"需确认的 unclear"(绝不静默判成可执行改动),HTTP 仍 200 带兜底建议。
|
|
"""
|
|
body = parse_job(raw)
|
|
if body is None:
|
|
self._send_json(400, {"error": "bad classify json"})
|
|
return
|
|
raw_text = (body.get("rawText") or "").strip()
|
|
if not raw_text:
|
|
self._send_json(400, {"error": "rawText required"})
|
|
return
|
|
try:
|
|
proposal = state.classify_fn(raw_text, body.get("sourceProject"))
|
|
log(f"判意图完成 category={proposal.get('category')}, needsConfirm={proposal.get('needsConfirm')}")
|
|
self._send_json(200, proposal)
|
|
except Exception as e: # noqa: BLE001 — 判意图异常不外抛,返需确认兜底
|
|
log(f"判意图异常 → 兜底 unclear: {type(e).__name__}: {e}")
|
|
self._send_json(200, cheap_classify._safe_fallback(f"endpoint 异常:{type(e).__name__}"))
|
|
|
|
def do_POST(self):
|
|
length = int(self.headers.get("Content-Length", "0") or "0")
|
|
raw = self.rfile.read(length) if length > 0 else b""
|
|
route = self.path.rstrip("/")
|
|
if route == "/classify":
|
|
self._handle_classify(raw)
|
|
return
|
|
if route != "/generate":
|
|
self._send_json(404, {"error": "not found", "path": self.path})
|
|
return
|
|
job = parse_job(raw)
|
|
if job is None:
|
|
self._send_json(400, {"accepted": False, "error": "bad job json"})
|
|
return
|
|
trace_id = job.get("traceId") or job.get("job_id")
|
|
accepted, code = try_enqueue(state, job)
|
|
log(f"收到 job trace_id={trace_id}, templateId={job.get('templateId')}, "
|
|
f"gameId={job.get('gameId')} → {'受理' if accepted else '队满拒'}({code})")
|
|
self._send_json(code, {"accepted": accepted, "job_id": job.get("job_id"),
|
|
"traceId": trace_id, "queued": state.queue.qsize()})
|
|
|
|
return Handler
|
|
|
|
|
|
def start_server(state: WorkerState, *, host: str = "127.0.0.1", port: int = DEFAULT_PORT):
|
|
"""起 HTTP server(serve 在 bg 线程);返回 (server, 实际端口)。port=0 由 OS 分配(单测用)。"""
|
|
server = ThreadingHTTPServer((host, port), make_handler(state))
|
|
actual_port = server.server_address[1]
|
|
threading.Thread(target=server.serve_forever, daemon=True, name="cheap-worker-http").start()
|
|
return server, actual_port
|
|
|
|
|
|
def main():
|
|
parser = argparse.ArgumentParser(description="cheap-worker §6.1 HTTP worker(便宜档生成即服务)")
|
|
parser.add_argument("--port", type=int, default=DEFAULT_PORT, help=f"监听端口(默认 {DEFAULT_PORT})")
|
|
parser.add_argument("--queue-maxsize", type=int, default=DEFAULT_QUEUE_MAXSIZE,
|
|
help=f"有界队列容量(默认 {DEFAULT_QUEUE_MAXSIZE};队满回 503)")
|
|
parser.add_argument("--callback-secret", default="", help="回调 HMAC 共享密钥(空=不签;须与后端同值)")
|
|
args = parser.parse_args()
|
|
|
|
state = WorkerState(callback_secret=args.callback_secret, queue_maxsize=args.queue_maxsize)
|
|
start_worker(state)
|
|
server, port = start_server(state, host="0.0.0.0", port=args.port)
|
|
log(f"cheap-worker service 已启动 → http://0.0.0.0:{port}/generate(队列={args.queue_maxsize}、串行消费;Ctrl-C 退出)")
|
|
try:
|
|
threading.Event().wait() # 阻塞主线程,serve + worker 在 bg 线程
|
|
except KeyboardInterrupt:
|
|
log("收到中断,退出")
|
|
server.shutdown()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|