M1 前置基座(便宜档源回传落库链 + base 源注入 + 血缘):
- derive_profile 兑现(便宜档基线 profile{realtime,discrete-choice,metric})解锁 sourceProject 回传
- worker parse_modify 提取 modify 区(mode/baseVersionId/modifyPatch/sourceProject)
- executor 据 baseVersionId 反查注入 base 源到 HTTP job(镜像 SaaGraphDispatcher.resolveBaseSourceProject)
- 回调落源回填 base_version_id 血缘(landSourceQuietly)
M2 worker 判意图(两段式第一段,NL→建议改动+风险):
- cheap_classify 把用户原话 + base 源可改面(assets.js 资产清单/core.js 集中数值/game-logic 玩法)
判成 {category,mode,target,payload,riskLevel,needsConfirm,clarify}
- mode×category 锁定(LLM 给的 mode 不采信)+ 危险回问硬编进 needsConfirm(只低风险确定性改免确认)
- worker /classify 同步端点(不入生成队列、不起 Chrome)
测试:Python test_a11_m1 5/5 + test_a11_m2_classify 9/9 + 回归全绿;
Java 执行器/回调/源服务 + SAA 回归(本地 maven 独立复跑 46 + 子代理 98)全绿。
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
351 lines
17 KiB
Python
351 lines
17 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):
|
|
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 判意图)
|
|
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。"""
|
|
trace_id = job.get("traceId") or job.get("job_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:
|
|
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
|
|
|
|
|
|
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()
|