lili 9a6feddfee feat(studio): A11 切片三 M5 三结构断言机器门(改对了没/有没有误伤/血缘可查)
九门只验"仍能玩",验不出"改对了没、有没有误伤";三断言补这层(plan① 线 154)。

cheap_assert.py:
- ① change_applied:改动真生效非 no-op(deterministic=manifest found 且 old≠new / regenerate=changed)
- ② nontarget_stable:base→new 真比对算内容差,非目标文件变=误伤——不信执行器自报标记,
  防执行器 bug 自证(兑现 Opus 评审 P1-2)
- ③ lineage_queryable:baseVersionId 在(worker 侧验必要条件,DB 真可查由 M1 回填 + 后端 e2e)
- three_assertions 合并 verdict

集成进 execute_*:两档加 base_version_id 参 + 附 assertions verdict + 收紧 status
(succeeded ⟺ 九门 ∧ ① ∧ ②;regenerate 原漏判 ①changed 已补);worker 透传 modify.baseVersionId。

测试:test_a11_m5_assertions 12/12 + 全 85 测试独立复跑全绿(含 M3/M4 桩接 kwarg)。
集成端到端 smoke(真九门):ROUND_MS 30000→19000 + baseVersionId=4096
→ status=succeeded、a1/a2/a3 全 true、allPass=true、落盘 19000。

受计费真后端 e2e(经 studio→aigc→worker→回调→新版本+D12计费+血缘)待 mini-desktop 全栈。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-29 02:34:06 -07:00

530 lines
27 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"),
base_version_id=modify.get("baseVersionId")) # M5 三断言③血缘透传
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"),
base_version_id=modify.get("baseVersionId")) # M5 三断言③血缘透传
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()