"""malfoy-meme-kit State 정의 (LangGraph용 TypedDict + Annotated reducer) - 키 정본: state_schema.json (x-reducer 이름이 아래 함수 이름과 1:1) - 도메인 12키: brief · plan · rules · assets · audio · blocks · renders · qa · feedback · budget · runtime · docs - LangGraph는 최상위 키에만 reducer를 건다 → 하위 맵(assets.images 등)은 최상위 reducer가 하위별로 병합 - 병렬 노드(Send 워커·섹션 워커)가 같은 키에 동시에 써도 아래 reducer로 합쳐짐(덮어쓰기 사고 없음) - 이 파일은 langgraph 없이도 import 됨(표준 라이브러리만 사용) """ from __future__ import annotations from typing import Annotated, Any, Literal, TypedDict, get_type_hints # ---------------------------------------------------------------- reducers def _d(x: Any) -> dict: return dict(x) if isinstance(x, dict) else {} def merge_dict(left: dict | None, right: dict | None) -> dict: """얕은 병합, 뒤가 이김. brief·plan.""" out = _d(left) out.update(_d(right)) return out def merge_by_id(left: dict | None, right: dict | None) -> dict: """id → 레코드. 같은 id면 레코드 필드 병합(뒤가 이김). blocks·renders·docs.""" out = _d(left) for k, v in _d(right).items(): if isinstance(v, dict) and isinstance(out.get(k), dict): out[k] = {**out[k], **v} else: out[k] = v return out def merge_nested_by_id(left: dict | None, right: dict | None) -> dict: """하위 맵마다 merge_by_id. assets(images·clips·mattes)·audio(tracks·vo·sfx·beats·candidates).""" out = _d(left) for sub, recs in _d(right).items(): out[sub] = merge_by_id(out.get(sub), recs) return out def upsert_by_id(left: list | None, right: list | dict | None) -> list: """리스트. 같은 id면 필드 병합, 없으면 뒤에 추가. rules·feedback(status 갱신용).""" out = [dict(x) for x in (left or [])] items = [right] if isinstance(right, dict) else (right or []) idx = {x.get("id"): i for i, x in enumerate(out) if x.get("id") is not None} for it in items: key = it.get("id") if key is not None and key in idx: out[idx[key]] = {**out[idx[key]], **it} else: idx[key] = len(out) out.append(dict(it)) return out def merge_qa(left: dict | None, right: dict | None) -> dict: """verdicts 리스트 누적 · retries 합산 · gates 덮어쓰기.""" l, r = _d(left), _d(right) retries = dict(l.get("retries", {})) for k, n in r.get("retries", {}).items(): retries[k] = retries.get(k, 0) + int(n) return { "verdicts": list(l.get("verdicts", [])) + list(r.get("verdicts", [])), "retries": retries, "gates": {**l.get("gates", {}), **r.get("gates", {})}, } def merge_budget(left: dict | None, right: dict | None) -> dict: """원장별 spent 합산 · entries 누적 · start/cap 등 나머지 덮어쓰기. 병렬 Flow 워커 5기가 각자 {"flow_credits": {"spent": 12}}를 내도 합이 맞음.""" out = {k: dict(v) for k, v in _d(left).items()} for ledger, upd in _d(right).items(): cur = out.setdefault(ledger, {}) for f, v in _d(upd).items(): if f == "spent": cur["spent"] = cur.get("spent", 0) + v elif f == "entries": cur["entries"] = list(cur.get("entries", [])) + list(v) else: cur[f] = v return out _RUNTIME_LISTS = ("trace", "incidents", "resume_log") def merge_runtime(left: dict | None, right: dict | None) -> dict: """trace·incidents·resume_log 누적 · slots·heavy·toolchain 병합.""" out = _d(left) for k, v in _d(right).items(): if k in _RUNTIME_LISTS: out[k] = list(out.get(k, [])) + list(v) elif isinstance(v, dict): out[k] = {**_d(out.get(k)), **v} else: out[k] = v return out REDUCERS = {f.__name__: f for f in (merge_dict, merge_by_id, merge_nested_by_id, upsert_by_id, merge_qa, merge_budget, merge_runtime)} # ---------------------------------------------------------------- 레코드 타입 (state_schema.json $defs와 같은 모양) Episode = Literal["V1", "V2", "V3", "V4", "V5"] class ImageRec(TypedDict, total=False): id: str; batch: str; used_in: list[str]; model: str; size: str; refs: list[str] prompt: str; file: str; final_file: str status: Literal["pending", "ok", "timeout", "failed", "renamed"]; elapsed_s: float; attempt: int class ClipRec(TypedDict, total=False): id: str; batch: str; mode: Literal["first_frame", "first_last", "ingredients"] first_frame: str; last_frame: str; refs: list[str]; seconds: int; resolution: str; aspect: str prompt: str; result_file: str; credits: int verdict: Literal["채택", "채택(주의)", "regen", "불합격", "pending"]; verdict_note: str regen_count: int; blue_bg: bool class MatteRec(TypedDict, total=False): clip_id: str; file: str; qc: str; backend: Literal["gpu", "cpu", "hybrid"] sec_per_frame: float; status: Literal["queued", "part", "done", "failed"] class Track(TypedDict, total=False): id: str; episode: Episode; file: str; source: str; request: str; bpm: float; lufs: float trim_from_raw_s: float; chosen: bool; lang: Literal["instrumental", "ko", "en"]; metrics: dict class VoTake(TypedDict, total=False): line_id: str; episode: Episode; voice_name: str; voice_id: str; model: str text_with_tags: str; file: str; asr: str; status: Literal["final", "alt", "rejected"] class SfxRec(TypedDict, total=False): id: str; prompt: str; duration: float; file: str class BeatGrid(TypedDict, total=False): file: str; bpm: float; beat_period_s: float; drop_s: float; big_hits_s: list[float] bars: int; source: Literal["plan", "measured"] class Block(TypedDict, total=False): id: str; file: str; worker: str; vars: dict; techniques: list[str]; lint: str class RenderRec(TypedDict, total=False): episode: Episode; passes: list[str]; drafts: list[str]; final: str; upload: str; contact: str duration_s: float; lufs: float; round: int status: Literal["building", "draft", "final", "approved", "rejected"]; reviewed: bool class Verdict(TypedDict, total=False): node: str; target: str; verdict: Literal["pass", "fail", "retry", "warn"]; reason: str; round: int; evidence: str class GateState(TypedDict, total=False): status: Literal["pass", "retry", "blocked"]; round: int; retry_items: list[str] class Feedback(TypedDict, total=False): id: str; quote: str; target_node: str; effect: str; rule: str status: Literal["open", "rewound", "resolved", "rejected"]; at: str class Rule(TypedDict, total=False): id: str; text: str; source: str; scope: str class Ledger(TypedDict, total=False): start: float; cap: float | str; spent: float; entries: list[dict] class Doc(TypedDict, total=False): id: str; path: str; kind: Literal["section", "merged", "verdict", "asset"]; lines: int; status: str # ---------------------------------------------------------------- 도메인 하위 구조 class Brief(TypedDict, total=False): reference: dict; meme_research: str; flow_research: str; pinterest: dict class Plan(TypedDict, total=False): episodes: dict[str, dict]; storyboard: dict; edl: dict beat_grid: dict[str, BeatGrid] image_batches: dict[str, list[dict]] # 배치 노드 id → jsonl 행 flow_jobs: dict[str, list[dict]] # 배치 노드 id → *_jobs.json 잡 credit_plan: dict class Assets(TypedDict, total=False): images: dict[str, ImageRec]; clips: dict[str, ClipRec]; mattes: dict[str, MatteRec] class Audio(TypedDict, total=False): tracks: dict[str, Track]; vo: dict[str, VoTake]; sfx: dict[str, SfxRec] beats: dict[str, BeatGrid]; candidates: dict[str, Track] class QA(TypedDict, total=False): verdicts: list[Verdict]; retries: dict[str, int]; gates: dict[str, GateState] class Budget(TypedDict, total=False): flow_credits: Ledger; eleven_credits: Ledger; codex: Ledger class Runtime(TypedDict, total=False): slots: dict[str, int]; heavy: dict; toolchain: dict incidents: list[dict]; resume_log: list[dict]; trace: list[str] # ---------------------------------------------------------------- 그래프 State class MalfoyState(TypedDict, total=False): brief: Annotated[Brief, merge_dict] plan: Annotated[Plan, merge_dict] rules: Annotated[list[Rule], upsert_by_id] assets: Annotated[Assets, merge_nested_by_id] audio: Annotated[Audio, merge_nested_by_id] blocks: Annotated[dict[str, Block], merge_by_id] renders: Annotated[dict[str, RenderRec], merge_by_id] qa: Annotated[QA, merge_qa] feedback: Annotated[list[Feedback], upsert_by_id] budget: Annotated[Budget, merge_budget] runtime: Annotated[Runtime, merge_runtime] docs: Annotated[dict[str, Doc], merge_by_id] STATE_KEYS = list(MalfoyState.__annotations__) REDUCER_OF = {k: t.__metadata__[0] for k, t in get_type_hints(MalfoyState, include_extras=True).items()} def apply_update(*updates: dict) -> dict: """부분 update 여러 개를 State reducer로 합침(한 노드가 여러 조각을 낼 때).""" out: dict = {} for u in updates: for k, v in (u or {}).items(): out[k] = REDUCER_OF[k](out.get(k), v) if k in out else v return out def initial_state() -> MalfoyState: """실제 실행 시작값: 예산·슬롯·heavy 세마포어 (pipeline.yaml budgets·concurrency와 같은 값).""" return { "budget": { "flow_credits": {"start": 1050, "spent": 0, "entries": []}, "eleven_credits": {"start": 0, "spent": 0, "entries": []}, "codex": {"spent": 0, "entries": []}, }, "runtime": { "slots": {"subagents": 8, "doc_subagents": 12, "heavy": 2, "matte_lines": 2, "flow": 5, "eleven": 2, "codex": 4}, "heavy": {"limit": 2, "in_use": 0, "script": "_tools/heavy.sh"}, "trace": [], }, "rules": [], "feedback": [], }