"""CLI entrypoint: python -m pipeline.cli [args] All SPEC_NEW §4 commands exist. Commands for phases not yet built exit loudly with a clear message instead of pretending to work (prime directive). Prefer the web UI for day-to-day use: `python -m pipeline.cli serve`. The CLI and the UI share pipeline/flow.py, so they never drift. The v2 video flow (SPEC_NEW §11.3): seed-test-video → create a video + 2 hook takes (different seeds) generate → generate the hook takes pick-seed → QC the takes, lock the winning seed, create body+cta generate → generate body+cta on the locked seed assemble → audio-aware smart stitch → final mp4 """ from __future__ import annotations import argparse import json import sys import traceback from typing import Any, Callable from . import db, flow # ── Job handler registry (used by `tick`) ──────────────────────────────────── def _handle_dummy(payload: dict[str, Any]) -> None: if payload.get("explode"): raise RuntimeError("dummy job asked to explode") def _handle_generate_segment(payload: dict[str, Any]) -> None: from . import veo veo.generate_and_download(payload["segment_id"]) def _handle_qc_segment(payload: dict[str, Any]) -> None: from . import qc qc.run_qc(payload["segment_id"]) def _handle_publish_video(payload: dict[str, Any]) -> None: from . import meta_client meta_client.publish(payload["video_id"]) JOB_HANDLERS: dict[str, Callable[[dict[str, Any]], None]] = { "dummy": _handle_dummy, "generate_segment": _handle_generate_segment, "qc_segment": _handle_qc_segment, "publish_video": _handle_publish_video, } # ── Commands ───────────────────────────────────────────────────────────────── def cmd_migrate(_: argparse.Namespace) -> int: applied = db.migrate() print("applied: " + ", ".join(applied) if applied else "up to date — nothing to apply") return 0 def cmd_tick(args: argparse.Namespace) -> int: processed = 0 while args.max_jobs is None or processed < args.max_jobs: with db.connect() as conn: job = db.claim_next_job(conn) if job is None: break job_id, job_type = str(job["id"]), job["type"] handler = JOB_HANDLERS.get(job_type) try: if handler is None: raise RuntimeError(f"no handler registered for job type '{job_type}'") handler(job["payload"]) except Exception as exc: # noqa: BLE001 — failures must land in jobs.error, never vanish err = f"{type(exc).__name__}: {exc}\n{traceback.format_exc(limit=5)}" with db.connect() as conn: db.fail_job(conn, job_id, err) print(f"job {job_id} ({job_type}) FAILED: {exc}", file=sys.stderr) else: with db.connect() as conn: db.complete_job(conn, job_id) print(f"job {job_id} ({job_type}) done") processed += 1 print(f"tick: processed {processed} job(s)") return 0 def cmd_seed_avatar(args: argparse.Namespace) -> int: from psycopg.types.json import Json persona = json.loads(args.persona) with db.connect() as conn: row = conn.execute( "insert into avatars (name, ref_image_path, persona) values (%s, %s, %s) returning id", (args.name, args.ref_image, Json(persona)), ).fetchone() print(f"avatar {row['id']} created") # type: ignore[index] return 0 def cmd_seed_test_video(args: argparse.Namespace) -> int: script = json.loads(__import__("pathlib").Path(args.script_json).read_text()) try: video_id = flow.create_video( args.avatar, script["hook"], script.get("body_segments", []), script.get("cta", ""), angle=script.get("angle", "test"), setting=script.get("setting"), takes=args.takes, single_take=args.single_take, seed=args.seed, ) except Exception as exc: # noqa: BLE001 — bad script/avatar → fail loud print(f"could not create video: {exc}", file=sys.stderr) return 1 if args.single_take: print(f"video {video_id} created (single-take) — next: generate {video_id}, then assemble {video_id}") else: print(f"video {video_id} created with hook takes — next: generate {video_id}, then pick-seed {video_id}") return 0 def cmd_generate(args: argparse.Namespace) -> int: try: flow.generate_pending(args.video_id, on_event=lambda m: print(m)) except Exception as exc: # noqa: BLE001 — fail loud, stop print(f"generation FAILED: {exc}", file=sys.stderr) return 1 body = db.fetch_all( "select 1 from segments where video_id = %s and kind in ('body','cta')", (args.video_id,) ) if not body: db.execute("update videos set status = 'pending_seed_pick' where id = %s", (args.video_id,)) print(f"hook takes ready — next: pick-seed {args.video_id}") else: print(f"all segments generated — next: assemble {args.video_id}") return 0 def cmd_pick_seed(args: argparse.Namespace) -> int: try: res = flow.pick_seed(args.video_id, take=args.take) except Exception as exc: # noqa: BLE001 print(f"pick-seed FAILED: {exc}", file=sys.stderr) return 1 print(f"locked seed {res['seed']} (hook take #{res['take']}); superseded {res['superseded']} take(s); " f"{res['body_count']} body+cta segment(s) created — next: generate {args.video_id}") return 0 def cmd_qc(args: argparse.Namespace) -> int: from . import qc result = qc.run_qc(args.segment_id) print(json.dumps(result, ensure_ascii=False, indent=2, default=str)) return 0 if result["verdict"] in ("pass", "flag") else 1 def cmd_assemble(args: argparse.Namespace) -> int: from . import assemble print(f"assembled: {assemble.assemble(args.video_id)}") return 0 def cmd_serve(args: argparse.Namespace) -> int: import uvicorn from .webapp import app print(f"UGC Pipeline UI → http://{args.host}:{args.port} (Ctrl-C to stop)") uvicorn.run(app, host=args.host, port=args.port, log_level="info") return 0 def cmd_ingest(args: argparse.Namespace) -> int: from . import ingest try: n = ingest.run(days=args.days) except Exception as exc: # noqa: BLE001 — fail loud print(f"ingest FAILED: {exc}", file=sys.stderr) return 1 print(f"ingest: upserted {n} metrics_daily row(s) via Windsor.ai (last {args.days}d)") return 0 def cmd_strategize(args: argparse.Namespace) -> int: from . import agents try: briefs = agents.strategize(days=args.days, n=args.n) except Exception as exc: # noqa: BLE001 — fail loud print(f"strategize FAILED: {exc}", file=sys.stderr) return 1 print(f"strategist produced {len(briefs)} brief(s):") for b in briefs: hyp = (b.get("hypothesis") or "").strip() print(f" [{b['mode']}] {b['angle']} — avatar {b.get('avatar_id')}" f"{' product ' + str(b['product_id']) if b.get('product_id') else ''}" f"{' :: ' + hyp[:80] if hyp else ''}") return 0 def cmd_write_scripts(args: argparse.Namespace) -> int: from . import agents pending = db.fetch_all( "select id from briefs where status = 'pending' order by created_at" ) if args.limit: pending = pending[: args.limit] if not pending: print("no pending briefs — run `strategize` first") return 0 made = 0 for b in pending: try: res = agents.write_script(str(b["id"])) except Exception as exc: # noqa: BLE001 — skip the bad brief, keep going print(f" brief {b['id']} FAILED: {exc}", file=sys.stderr) continue made += 1 print(f" brief {b['id']} -> video {res['video_id']} (script {res['script_id']})") print(f"write-scripts: {made}/{len(pending)} brief(s) scripted " f"— next: generate / pick-seed / assemble each video, or use the UI") return 0 def cmd_publish(args: argparse.Namespace) -> int: from . import meta_client try: row = meta_client.publish(args.video_id) except Exception as exc: # noqa: BLE001 — fail loud print(f"publish FAILED: {exc}", file=sys.stderr) return 1 print(f"published PAUSED ad {row['ad_id']} (creative {row['creative_id']}, " f"meta_video {row['meta_video_id']}) — name {row['ad_name']}") return 0 def cmd_daily(args: argparse.Namespace) -> int: """The daily learning loop (cron target): ingest -> strategize -> write-scripts -> report. Video generation/QC/assembly run via `tick` workers; publishing stays a human approval gate (SPEC_NEW §8 Phase 5).""" from . import agents, ingest, rules print("== daily loop ==") try: print(f"ingest: {ingest.run(days=args.days)} metrics_daily row(s)") briefs = agents.strategize(days=args.days, n=args.n) print(f"strategize: {len(briefs)} brief(s)") scripted = 0 for b in briefs: try: agents.write_script(str(b["id"])) scripted += 1 except Exception as exc: # noqa: BLE001 — skip the bad brief print(f" write_script {b['id']} failed: {exc}", file=sys.stderr) print(f"write-scripts: {scripted}/{len(briefs)} scripted") print(f"kill candidates: {len(rules.kill_candidates())}; " f"winner candidates: {len(rules.winner_candidates())}") except Exception as exc: # noqa: BLE001 — fail loud, non-zero for cron print(f"daily loop FAILED: {exc}", file=sys.stderr) return 1 return 0 def cmd_report(_: argparse.Namespace) -> int: from . import rules failed = db.fetch_all( "select id, type, error, created_at from jobs where status = 'failed' order by created_at desc limit 50" ) print(f"failed jobs: {len(failed)}") for j in failed: first_line = (j["error"] or "").splitlines()[0] if j["error"] else "" print(f" {j['created_at']:%Y-%m-%d %H:%M} {j['id']} [{j['type']}] {first_line}") try: kills = rules.kill_candidates() winners = rules.winner_candidates() except Exception as exc: # noqa: BLE001 — report should not crash on a metrics issue print(f"(kill/winner unavailable: {exc})", file=sys.stderr) return 0 print(f"\nkill candidates: {len(kills)}") for k in kills: print(f" {k['ad_name']} [{k.get('angle')}] — {k['reason']}") print(f"\nwinner candidates: {len(winners)}") for w in winners: cpa = float(w["cpa"]) if w.get("cpa") is not None else 0.0 print(f" {w['ad_name']} [{w.get('angle')}] — {int(w['purchases'])} purchases, CPA {cpa:.0f}") return 0 def cmd_gen_startframes(args: argparse.Namespace) -> int: """Nano-Banana: generate new startframes for a product from base startframe(s) (§16.3).""" from . import image_gen try: ids = image_gen.generate_startframe_batch(args.base, args.product, n=args.n) except Exception as exc: # noqa: BLE001 — fail loud print(f"gen-startframes FAILED: {exc}", file=sys.stderr) return 1 print(f"generated {len(ids)} startframe(s) for product {args.product}:") for i in ids: print(f" {i}") return 0 def cmd_analyze(args: argparse.Namespace) -> int: """CS#1 'what worked' → a creative_insights row (§16.4).""" from . import insights try: row = insights.analyze_winners(days=args.days) except Exception as exc: # noqa: BLE001 — fail loud print(f"analyze FAILED: {exc}", file=sys.stderr) return 1 print(f"creative_insights {row['id']} ({args.days}d): " f"angles={row.get('winning_angles')} audience={row.get('winning_audience')}") return 0 def cmd_make_body(args: argparse.Namespace) -> int: """Create a reusable body for variation testing (§16.6).""" from . import variations try: body_id = variations.create_body(args.avatar, args.line, product_id=args.product) except Exception as exc: # noqa: BLE001 — fail loud print(f"make-body FAILED: {exc}", file=sys.stderr) return 1 print(f"body {body_id} created with {len(args.line)} line(s) — next: " f"gen-hooks {body_id} --line ..., gen-ctas {body_id} --line ..., build-variations {body_id}") return 0 def cmd_gen_hooks(args: argparse.Namespace) -> int: from . import variations ids = variations.generate_hook_pool(args.body_id, args.line) print(f"added {len(ids)} hook candidate(s) to body {args.body_id}") return 0 def cmd_gen_ctas(args: argparse.Namespace) -> int: from . import variations ids = variations.generate_cta_pool(args.body_id, args.line) print(f"added {len(ids)} CTA candidate(s) to body {args.body_id}") return 0 def cmd_build_variations(args: argparse.Namespace) -> int: """Generate the body+pool clips once, then build every hook×CTA variation video (§16.6).""" from . import variations try: variations.generate_body_clips(args.body_id) # body clip generated ONCE (dry-run copies fixtures) built = variations.build_all_variations(args.body_id) except Exception as exc: # noqa: BLE001 — fail loud print(f"build-variations FAILED: {exc}", file=sys.stderr) return 1 print(f"built {len(built)} variation video(s) from body {args.body_id}:") for b in built: print(f" variation {b['variation_id']} -> video {b['video_id']}") return 0 def build_parser() -> argparse.ArgumentParser: p = argparse.ArgumentParser(prog="pipeline", description="Swedish UGC ad pipeline (v2)") sub = p.add_subparsers(dest="command", required=True) sub.add_parser("migrate", help="apply pending SQL migrations").set_defaults(fn=cmd_migrate) tick = sub.add_parser("tick", help="process due jobs") tick.add_argument("--max-jobs", type=int, default=None) tick.set_defaults(fn=cmd_tick) seed = sub.add_parser("seed-avatar", help="insert an avatar") seed.add_argument("--name", required=True) seed.add_argument("--persona", required=True, help="persona JSON") seed.add_argument("--ref-image", default=None, help="startframe image path (reference for every segment)") seed.set_defaults(fn=cmd_seed_avatar) stv = sub.add_parser("seed-test-video", help="create video + hook takes from a script JSON file") stv.add_argument("script_json") stv.add_argument("--avatar", required=True, help="avatar id") stv.add_argument("--takes", type=int, default=None, help="hook seed candidates (default SEED_TAKES)") stv.add_argument("--single-take", action="store_true", help="insert all segments on one seed") stv.add_argument("--seed", type=int, default=None, help="explicit seed for --single-take") stv.set_defaults(fn=cmd_seed_test_video) gen = sub.add_parser("generate", help="generate all pending segments for a video via Veo") gen.add_argument("video_id") gen.set_defaults(fn=cmd_generate) pick = sub.add_parser("pick-seed", help="lock the hook seed and create body+cta") pick.add_argument("video_id") pick.add_argument("--take", type=int, default=None, help="explicit take to use (default: auto-pick by QC)") pick.set_defaults(fn=cmd_pick_seed) qcp = sub.add_parser("qc", help="run QC on a segment") qcp.add_argument("segment_id") qcp.set_defaults(fn=cmd_qc) ing = sub.add_parser("ingest", help="pull ad metrics via Windsor.ai into metrics_daily") ing.add_argument("--days", type=int, default=14, help="lookback window (default 14)") ing.set_defaults(fn=cmd_ingest) strat = sub.add_parser("strategize", help="run the strategist agent (produce briefs)") strat.add_argument("--days", type=int, default=14, help="metrics lookback window (default 14)") strat.add_argument("--n", type=int, default=None, help="number of briefs (default 5)") strat.set_defaults(fn=cmd_strategize) ws = sub.add_parser("write-scripts", help="run the scriptwriter on pending briefs") ws.add_argument("--limit", type=int, default=None, help="max briefs to script this run") ws.set_defaults(fn=cmd_write_scripts) daily = sub.add_parser("daily", help="the daily learning loop: ingest -> strategize -> write-scripts -> report") daily.add_argument("--days", type=int, default=14, help="metrics lookback window (default 14)") daily.add_argument("--n", type=int, default=None, help="number of briefs (default 5)") daily.set_defaults(fn=cmd_daily) # ── Creative supply (§16) ── gsf = sub.add_parser("gen-startframes", help="Nano-Banana: new startframes for a product from base(s)") gsf.add_argument("--product", required=True, help="product id") gsf.add_argument("--base", action="append", required=True, help="base startframe id (repeatable)") gsf.add_argument("--n", type=int, default=None, help="how many to generate (default 1)") gsf.set_defaults(fn=cmd_gen_startframes) anz = sub.add_parser("analyze", help="CS#1: what worked → creative_insights") anz.add_argument("--days", type=int, default=14, help="metrics lookback window (default 14)") anz.set_defaults(fn=cmd_analyze) mb = sub.add_parser("make-body", help="create a reusable body for variation testing") mb.add_argument("--avatar", required=True, help="startframe/avatar id") mb.add_argument("--product", default=None, help="product id") mb.add_argument("--line", action="append", required=True, help="body line (repeatable)") mb.set_defaults(fn=cmd_make_body) gh = sub.add_parser("gen-hooks", help="add hook candidates to a body") gh.add_argument("body_id") gh.add_argument("--line", action="append", required=True, help="hook line (repeatable)") gh.set_defaults(fn=cmd_gen_hooks) gc = sub.add_parser("gen-ctas", help="add CTA candidates to a body") gc.add_argument("body_id") gc.add_argument("--line", action="append", required=True, help="CTA line (repeatable)") gc.set_defaults(fn=cmd_gen_ctas) bv = sub.add_parser("build-variations", help="generate clips once + build every hook×CTA variation video") bv.add_argument("body_id") bv.set_defaults(fn=cmd_build_variations) asm = sub.add_parser("assemble", help="audio-aware smart stitch + loudnorm + captions") asm.add_argument("video_id") asm.set_defaults(fn=cmd_assemble) pub = sub.add_parser("publish", help="publish a video to Meta as a PAUSED ad") pub.add_argument("video_id") pub.set_defaults(fn=cmd_publish) srv = sub.add_parser("serve", help="run the web UI") srv.add_argument("--host", default="127.0.0.1") srv.add_argument("--port", type=int, default=8000) srv.set_defaults(fn=cmd_serve) sub.add_parser("report", help="failed jobs + kill/winner candidates").set_defaults(fn=cmd_report) return p def main(argv: list[str] | None = None) -> int: args = build_parser().parse_args(argv) from .config import get_settings get_settings() # fail loudly on missing DATABASE_URL before doing anything return args.fn(args) if __name__ == "__main__": raise SystemExit(main())