固化 Match-3 生产者、视觉、音频与双 Judge 证据闭包。 将《山海行纪》r1.1 绑定新的不可变 release,并以生产预检现场核验 bundle、Registry/2 和 25 项 Writer 快照。 同步地图1平衡锁值、跨游戏回归修复、验收契约与 SoT 证据。
714 lines
38 KiB
Python
714 lines
38 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 os
|
||
import queue
|
||
import threading
|
||
import urllib.error
|
||
import urllib.request
|
||
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
||
from pathlib import Path
|
||
|
||
import cheap_classify
|
||
import cheap_otlp_sink # 面四三跳传播:worker span 作用域 + W3C carrier 提取(顶层零重依赖,otel 惰性)
|
||
import cheap_verify
|
||
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)
|
||
|
||
|
||
def _mask_token(tok) -> str:
|
||
"""token 日志脱敏(WU2 §3.8):前后各留 4 位、中间省略;过短/空则整体隐藏。绝不整条打 userToken。"""
|
||
if not tok:
|
||
return "<none>"
|
||
s = str(tok)
|
||
return "****" if len(s) <= 8 else f"{s[:4]}…{s[-4:]}"
|
||
|
||
|
||
def _source_project_scaffold_template(source_project) -> str | None:
|
||
"""从受信 sourceProject 血缘读取模板路由;历史 2.0 缺字段时返回 None,绝不从 brief/genre 猜。"""
|
||
try:
|
||
value = json.loads(source_project) if isinstance(source_project, str) else source_project
|
||
except (TypeError, ValueError):
|
||
return None
|
||
if not isinstance(value, dict):
|
||
return None
|
||
template = value.get("scaffoldTemplate")
|
||
return template if template in result_out._SOURCE_PROJECT_SCAFFOLDS else None
|
||
|
||
|
||
def _source_project_origin(source_project) -> tuple[str | None, str | None]:
|
||
"""读取并复核 base SourceProject 的不可变原题面;缺失或 hash 不自洽时返回空。"""
|
||
try:
|
||
value = json.loads(source_project) if isinstance(source_project, str) else source_project
|
||
except (TypeError, ValueError):
|
||
return None, None
|
||
if not isinstance(value, dict):
|
||
return None, None
|
||
brief = value.get("originBrief")
|
||
brief_hash = value.get("originBriefHash")
|
||
if (not isinstance(brief, str) or not brief.strip() or not isinstance(brief_hash, str)
|
||
or hashlib.sha256(brief.encode("utf-8")).hexdigest() != brief_hash):
|
||
return None, None
|
||
return brief, brief_hash
|
||
|
||
|
||
def _accepted_scaffold_template(summary) -> str | None:
|
||
"""只从可信 v3 封存对象恢复 create 模板;未接受、旧模式或字段缺失均不产血缘。"""
|
||
acceptance = summary.get("acceptanceV3") if isinstance(summary, dict) else None
|
||
if not isinstance(acceptance, dict) or not cheap_verify.is_v3_publishable(acceptance):
|
||
return None
|
||
template = acceptance.get("templateRoute")
|
||
return template if template in result_out._SOURCE_PROJECT_SCAFFOLDS else None
|
||
|
||
|
||
def _source_project_build_dir(summary: dict, game_dir) -> Path:
|
||
"""active v3 从验收同源 staged/src 组源;legacy 继续读取既有 release 工作目录。"""
|
||
acceptance = summary.get("acceptanceV3") if isinstance(summary, dict) else None
|
||
if isinstance(acceptance, dict) and cheap_verify.is_v3_publishable(acceptance):
|
||
game_id = str(acceptance.get("gameId") or "")
|
||
if cheap_verify._is_safe_game_id_v3(game_id):
|
||
return Path(game_dir).parent / "_wg1-gen" / game_id
|
||
return Path(game_dir)
|
||
|
||
|
||
# ---------- §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 只消费不反查)
|
||
}
|
||
|
||
|
||
def _freeze_expected_acceptance_mode(job: dict) -> str | None:
|
||
"""在 HTTP 入队边界冻结生产验收代际,防 run-summary 全 marker 丢失后被当成 legacy。
|
||
|
||
运行时配置是主权威;若配置仍是历史模式但可信请求显式要求 v3/v3_shadow,则按请求提高门槛。
|
||
只冻结新协议两态,真 legacy 保持 None,继续走原兼容窗口。
|
||
"""
|
||
requested = job.get("acceptanceMode")
|
||
try:
|
||
configured = (cheap_verify._acceptance_v3_cfg() or {}).get("mode")
|
||
except Exception: # noqa: BLE001 —— 配置读取失败时仍可尊重可信请求;无请求则不在此猜代际
|
||
configured = None
|
||
if configured in result_out._V3_MODES:
|
||
return configured
|
||
return requested if requested in result_out._V3_MODES else None
|
||
|
||
|
||
# ---------- 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")
|
||
# 真后端来的 gameId 是 Java Long → JSON 数字 → Python int;cheap_run 的 subprocess shell-out(scaffold/build/stage/smoke/play)
|
||
# 与路径段都要求 str,int 会抛 TypeError(expected str, not int) → 生成静默判 llm_error。此处在 job 消费边界统一归一为 str。
|
||
game_id = str(game_id) if game_id is not None else None
|
||
brief = job.get("brief") or ""
|
||
summary = asyncio.run(cheap_studio.run_studio(
|
||
game_id, brief, acceptance_task_trace_id=job.get("traceId") or job.get("job_id"),
|
||
# 兼容进程内入口只透传可信 job 显式值;None 保持普通 puzzle 不自动选择 profile。
|
||
interaction_profile_id=job.get("interactionProfileId")))
|
||
return summary, cheap_run.game_dir(game_id)
|
||
|
||
|
||
def _service_run_fn(job: dict):
|
||
"""默认 run_fn(阶段一②):驱动 cheap Service /chat 跑一局(单 POST + 洋葱内续修),返 (run-summary, game_dir)。
|
||
|
||
取代旧 _default_run_fn 的进程内 run_studio:续修/门判/软预算全在 Service 端 middleware(阶段一① 复用),
|
||
消费方只发一次 POST。懒导入 cheap_service_driver(牵 httpx/tier2),单测注入 run_fn 不触发。
|
||
"""
|
||
import asyncio
|
||
|
||
# 先接 sys.path 兜底(tier2/gen-worker → 顶层包 worker 可 import),再导 driver——
|
||
# cheap_service_driver 运行期 `from worker import client` 依赖它;CLI 路(cheap_studio)自带
|
||
# _bootstrap、本门面路此前漏接,S1 真验实证 ModuleNotFoundError 兜底 failed 回调(2026-07-03)。
|
||
import _bootstrap # noqa: F401
|
||
import cheap_service_driver
|
||
|
||
base_url = os.environ.get("CHEAP_SERVICE_URL", "http://127.0.0.1:8300")
|
||
# WU2 §3.7 F:整条 job(含 §6.1 追加的 userToken)原样透传给 driver;driver 从 job.get("userToken") 取
|
||
# per-user token 装 new-api 凭据、缺失回落全局 env key(见 cheap_service_driver._resolve_key)。
|
||
# 此处不拆字段单独传递,避免冗余管道——透传整 job 即满足「worker 取 job token」。
|
||
return asyncio.run(cheap_service_driver.drive_cheap_generation(job, base_url=base_url))
|
||
|
||
|
||
# ---------- 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 _service_run_fn # job -> (summary, game_dir)(阶段一②:默认驱动 Service /chat)
|
||
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")
|
||
expected_acceptance_mode = job.get("_expectedAcceptanceMode")
|
||
|
||
# ── A11 调整路分流(modify):有 modifyMode → deterministic 确定性执行 / regenerate-module 留 M4。──
|
||
modify = parse_modify(job)
|
||
if modify is not None:
|
||
return _process_modify_job(state, job, modify)
|
||
|
||
# ↓↓↓ 以下 create/生成路保持原样、零改动 ↓↓↓
|
||
# 面四断点①(span):在入站 W3C context 下开 worker 这跳的 SERVER span,包住真生成调用(run_fn 内经
|
||
# driver 打 :8300)。span attach 为当前 context 后,driver 的 httpx inject 与写工厂 sidecar 的 traceparent
|
||
# 都据它继承,使 cheap-service 生成 span 挂到本 span(→ Java span)下当子 span,串成单条 trace 树。
|
||
# 业务 traceId 并存(§10 决策5):gameId 作 conversation.id 与 cheap-service 生成 span 同键、业务 traceId 另存属性。
|
||
# worker_span_scope 全 best-effort:OTLP 未开 / otel 不可用时退化成透传入站 context 或什么都不做,绝不咬生成。
|
||
_game_id = job.get("gameId") or job.get("job_id")
|
||
with cheap_otlp_sink.worker_span_scope(job.get("_otelCarrier"),
|
||
conversation_id=_game_id, business_trace_id=trace_id):
|
||
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:
|
||
accepted_template = _accepted_scaffold_template(summary)
|
||
origin_brief = str(job.get("brief") or "") if accepted_template else None
|
||
origin_brief_hash = (hashlib.sha256(origin_brief.encode("utf-8")).hexdigest()
|
||
if origin_brief else None)
|
||
source_project = result_out.build_source_project(
|
||
_source_project_build_dir(summary, game_dir), profile,
|
||
scaffold_template=accepted_template, origin_brief=origin_brief,
|
||
origin_brief_hash=origin_brief_hash)
|
||
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,
|
||
expected_acceptance_mode=expected_acceptance_mode)
|
||
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]
|
||
summary = {
|
||
"ok": pass_, "gameId": game_id, "stage": result.get("stage", "play"), "attempts": 1,
|
||
"verdict": {"pass": pass_, "failedGates": failed},
|
||
"verdictFull": verdict_full,
|
||
}
|
||
acceptance = result.get("acceptance") if isinstance(result.get("acceptance"), dict) else {}
|
||
sealed = acceptance.get("acceptanceV3")
|
||
if isinstance(sealed, dict):
|
||
# deterministic 执行器把统一验收收在 result.acceptance;回调摘要必须恢复完整封存对象与
|
||
# 四个平铺身份,禁止因包装层丢字段而静默退回 legacy verdict。
|
||
import cheap_studio # noqa: PLC0415
|
||
|
||
first_pass = acceptance.get("acceptanceV3FirstPass")
|
||
summary = cheap_studio.apply_acceptance_v3(
|
||
summary, sealed, first_pass=first_pass if isinstance(first_pass, dict) else None)
|
||
# 三断言/执行状态仍是外层硬门。验收 accept 但改动 no-op 或越界时,故意制造与 compatibility.ok
|
||
# 的不一致,让 result_out fail-closed,不能只凭玩法仍可玩就把失败的修改发布出去。
|
||
summary["ok"] = pass_
|
||
else:
|
||
# v3 入口失败虽无 canonical payload,也要保留代际标记,确保 result_out 不回落旧九门。
|
||
for key in ("acceptanceVersion", "accepted", "ok", "publishFrozen", "failureReason", "failureLayer"):
|
||
if key in acceptance:
|
||
summary[key] = acceptance[key]
|
||
return summary
|
||
|
||
|
||
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")
|
||
expected_acceptance_mode = job.get("_expectedAcceptanceMode")
|
||
mode = modify.get("mode")
|
||
game_id = job.get("gameId") or job.get("job_id")
|
||
# 真后端来的 gameId 是 Java Long → JSON 数字 → Python int;cheap_run 的 subprocess shell-out(scaffold/build/stage/smoke/play)
|
||
# 与路径段都要求 str,int 会抛 TypeError(expected str, not int) → 生成静默判 llm_error。此处在 job 消费边界统一归一为 str。
|
||
game_id = str(game_id) if game_id is not None else None
|
||
|
||
# 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"),
|
||
task_trace_id=trace_id) # M5 三断言③血缘 + v3 当前任务绑定透传
|
||
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:
|
||
origin_brief, origin_brief_hash = _source_project_origin(modify.get("sourceProject"))
|
||
source_project = result_out.build_source_project(
|
||
_source_project_build_dir(summary, game_dir), profile,
|
||
scaffold_template=_source_project_scaffold_template(modify.get("sourceProject")),
|
||
origin_brief=origin_brief, origin_brief_hash=origin_brief_hash)
|
||
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,
|
||
expected_acceptance_mode=expected_acceptance_mode)
|
||
# 改动清单 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
|
||
# v3 不再读取 verdict.pass 作为发布权威,必须把模块重生成外层三断言终态同步写回 ok。
|
||
# 当玩法验收本身 accept、但目标 no-op 或误伤非目标时,ok=false 会与封存 compatibility.ok=true
|
||
# 形成显式矛盾,由 result_out fail-closed;否则旧 run_summary.ok=true 会把执行器 failed 翻成 succeeded。
|
||
s["ok"] = pass_
|
||
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")
|
||
expected_acceptance_mode = job.get("_expectedAcceptanceMode")
|
||
game_id = job.get("gameId") or job.get("job_id")
|
||
# 真后端来的 gameId 是 Java Long → JSON 数字 → Python int;cheap_run 的 subprocess shell-out(scaffold/build/stage/smoke/play)
|
||
# 与路径段都要求 str,int 会抛 TypeError(expected str, not int) → 生成静默判 llm_error。此处在 job 消费边界统一归一为 str。
|
||
game_id = str(game_id) if game_id is not None else None
|
||
|
||
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"),
|
||
task_trace_id=trace_id) # M5 三断言③血缘 + v3 当前任务绑定透传
|
||
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:
|
||
origin_brief, origin_brief_hash = _source_project_origin(modify.get("sourceProject"))
|
||
source_project = result_out.build_source_project(
|
||
_source_project_build_dir(summary, game_dir), profile,
|
||
scaffold_template=_source_project_scaffold_template(modify.get("sourceProject")),
|
||
origin_brief=origin_brief, origin_brief_hash=origin_brief_hash)
|
||
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,
|
||
expected_acceptance_mode=expected_acceptance_mode)
|
||
# 改动清单 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",
|
||
expected_acceptance_mode=job.get("_expectedAcceptanceMode"))
|
||
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
|
||
# 生产 HTTP 是三条执行路的共同可信入口;覆盖任何同名入站私有字段,避免调用方自行降级。
|
||
job["_expectedAcceptanceMode"] = _freeze_expected_acceptance_mode(job)
|
||
# 面四断点①(capture):抓入站 W3C traceparent 存进 job,供 worker 线程建 worker span(见 process_job)。
|
||
# do_POST 只入队即 202 投递握手,真生成在 worker 线程发生;self.headers 仅这里可得,故此处只抓 header、
|
||
# 不建 span(span 建在真处理该 job 的 worker 线程,才能覆盖「调 cheap-service 生成」这一跳并当下游的父)。
|
||
# 纯增量 best-effort:无 traceparent / 抓取异常都不影响入队与生成(设计 §8 旁路铁律)。
|
||
try:
|
||
_tp = self.headers.get("traceparent")
|
||
if _tp:
|
||
_carrier = {"traceparent": _tp}
|
||
_ts = self.headers.get("tracestate")
|
||
if _ts:
|
||
_carrier["tracestate"] = _ts
|
||
job["_otelCarrier"] = _carrier
|
||
except Exception: # noqa: BLE001 —— 抓 traceparent 是旁路,失败绝不咬入队/生成
|
||
pass
|
||
trace_id = job.get("traceId") or job.get("job_id")
|
||
accepted, code = try_enqueue(state, job)
|
||
# WU2 §3.7/§3.8:记 userToken 是否随 job 带来(脱敏)——便于排「member create 该带 token 却没带」。
|
||
log(f"收到 job trace_id={trace_id}, templateId={job.get('templateId')}, "
|
||
f"gameId={job.get('gameId')}, userToken={_mask_token(job.get('userToken'))} "
|
||
f"→ {'受理' 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()
|