#!/usr/bin/env python3 """malfoy-meme-kit LangGraph StateGraph 골격 구조 정본은 pipeline.yaml(roster.json v2의 77노드·113배선·반려 15). 이 파일은 그 yaml을 읽어 그래프를 만든다. - 노드 함수는 스테이지별 stub. 실제 도구 호출 자리는 TODO + 스크립트 경로( 기준) - Send fan-out: 이미지 배치(codex 행 단위) · Flow 배치(잡 단위) · 섹션 워커(fx·lay·pr·pipe·news 그룹) - conditional edge: 테스트 게이트(flow_test) · 판정 루프(컨택트시트·후보 실측·이름 대조) · 재시도 1회(codex timeout·Flow _r2) - interrupt: user_review 노드에서 interrupt() → 승인 / 반려(quote·target_node) → 드라이버가 rewind_to(target) - 동시성: 노드별 슬롯 세마포어(heavy 2 · flow 5 · eleven 2 · codex 4 · matte 2 · 서브에이전트 8) 사용 python3 graph.py --dry-run # 구조(노드·엣지)만 출력. langgraph 없으면 구조만, 있으면 compile 후 get_graph() 요약까지 python3 graph.py --dry-run --mermaid g.mmd # + Mermaid 저장(langgraph 필요) python3 graph.py --simulate # stub으로 끝까지 실행(외부 호출 없음, 반려 자동 승인) python3 graph.py --simulate --reject c_v2_fix # 첫 검토에서 반려 1건을 넣어 되감기 경로까지 시험 필요: PyYAML. 선택: langgraph(>=1.0) — 설치는 /mnt/d/_cache 아래 venv 권장 """ from __future__ import annotations import argparse import re import sys import threading from collections import defaultdict from contextlib import contextmanager from pathlib import Path HERE = Path(__file__).resolve().parent sys.path.insert(0, str(HERE)) from state import MalfoyState, apply_update, initial_state # noqa: E402 START_, END_ = "__start__", "__end__" USER_REVIEW = "user_review" def load_pipeline(path: Path = HERE / "pipeline.yaml") -> dict: try: import yaml except ImportError: sys.exit("PyYAML 필요: pip install pyyaml") return yaml.safe_load(path.read_text(encoding="utf-8")) # ===================================================================== # 1) 구조 계산 — langgraph 없이 동작 (dry-run의 본체) # ===================================================================== def build_structure(P: dict) -> dict: nodes = {n["id"]: n for n in P["nodes"]} groups = P["graph"]["section_groups"] member_of = {m: g for g, gd in groups.items() for m in gd["members"]} def entry(n: str) -> str: # 배선이 실제로 들어가는 그래프 노드 if n in member_of: return f"{member_of[n]}__dispatch" if "fanout" in nodes[n]: return nodes[n]["fanout"]["dispatch"] return n def exit_(n: str) -> str: # 배선이 실제로 나가는 그래프 노드 return nodes[n]["gate"]["pass_node"] if "gate" in nodes[n] else n gnodes: dict[str, dict] = {} # name → {kind, roster, stage} for nid, n in nodes.items(): if "fanout" in n: gnodes[n["fanout"]["dispatch"]] = {"kind": "dispatch", "roster": nid, "stage": n["stage"]} gnodes[n["fanout"]["worker"]] = {"kind": "worker", "roster": nid, "stage": n["stage"]} gnodes[nid] = {"kind": "roster", "roster": nid, "stage": n["stage"]} if "gate" in n: gnodes[n["gate"]["pass_node"]] = {"kind": "pass", "roster": nid, "stage": n["stage"]} for g, gd in groups.items(): gnodes[f"{g}__dispatch"] = {"kind": "group_dispatch", "roster": None, "stage": nodes[gd["members"][0]]["stage"]} gnodes[USER_REVIEW] = {"kind": "interrupt", "roster": None, "stage": "-"} preds: dict[str, set] = defaultdict(set) realized = [] for a, b in P["edges"]: s, d = exit_(a), entry(b) preds[d].add(s) realized.append((a, b, s, d)) static, waiting, cond = [], [], [] for d, ss in preds.items(): if len(ss) == 1: static.append((next(iter(ss)), d)) else: waiting.append((sorted(ss), d)) for nid, n in nodes.items(): if "fanout" in n: f = n["fanout"] cond.append({"src": f["dispatch"], "router": "send_items", "dests": [f["worker"], nid], "send": True}) static.append((f["worker"], nid)) if "gate" in n: g = n["gate"] retry_dest = n["fanout"]["worker"] if g["on_fail"] == "send_failed_items_to_worker" else nid cond.append({"src": nid, "router": "gate", "dests": [retry_dest, g["pass_node"]], "send": retry_dest != nid, "check": g["check"], "max": g["max_retries"]}) for g, gd in groups.items(): cond.append({"src": f"{g}__dispatch", "router": "send_sections", "dests": list(gd["members"]), "send": True}) entries = {entry(n) for n in nodes} starts = sorted(e for e in entries if not preds.get(e)) for s in starts: static.append((START_, s)) for d in P["graph"]["deliverables"]: static.append((exit_(d), USER_REVIEW)) static.append((USER_REVIEW, END_)) has_out = {s for s, _ in static} | {s for ss, _ in waiting for s in ss} | {c["src"] for c in cond} # Send로만 불리는 섹션 멤버는 자기 정적 엣지로 판단 for nid in nodes: e = exit_(nid) if e not in has_out: static.append((e, END_)) has_out.add(e) review_targets = sorted({i["to"] for i in P["interrupts"]}) return {"nodes": nodes, "gnodes": gnodes, "static": static, "waiting": waiting, "cond": cond, "starts": starts, "realized": realized, "entry": entry, "exit": exit_, "groups": groups, "member_of": member_of, "review_targets": review_targets} def print_structure(P: dict, S: dict) -> None: kinds = defaultdict(list) for name, m in S["gnodes"].items(): kinds[m["kind"]].append(name) print(f"# {P['title']}") print(f"roster: 노드 {len(P['nodes'])} · 배선 {len(P['edges'])} · 반려 {len(P['interrupts'])}") print("그래프 노드 " + str(len(S["gnodes"])) + " = " + " · ".join(f"{k} {len(v)}" for k, v in kinds.items())) print("\n## 스테이지별 roster 노드") for st in P["stages"]: print(f"{st['id']} {st['name']} ({len(st['nodes'])}): " + ", ".join(st["nodes"])) print("\n## 가상 노드") for k in ("dispatch", "worker", "pass", "group_dispatch", "interrupt"): print(f"{k} ({len(kinds[k])}): " + ", ".join(kinds[k])) print(f"\n## 정적 엣지 ({len(S['static'])})") by_src = defaultdict(list) for s, d in S["static"]: by_src[s].append(d) for s in sorted(by_src, key=lambda x: (x != START_, x)): print(f"{s} -> {', '.join(by_src[s])}") print(f"\n## 대기 엣지 fan-in ({len(S['waiting'])})") for ss, d in S["waiting"]: print(f"[{', '.join(ss)}] => {d}") print(f"\n## 조건 엣지 ({len(S['cond'])})") for c in S["cond"]: tag = "Send" if c["send"] else "goto" extra = f" # {c['check']} · 재시도 ≤{c['max']}" if c["router"] == "gate" else "" print(f"{c['src']} --{c['router']}/{tag}--> {{{', '.join(c['dests'])}}}{extra}") print(f"\n## interrupt (반려 {len(P['interrupts'])} → 되감기 대상)") for i in P["interrupts"]: print(f"{i['id']} → {i['to']} (entry {S['entry'](i['to'])}) : {i['effect']}") print(f"검토 노드: {USER_REVIEW} ← " + ", ".join(S["exit"](d) for d in P["graph"]["deliverables"])) print("review=all 이면 interrupt_after: " + ", ".join(S["review_targets"])) # 자체 점검 gn = S["gnodes"] missing = [n["id"] for n in P["nodes"] if n["id"] not in gn] unreal = [r for r in S["realized"] if r[2] not in gn or r[3] not in gn] print(f"\n점검: roster 노드 누락 {len(missing)} · 배선 실현 {len(S['realized']) - len(unreal)}/{len(P['edges'])} · 시작 노드 {len(S['starts'])}") # ===================================================================== # 2) 동시성 슬롯 (pipeline.yaml concurrency) # ===================================================================== _SLOTS: dict[str, threading.BoundedSemaphore] = {} def init_slots(P: dict) -> None: for name, c in P["concurrency"].items(): _SLOTS[name] = threading.BoundedSemaphore(int(c["limit"])) @contextmanager def slot(meta: dict): c = meta.get("concurrency") or {} sem = _SLOTS.get(c.get("slot", "")) if sem is None: yield return with sem: # heavy 2 · flow 5 · eleven 2 · codex 4 · matte 2 · subagents 8 yield # ===================================================================== # 3) 스테이지별 stub 핸들러 — 실제 도구 호출 자리는 TODO # 시그니처: handler(meta, state_or_payload) -> 부분 update(dict) # ===================================================================== VERBOSE = False REVIEW_MODE = "deliverables" def todo(meta: dict, what: str) -> None: if VERBOSE: t = meta.get("tool", {}) print(f" TODO[{meta['id']}] {what} :: {t.get('name')} {t.get('script', '')} {t.get('input', '')}".rstrip()) def _ep(node_id: str) -> str | None: m = re.match(r"c_v(\d)", node_id) return f"V{m.group(1)}" if m else None def h_research(meta, s): """S0 리서치. TODO: Aside repl(브라우저) 조사, 웹 조사는 Aside만(firecrawl 금지), curl 다운로드 → brief.*""" todo(meta, "리서치 산출 → " + meta["out"]) key = meta["writes"][0].split(".", 1)[1] return {"brief": {key: meta["out"]}} def h_plan(meta, s): """S0 스토리보드·EDL. TODO: STORYBOARD.md·V1-EDL.md·V3-STORYBOARD.md·*_jobs.json·jsonl 작성(메인/opus 서브)""" todo(meta, "plan.* 작성") return {"plan": {"storyboard": {**(s.get("plan", {}).get("storyboard") or {}), meta["id"]: meta["out"]}}} def h_codex_batch(meta, s): """S1 codex 배치 수집·판정. TODO: _tools/imagegen_runner.py (PARALLEL=4, TIMEOUT=300/600) 결과 확인""" todo(meta, "배치 결과 대조") return {} def h_codex_job(meta, p): """S1 codex 1장. TODO: codex exec $imagegen 1행 실행(-i 참조 첨부), 회수는 move, 있는 파일 건너뜀""" item = p["item"] todo(meta, f"이미지 {item['id']} attempt {p['attempt']}") return {"assets": {"images": {item["id"]: {"id": item["id"], "batch": meta["id"], "status": "ok", "attempt": p["attempt"]}}}} def h_name_check(meta, s): """S1 내용↔파일명 대조. TODO: 컨택트시트 생성 → 육안 대조 → 임시 이름 거쳐 mv 교정 → 재대조""" todo(meta, "컨택트시트 대조") return {} def h_flow_test(meta, s): """S2 360p 테스트. TODO: 06_flow/FLOW-UI.md 절차(Aside), T1·T2, 크레딧 상한 30, scene>0.25 검출""" todo(meta, "360p 테스트 + 게이트") return {} def h_flow_batch(meta, s): """S2 Flow 배치 수집·클립 판정. TODO: 06_flow/<배치>/LOG.md 행 단위 판정(채택/주의/regen)""" todo(meta, "클립 판정") return {} def h_flow_job(meta, p): """S2 Flow 잡 1개. TODO: Aside로 첫 프레임 에셋 업로드(세션 폴더 경유) → 720p 9:16 x1 생성 → 다운로드 → id 대조""" item = p["item"] todo(meta, f"클립 {item['id']} attempt {p['attempt']}") return {"assets": {"clips": {item["id"]: {"id": item["id"], "batch": meta["id"], "mode": item.get("mode", "first_frame"), "verdict": "pending", "regen_count": p["attempt"] - 1}}}} def h_matte_dev(meta, s): """S3 매팅 방식·GPU 포팅. TODO: 05_matte/matte.sh(--gpu), Windows onnxruntime-directml, 캐시는 D:\\_cache""" todo(meta, "매팅 도구 준비") return {"runtime": {"toolchain": {meta["id"]: meta["tool"].get("script")}}} def h_matte_queue(meta, s): """S3 매팅 큐 1줄. TODO: _tools/matte_queue.sh SRC DST — .part.webm → 완료 시 이름 변경, GPU 실패 컷만 CPU""" todo(meta, "매팅 큐") return {"assets": {"mattes": {f"{meta['id']}#batch": {"clip_id": meta["tool"].get("input", ""), "file": meta["out"], "status": "done"}}}} def h_eleven_music(meta, s): """S4 곡 후보 → 실측 → 마스터. TODO: gen_music*.py(동시 2) → analyze*.py(BPM·드롭·브레이크) → (ASR) → build_master*.py""" todo(meta, "곡 후보·실측·마스터") return {"audio": {"tracks": {meta["id"]: {"id": meta["id"], "file": meta["out"], "chosen": True}}}} def h_eleven_tts(meta, s): """S4 VO. TODO: tts.py 오디션 → 테이크 2개씩 → ASR 오인식 _rejected/""" todo(meta, "VO 생성·ASR 필터") return {"audio": {"vo": {meta["id"]: {"line_id": meta["id"], "file": meta["out"], "status": "final"}}}} def h_eleven_sfx(meta, s): """S4 SFX. TODO: sfx.py + sfx_prompts.json (동시 2 권장)""" todo(meta, "SFX 생성") return {"audio": {"sfx": {meta["id"]: {"id": meta["id"], "file": meta["out"]}}}} def h_motion_block(meta, s): """S5 코드 배경 블록. TODO: HyperFrames composition 작성 → hyperframes check(lint 0), layer 변수 bg/fg""" todo(meta, "블록 작성") return {"blocks": {meta["id"]: {"id": meta["id"], "file": meta["tool"].get("script", meta["out"]), "lint": "pass"}}} def h_render_bench(meta, s): """S5 render 가속 실측. TODO: _tools/render-win.sh, 결론 _tools/RENDER-FAST.md(--workers가 효과, GPU 이득 0)""" todo(meta, "render 실측") return {"runtime": {"toolchain": {"render": {"workers": {"emergency": 2, "two_slots": 4, "solo": 6}}}}} def h_semaphore(meta, s): """S6 heavy.sh 세마포어. TODO: _tools/heavy.sh <명령> (flock 2슬롯 + nice, env -u TMPDIR)""" todo(meta, "heavy 세마포어 기동") return {"runtime": {"heavy": {"limit": 2, "in_use": 0, "script": "_tools/heavy.sh"}}} def h_composite(meta, s): """S6 편별 합성. TODO: 배경·전경 패스(HyperFrames, heavy.sh) → 알파 합성(파이썬·ffmpeg) → draft 540p → 컨택트시트 → 최종 1080x1920 -14 LUFS""" todo(meta, "합성·드래프트·최종") ep = _ep(meta["id"]) return {"renders": {ep: {"episode": ep, "final": meta["out"].split(" ")[0], "status": "final", "reviewed": False}}} def h_fix(meta, s): """S6 수정 재합성(2편 하단 자막 제거·-13.3 LUFS·업로드본). TODO: 07_render/v2/tools/comp.py 재실행""" todo(meta, "수정 재합성") return {"renders": {"V2": {"episode": "V2", "status": "final", "reviewed": False, "upload": "07_render/V2_AGENT_SHOWCASE_upload.mp4"}}} def h_section_worker(meta, p): """S7·S8 섹션 워커(서브에이전트 1기 = 출력 파일 1개). TODO: 지시서(SCHEMA) + 근거 파일 읽고 섹션 작성""" todo(meta, "섹션 → " + meta["out"]) return {"docs": {meta["id"]: {"id": meta["id"], "path": meta["out"], "kind": "section"}}} def h_judge(meta, s): """S8 문체 판정(gn-voice-judge). TODO: 등급 + 문장 단위 수정 지시""" todo(meta, "문체 판정") return {"qa": {"verdicts": [{"node": meta["id"], "verdict": "warn", "reason": "C → 문장 수정 지시"}]}} def h_merge(meta, s): """S7·S8 병합. TODO: 섹션 병합 + 누락 0 대조 + 비밀값 검사""" todo(meta, "병합 → " + meta["out"]) return {"docs": {meta["id"]: {"id": meta["id"], "path": meta["out"], "kind": "merged"}}} HANDLERS = { "research": h_research, "plan": h_plan, "codex_batch": h_codex_batch, "name_check": h_name_check, "flow_test": h_flow_test, "flow_batch": h_flow_batch, "matte_dev": h_matte_dev, "matte_queue": h_matte_queue, "eleven_music": h_eleven_music, "eleven_tts": h_eleven_tts, "eleven_sfx": h_eleven_sfx, "motion_block": h_motion_block, "render_bench": h_render_bench, "semaphore": h_semaphore, "composite": h_composite, "fix": h_fix, "section_worker": h_section_worker, "judge": h_judge, "merge": h_merge, } WORKER_HANDLERS = {"codex_batch": h_codex_job, "flow_batch": h_flow_job} # ===================================================================== # 4) 게이트(판정 루프·재시도) — stub 판정은 pipeline.yaml의 관측 횟수(sim_*)를 재현 # ===================================================================== def gate_eval(meta: dict, s: dict) -> dict: nid, g = meta["id"], meta["gate"] done = (s.get("qa", {}).get("retries") or {}).get(nid, 0) rnd = done + 1 want = g.get("sim_retries", 0) retry_items: list[str] = [] if "fanout" in meta and meta["fanout"].get("sim_regen"): want = 1 # 아이템 재생성은 1회만 retry_items = [f"{nid}#{i}" for i in range(meta["fanout"]["sim_regen"])] if done < min(want, g["max_retries"]): return {"qa": {"gates": {nid: {"status": "retry", "round": rnd, "retry_items": retry_items}}, "retries": {nid: 1}, "verdicts": [{"node": nid, "verdict": "retry", "round": rnd, "reason": g["check"]}]}} return {"qa": {"gates": {nid: {"status": "pass", "round": rnd, "retry_items": []}}, "verdicts": [{"node": nid, "verdict": "pass", "round": rnd}]}} def _merge_updates(*ups: dict) -> dict: return apply_update(*ups) def injected(state: dict, config: dict) -> dict: """되감기 때 config로 들어온 반려 규칙·피드백을 State에 반영(없는 id만, upsert라 중복 안전).""" inj = (config or {}).get("configurable", {}).get("inject") or {} have = {r.get("id") for r in state.get("rules") or []} rules = [r for r in inj.get("rules", []) if r.get("id") not in have] if not rules: return {} return {"rules": rules, "feedback": inj.get("feedback", [])} def make_node(meta: dict, P: dict): handler = HANDLERS[meta["handler"]] def fn(state, config): with slot(meta): up = handler(meta, state) ups = [injected(state, config), up, {"runtime": {"trace": [meta["id"]]}}] if "gate" in meta and not (up.get("qa", {}).get("gates") or {}).get(meta["id"]): ups.append(gate_eval(meta, state)) return _merge_updates(*ups) fn.__name__ = meta["id"] fn.__doc__ = handler.__doc__ return fn def make_worker(meta: dict): handler = WORKER_HANDLERS[meta["handler"]] def fn(payload): with slot(meta): up = handler(meta, payload) return _merge_updates(up, {"runtime": {"trace": [meta["fanout"]["worker"]]}}) fn.__name__ = meta["fanout"]["worker"] return fn def make_section(meta: dict): def fn(payload): with slot(meta): up = h_section_worker(meta, payload) return _merge_updates(up, {"runtime": {"trace": [meta["id"]]}}) fn.__name__ = meta["id"] return fn def passthrough(name: str): def fn(state): return {"runtime": {"trace": [name]}} fn.__name__ = name return fn # ===================================================================== # 5) LangGraph 조립 # ===================================================================== def build_graph(P: dict, S: dict): from langgraph.graph import StateGraph from langgraph.types import Send, interrupt g = StateGraph(MalfoyState) nodes = S["nodes"] for name, m in S["gnodes"].items(): if m["kind"] == "roster": meta = nodes[name] if meta.get("section_group"): g.add_node(name, make_section(meta)) # Send로만 호출되는 섹션 워커 else: g.add_node(name, make_node(meta, P)) elif m["kind"] == "worker": g.add_node(name, make_worker(nodes[m["roster"]])) elif m["kind"] in ("dispatch", "group_dispatch", "pass"): g.add_node(name, passthrough(name)) def user_review(state): """사용자 검토 = interrupt. resume 값: {"approved": True} 또는 {"quote", "target_node", ("id","rule")}""" pending = sorted(ep for ep, r in (state.get("renders") or {}).items() if r.get("status") == "final" and not r.get("reviewed")) if REVIEW_MODE == "none": return {"renders": {ep: {"reviewed": True} for ep in pending}, "runtime": {"trace": [USER_REVIEW]}} ans = interrupt({"kind": USER_REVIEW, "pending": pending, "resume": "{'approved': True} | {'quote': ..., 'target_node': ...}"}) trace = {"runtime": {"trace": [USER_REVIEW]}} if ans.get("approved"): return {"renders": {ep: {"reviewed": True, "status": "approved"} for ep in pending}, **trace} fid = ans.get("id") or f"F-{len(state.get('feedback') or []) + 1:02d}" fb = {"id": fid, "quote": ans["quote"], "target_node": ans["target_node"], "status": "open", "rule": ans.get("rule", ans["quote"]), "at": USER_REVIEW} return {"feedback": [fb], "rules": [{"id": f"R-{fid}", "text": fb["rule"], "source": fid}], **trace} g.add_node(USER_REVIEW, user_review) # 정적·대기 엣지 for s, d in S["static"]: g.add_edge(s, d) for ss, d in S["waiting"]: g.add_edge(ss, d) # 조건 엣지: Send fan-out · 게이트 def items_router(nid: str): meta = nodes[nid] over = meta["fanout"]["over"].split(".") # plan.image_batches. def route(state): items = (state.get(over[0]) or {}).get(over[1], {}).get(over[2]) if not items: # stub: 관측 개수만큼 가짜 아이템 items = [{"id": f"{nid}#{i}"} for i in range(meta["fanout"].get("sim_items", 1))] return [Send(meta["fanout"]["worker"], {"item": it, "attempt": 1}) for it in items] or nid route.__name__ = f"send_items_{nid}" return route def sections_router(gname: str): members = S["groups"][gname]["members"] def route(state): return [Send(m, {"section": m, "rules": state.get("rules", [])}) for m in members] route.__name__ = f"send_sections_{gname}" return route def gate_router(nid: str): meta = nodes[nid] gate = meta["gate"] def route(state): gs = (state.get("qa", {}).get("gates") or {}).get(nid, {}) if gs.get("status") != "retry": return gate["pass_node"] if gate["on_fail"] == "send_failed_items_to_worker": return [Send(meta["fanout"]["worker"], {"item": {"id": i}, "attempt": 2}) for i in gs.get("retry_items") or []] \ or gate["pass_node"] return nid # 자기 루프(재대조·재테스트·후보 재생성·드래프트 라운드) route.__name__ = f"gate_{nid}" return route for c in S["cond"]: if c["router"] == "send_items": g.add_conditional_edges(c["src"], items_router(nodes_by_dispatch(nodes)[c["src"]]), c["dests"]) elif c["router"] == "send_sections": g.add_conditional_edges(c["src"], sections_router(c["src"].split("__")[0]), c["dests"]) elif c["router"] == "gate": g.add_conditional_edges(c["src"], gate_router(c["src"]), c["dests"]) return g def nodes_by_dispatch(nodes: dict) -> dict: return {n["fanout"]["dispatch"]: nid for nid, n in nodes.items() if "fanout" in n} def compile_graph(P: dict, S: dict, review: str = "deliverables"): from langgraph.checkpoint.memory import InMemorySaver g = build_graph(P, S) kw = {"checkpointer": InMemorySaver()} # 실사용: SqliteSaver("malfoy.db") — WSL 재부팅 후 재개 if review == "all": kw["interrupt_after"] = S["review_targets"] return g.compile(**kw) # ===================================================================== # 6) 드라이버: 검토·반려·되감기 # ===================================================================== def rewind_to(app, config: dict, target_entry: str, rule: dict | None, fb: dict) -> dict: """체크포인트 이력에서 target_entry가 next인 가장 최근 지점을 찾아 그 checkpoint_id config를 반환(포크 재실행용). 반려 규칙·피드백은 configurable.inject로 실어 보내고, 재실행되는 첫 노드가 State에 넣는다(injected). update_state(as_node=...)를 쓰지 않는 이유: 대기 엣지·Send 라우터를 다시 발화시켜 중복 실행을 만들 수 있음. 되감은 지점 이후 산출은 다시 돈다 → 실제 노드는 '산출 파일 있으면 건너뜀'(멱등)으로 짤 것.""" for snap in app.get_state_history(config): if target_entry in snap.next: inj = {"rules": [rule] if rule else [], "feedback": [{**fb, "status": "rewound"}]} conf = {**snap.config["configurable"], "inject": inj} return {"configurable": conf} raise RuntimeError(f"되감을 체크포인트 없음: {target_entry}") def drive(app, P: dict, S: dict, review: str, reject: str | None, quiet: bool = False) -> dict: from langgraph.types import Command config = {"configurable": {"thread_id": "malfoy-sim"}, "recursion_limit": 1000, "max_concurrency": 8} app.invoke(initial_state(), config) pending_reject = reject rewinds: list[dict] = [] seen_trace = 0 guard = 0 while True: guard += 1 if guard > 200: raise RuntimeError("드라이버 루프 과다") snap = app.get_state(config) if not snap.next: break intr = [i for t in snap.tasks for i in (t.interrupts or ())] trace = (snap.values.get("runtime") or {}).get("trace", []) just_ran = trace[seen_trace:] seen_trace = len(trace) target = S["entry"](pending_reject) if pending_reject else None if target and any(target in h.next for h in app.get_state_history(config)): hit = [i for i in P["interrupts"] if i["to"] == pending_reject] fb = {"id": hit[0]["id"] if hit else "F-sim", "quote": hit[0]["quote"] if hit else "sim reject", "target_node": pending_reject, "status": "open", "effect": hit[0]["effect"] if hit else ""} rule = {"id": f"R-{fb['id']}", "text": fb["effect"] or fb["quote"], "source": fb["id"]} if not quiet: print(f" 반려 {fb['id']} '{fb['quote']}' → rewind {target}") fork = rewind_to(app, config, target, rule, fb) n_before = len((app.get_state(fork).values.get("runtime") or {}).get("trace", [])) rewinds.append({"target": target, "trace_at_fork": n_before}) app.invoke(None, {**fork, "recursion_limit": 1000, "max_concurrency": 8}) # 이후는 스레드 최신(포크 머리)에서 계속, inject는 유지 config = {"configurable": {"thread_id": fork["configurable"]["thread_id"], "inject": fork["configurable"]["inject"]}, "recursion_limit": 1000, "max_concurrency": 8} pending_reject = None seen_trace = 0 continue if intr: # user_review interrupt() → 자동 승인 app.invoke(Command(resume={"approved": True}), config) else: # interrupt_after(review=all) 정지 → 자동 승인 = 그대로 진행 if not quiet: print(f" interrupt_after 정지: {[n for n in just_ran if n in S['review_targets']]} → 승인") app.invoke(None, config) final = dict(app.get_state(config).values) final["_rewinds"] = rewinds return final # ===================================================================== # 7) CLI # ===================================================================== def main() -> int: global VERBOSE, REVIEW_MODE ap = argparse.ArgumentParser(description=__doc__.split("\n")[0]) ap.add_argument("--dry-run", action="store_true", help="구조만 출력하고 종료(외부 호출 없음)") ap.add_argument("--simulate", action="store_true", help="stub으로 끝까지 실행(외부 호출 없음, langgraph 필요)") ap.add_argument("--review", choices=["none", "deliverables", "all"], default="deliverables") ap.add_argument("--reject", help="시뮬레이션에서 첫 검토 때 넣을 반려의 target 노드(예: c_v2_fix)") ap.add_argument("--mermaid", help="compile 후 Mermaid 저장 경로") ap.add_argument("--verbose", action="store_true", help="stub TODO 줄 출력") a = ap.parse_args() VERBOSE = a.verbose REVIEW_MODE = a.review P = load_pipeline() S = build_structure(P) init_slots(P) if a.dry_run or not a.simulate: print_structure(P, S) try: import langgraph # noqa: F401 except ImportError: print("\nlanggraph 미설치 → 구조만 출력하고 종료") return 0 app = compile_graph(P, S, a.review) gg = app.get_graph() names = set(gg.nodes) - {START_, END_} print(f"\n## compile 결과 (langgraph get_graph)") print(f"노드 {len(names)} (구조 계산 {len(S['gnodes'])}, 일치 {names == set(S['gnodes'])}) · 엣지 {len(gg.edges)} " f"(조건 {sum(1 for e in gg.edges if e.conditional)})") if a.mermaid: Path(a.mermaid).write_text(gg.draw_mermaid(), encoding="utf-8") print(f"Mermaid → {a.mermaid}") return 0 app = compile_graph(P, S, a.review) final = drive(app, P, S, a.review, a.reject) trace = final["runtime"]["trace"] ran = set(trace) roster = [n["id"] for n in P["nodes"]] miss = [n for n in roster if n not in ran] print(f"시뮬레이션 완료: 실행 {len(trace)}회 · roster {len(roster) - len(miss)}/{len(roster)} 실행 · 누락 {miss}") reruns = {n: trace.count(n) for n in roster if trace.count(n) > 1} print(f"2회 이상 실행(판정 루프·되감기): {reruns}") print(f"qa.retries: {final['qa']['retries']}") print(f"feedback: {[(f['id'], f['target_node'], f['status']) for f in final.get('feedback', [])]} · rules: {[r['id'] for r in final.get('rules', [])]}") for r in final["_rewinds"]: print(f"되감기: {r['target']} 지점(trace {r['trace_at_fork']}번째)부터 포크 재실행 → 재실행 {len(trace) - r['trace_at_fork']}회") print("renders: " + ", ".join(f"{k}={v.get('status')}" for k, v in sorted(final["renders"].items()))) return 0 if not miss else 1 if __name__ == "__main__": sys.exit(main())