diff --git a/README.md b/README.md index 926e06d..3694a07 100644 --- a/README.md +++ b/README.md @@ -148,6 +148,17 @@ bash scripts/doctor.sh # diagnóstico Com mais de uma empresa no recorte, o `report` não imprime total consolidado: sai um resumo por org. Somar empresas distintas num número só produz um valor que não serve a nenhuma delas. Para consolidar de propósito, `--all-orgs`. +Apps do Workbench (Crystal Ball e outros manifestos de `templates/apps/`): + +```bash +studiolocal workflow install # registra todo o catálogo de apps no root +studiolocal workflow install crystal-ball --force # reinstala um, sobrescrevendo edição manual +studiolocal runs show 42 --json # registro completo de uma run, com o contrato da pausa +studiolocal refine 42 --alvo prompt_video --queixa "o gancho está fraco" # ajusta fora do pipeline +``` + +`workflow install` roda sozinho no fim de `studiolocal install`, sem `--force`: uma máquina nova já nasce com o catálogo, e um YAML editado à mão no root nunca é sobrescrito em silêncio. `refine` nunca toca `status`/`finished_at` da run — só marca a nota como desatualizada. + ## Documentação - [docs/DESIGN.md](docs/DESIGN.md): plano de arquitetura V1 completo diff --git a/lib/cli.py b/lib/cli.py index fbfbcfd..055903f 100644 --- a/lib/cli.py +++ b/lib/cli.py @@ -34,7 +34,14 @@ from .report_generator import ReportGenerator from .session_manager import NoActiveSession, SessionManager from .tracker import Tracker -from .workflow_runner import WorkflowError, WorkflowPaused, WorkflowRunner, WorkflowSpec +from .workflow_runner import ( + WorkflowError, + WorkflowPaused, + WorkflowRunner, + WorkflowSpec, + _contrato_da_pausa, + contrato_persistido, +) from .workspace_detector import detect console = Console() @@ -279,6 +286,18 @@ def install(force: bool, claude_md_dir: str | None) -> None: f" Para trazê-lo para este root: studiolocal db adopt --primary {cand}" ) + # Uma máquina nova já nasce com o catálogo de apps (Crystal Ball, etc.): + # sem --force, e sem imprimir por manifesto — bootstrap não é o lugar de + # listar o catálogo inteiro, e um YAML que o David já tenha editado à mão + # no root não pode ser sobrescrito em silêncio. Nunca chamar isto fora de + # `install`: em `_ctx()`, ou em qualquer caminho que rode a cada invocação + # da CLI, sobrescreveria em silêncio uma edição manual a cada comando. + catalogo_tracker = Tracker(target / "tracker.db") + try: + _instalar_catalogo_de_apps(catalogo_tracker, AssetStore(target), slug=None, force=False) + finally: + catalogo_tracker.close() + def _patch_claude_md(workspace_root: Path, mode: str, data_root_path: Path) -> None: md = workspace_root / "CLAUDE.md" @@ -957,7 +976,11 @@ def workflow_save(name: str, from_yaml: str, description: str) -> None: target.write_text(src.read_text()) active = tracker.find_active_session() sid = active["id"] if active else None - wid = tracker.create_workflow( + # `upsert`, não `create`: salvar de novo por cima de um slug que já existe + # (o mesmo `workflow save` rodado outra vez, ou um app do catálogo que a + # pessoa também salvou à mão) não pode estourar `IntegrityError` cru contra + # o UNIQUE de `slug`. + wid = tracker.upsert_workflow( slug=slug, name=name, yaml_path=str(target.relative_to(store.root)), @@ -967,6 +990,145 @@ def workflow_save(name: str, from_yaml: str, description: str) -> None: click.echo(f"✓ Workflow `{slug}` (#{wid}) salvo em {target}") +# --- workflow install (catálogo de apps) ----------------------------------- + + +def _apps_dir() -> Path: + return resource("templates") / "apps" + + +def _manifestos_de_apps() -> list[Path]: + d = _apps_dir() + return sorted(d.glob("*.yaml")) if d.exists() else [] + + +def _slug_declarado(path: Path) -> str | None: + try: + return (yaml.safe_load(path.read_text()) or {}).get("slug") + except (OSError, yaml.YAMLError): + return None + + +def _validar_manifesto_de_app(path: Path) -> WorkflowSpec: + """Valida schema + o contrato de cada pausa `human_pick`, ANTES de gravar. + + O schema é `WorkflowSpec.from_yaml`. O contrato da pausa é a MESMA + validação que o runner aplica em tempo de execução (`_contrato_da_pausa`), + não uma reimplementação: um manifesto com `aceita: [item_da_lista]` e sem + `from`, ou `aceita: [texto_livre]` sem `campo_livre`, tem que ser rejeitado + aqui, não descoberto na primeira run que pausa. + + Não resolve `from` (que depende de `steps.*`, inexistente em tempo de + instalação): usa a PRESENÇA da chave como o bastante, porque a regra que + importa aqui é sobre a FORMA do passo, não o conteúdo que só existe em + runtime. + """ + spec = WorkflowSpec.from_yaml(path) + for step in spec.steps: + if step.get("kind") == "human_pick": + from_placeholder = ["_"] if step.get("from") is not None else [] + _contrato_da_pausa(step, from_placeholder) + return spec + + +def _instalar_um_manifesto( + tracker: Tracker, store: AssetStore, path: Path, force: bool +) -> tuple[str, str]: + """Valida e instala UM manifesto. Devolve (resultado, mensagem). + + `resultado` é um de `erro | preservado | ok`. Nada é gravado — nem o + arquivo em `workflows/`, nem a linha no banco — quando não é `ok`. + """ + try: + spec = _validar_manifesto_de_app(path) + except (WorkflowError, yaml.YAMLError) as e: + return "erro", f"{path.name}: {e}" + + target = store.root / "workflows" / f"{spec.slug}.yaml" + conteudo_novo = path.read_text() + if target.exists() and target.read_text() != conteudo_novo and not force: + return ( + "preservado", + f"`{spec.slug}` já existe em {target} e diverge do template — " + "preservado. Use --force para sobrescrever.", + ) + + target.parent.mkdir(parents=True, exist_ok=True) + target.write_text(conteudo_novo) + wid = tracker.upsert_workflow( + slug=spec.slug, + name=spec.name, + yaml_path=str(target.relative_to(store.root)), + description=spec.description, + ) + return "ok", f"✓ `{spec.slug}` instalado em {target} (#{wid})" + + +def _instalar_catalogo_de_apps( + tracker: Tracker, store: AssetStore, slug: str | None, force: bool +) -> list[tuple[str, str]]: + manifestos = _manifestos_de_apps() + if slug: + manifestos = [p for p in manifestos if _slug_declarado(p) == slug] + return [_instalar_um_manifesto(tracker, store, p, force) for p in manifestos] + + +@workflow.command("install") +@click.argument("slug", required=False) +@click.option("--force", is_flag=True, help="Sobrescreve YAML editado à mão no root.") +def workflow_install(slug: str | None, force: bool) -> None: + """Instala manifestos de app de `templates/apps/` no root. + + Sem SLUG, instala todos. Um YAML já presente no root que diverge do + template é preservado, a menos que `--force` seja passado — mesma postura + defensiva de `install` com estado antigo: avisar e não sobrescrever em + silêncio o que foi editado à mão. + """ + tracker, store, _, _ = _ctx() + manifestos = _manifestos_de_apps() + if not manifestos: + click.echo(f"Nenhum manifesto de app em {_apps_dir()}.") + return + if slug: + alvo = [p for p in manifestos if _slug_declarado(p) == slug] + if not alvo: + click.echo(f"Nenhum manifesto de app com slug `{slug}` em {_apps_dir()}.", err=True) + sys.exit(1) + manifestos = alvo + + houve_erro = False + for path in manifestos: + resultado, msg = _instalar_um_manifesto(tracker, store, path, force) + if resultado == "erro": + click.echo(f"✗ {msg}", err=True) + houve_erro = True + else: + click.echo(msg) + if houve_erro: + sys.exit(1) + + +def _contrato_da_run( + tracker: Tracker, run_id: int, estado: dict[str, Any] | None = None +) -> dict[str, Any] | None: + """O contrato da pausa de uma run, a partir do `state` gravado. + + Fonte única para os três lugares que publicam pausa (`run --json`, + `runs --json`, `runs show`): todos leem `contrato_persistido`, nunca + reconstroem a regra do fallback por conta própria. + + `estado`, quando quem chama já carregou o `state` da mesma run (ex: + `run_cmd` no ramo `except WorkflowPaused`, que já leu tudo em + `_carregar_resumo`), evita uma segunda consulta e um segundo parse do + mesmo JSON. Sem ele, consulta e faz o parse por conta própria — o + comportamento de sempre para quem só tem o `run_id` em mãos. + """ + if estado is None: + row = tracker.query("SELECT state FROM runs WHERE id = ?", (run_id,))[0] + estado = json.loads(row["state"] or "{}") + return contrato_persistido(estado) + + @main.command("run") @click.argument("workflow_slug") @click.option("--project", "project_slug", required=True) @@ -994,16 +1156,25 @@ def run_cmd( spec = WorkflowSpec.from_yaml(store.root / wrow["yaml_path"]) runner = WorkflowRunner(tracker, registry, _providers(registry), store) - def _resumo(run_id: int) -> dict[str, Any]: + def _carregar_resumo(run_id: int) -> tuple[dict[str, Any], dict[str, Any]]: + """(resumo, estado) numa consulta só, para quem vier depois (o ramo + `except WorkflowPaused`) reusar o `estado` já parseado em vez de + reconsultar a mesma run para extrair o contrato da pausa.""" row = tracker.query("SELECT * FROM runs WHERE id = ?", (run_id,))[0] estado = json.loads(row["state"] or "{}") - return { + resumo = { "run_id": run_id, "status": row["status"], "session_id": row["session_id"], "cost_brl": row["cost_brl"], "step_outputs": estado.get("step_outputs", {}), + "nota_desatualizada": estado.get("nota_desatualizada", False), } + return resumo, estado + + def _resumo(run_id: int) -> dict[str, Any]: + resumo, _ = _carregar_resumo(run_id) + return resumo try: if resume_run_id: @@ -1025,12 +1196,15 @@ def _resumo(run_id: int) -> dict[str, Any]: click.echo(f"✓ Run #{run_id} concluída") except WorkflowPaused as p: if as_json: - saida = _resumo(p.run_id) + saida, estado = _carregar_resumo(p.run_id) + contrato = _contrato_da_run(tracker, p.run_id, estado=estado) or {} saida.update( { "step_id": p.step_id, "prompt_to_user": p.prompt_to_user, "options": p.options, + "aceita": contrato.get("aceita", p.aceita), + "campo_livre": contrato.get("campo_livre", p.campo_livre), } ) click.echo(json.dumps(saida, ensure_ascii=False)) @@ -1051,12 +1225,62 @@ def _resumo(run_id: int) -> dict[str, Any]: sys.exit(2) -@main.command("runs") +def _pendente(row) -> str | None: + try: + return (json.loads(row["state"] or "{}")).get("pending_step") + except json.JSONDecodeError: + return None + + +def _pendencia_json(row) -> dict[str, Any]: + """`pending_step` e as quatro chaves `pending_*` de `runs --json`, a + partir de UM parse do `state` e do contrato persistido + (`contrato_persistido`). Todas `None` (e a contagem `0`) quando a run não + está pausada. + + Antes, `_pendente(row)` e este helper faziam cada um o seu próprio parse + do mesmo `state` para a mesma linha — esta função assume as duas + responsabilidades para que `runs --json` pague o parse uma vez só por + run. `_pendente` continua existindo à parte para o branch tabular + (não-JSON), que só precisa do step e não chama este helper. + + De propósito não entra a lista de opções inteira: é endpoint de listagem, + e cinco opções por run em toda página inflaria a resposta sem necessidade. + `nota_desatualizada` entra porque é booleano, já veio do mesmo parse, e a + tela lê exatamente esta chave para acender o aviso de nota potencialmente + obsoleta depois de um refine. + """ + try: + estado = json.loads(row["state"] or "{}") + except json.JSONDecodeError: + estado = {} + contrato = contrato_persistido(estado) + return { + "pending_step": estado.get("pending_step"), + "pending_aceita": contrato["aceita"] if contrato else None, + "pending_campo_livre": contrato["campo_livre"] if contrato else None, + "pending_prompt": contrato["prompt_to_user"] if contrato else None, + "pending_opcoes": len(contrato["options"]) if contrato else 0, + "nota_desatualizada": estado.get("nota_desatualizada", False), + } + + +@main.group(invoke_without_command=True) @click.option("--project", "project_slug", help="Filtra por project (slug).") @click.option("--paused", is_flag=True, help="Só as pausadas esperando humano.") @click.option("--json", "as_json", is_flag=True) -def runs_cmd(project_slug: str | None, paused: bool, as_json: bool) -> None: - """Lista runs. Sem isto, o run_id só existe no stdout de quem disparou.""" +@click.pass_context +def runs(ctx: click.Context, project_slug: str | None, paused: bool, as_json: bool) -> None: + """Lista runs (default) ou, com `show `, o registro completo de uma. + + Sem isto, o run_id só existe no stdout de quem disparou. + """ + if ctx.invoked_subcommand is not None: + return + _runs_listar(project_slug, paused, as_json) + + +def _runs_listar(project_slug: str | None, paused: bool, as_json: bool) -> None: tracker, _, _, _ = _ctx() sql = ( "SELECT r.id, r.status, r.session_id, r.cost_brl, r.started_at, " @@ -1075,12 +1299,6 @@ def runs_cmd(project_slug: str | None, paused: bool, as_json: bool) -> None: sql += " ORDER BY r.id DESC" rows = tracker.query(sql, tuple(args)) - def _pendente(row) -> str | None: - try: - return (json.loads(row["state"] or "{}")).get("pending_step") - except json.JSONDecodeError: - return None - if as_json: click.echo( json.dumps( @@ -1093,7 +1311,7 @@ def _pendente(row) -> str | None: "session_id": r["session_id"], "cost_brl": r["cost_brl"], "started_at": r["started_at"], - "pending_step": _pendente(r), + **_pendencia_json(r), } for r in rows ], @@ -1115,6 +1333,106 @@ def _pendente(row) -> str | None: console.print(table) +@runs.command("show") +@click.argument("run_id", type=int) +@click.option("--json", "as_json", is_flag=True) +def runs_show(run_id: int, as_json: bool) -> None: + """Registro completo de uma run, com o campo `pausa` (contrato ou `null`). + + `pausa`, quando não é `null`, tem exatamente as chaves do ramo pausado de + `run --json`: `step_id`, `prompt_to_user`, `options`, `aceita`, + `campo_livre` — a mesma `contrato_persistido` que os outros dois lugares. + + `nota_desatualizada` (`estado.get("nota_desatualizada", False)`) também + sai aqui: depois de um `refine`, o state grava a chave, mas até agora + nenhuma lista fixa de saída a incluía — a tela que lê este JSON para + mostrar o aviso de nota potencialmente obsoleta nunca a via, mesmo com o + dado presente no banco. + """ + tracker, _, _, _ = _ctx() + rows = tracker.query( + "SELECT r.*, w.slug AS workflow, p.slug AS project FROM runs r " + "JOIN workflows w ON w.id = r.workflow_id " + "JOIN projects p ON p.id = r.project_id WHERE r.id = ?", + (run_id,), + ) + if not rows: + click.echo(f"Run #{run_id} não existe.", err=True) + sys.exit(1) + row = rows[0] + estado = json.loads(row["state"] or "{}") + saida = { + "run_id": row["id"], + "status": row["status"], + "workflow": row["workflow"], + "project": row["project"], + "session_id": row["session_id"], + "cost_brl": row["cost_brl"], + "started_at": row["started_at"], + "finished_at": row["finished_at"], + "inputs": json.loads(row["inputs"] or "{}"), + "step_outputs": estado.get("step_outputs", {}), + "pausa": contrato_persistido(estado), + "nota_desatualizada": estado.get("nota_desatualizada", False), + } + if as_json: + click.echo(json.dumps(saida, ensure_ascii=False)) + return + console.print(saida) + + +# --- refine ------------------------------------------------------------ + + +@main.command("refine") +@click.argument("run_id", type=int) +@click.option("--alvo", required=True, help="`prompt_video` ou `cena:N`.") +@click.option("--queixa", required=True, help="O que não ficou bom, em texto livre.") +@click.option("--json", "as_json", is_flag=True) +def refine_cmd(run_id: int, alvo: str, queixa: str, as_json: bool) -> None: + """Ajusta um prompt já pronto (vídeo ou uma cena) a partir de uma queixa. + + Fora do pipeline declarado: entra pela conversa, depois que o pacote já + existe. Nunca toca `status`/`finished_at` da run — só marca + `nota_desatualizada`. Recusa run `running`/`failed` (exit 1) e uma mudança + de status concorrente à escrita (exit 2). + """ + from .env_loader import load_gemini_key + from .refinamento import RefinamentoConcorrente, RefinamentoError, refinar + from .workspace_detector import detect + + tracker, _, _, _ = _ctx() + ctx = detect(data_root=_root().path) + try: + gemini_key = load_gemini_key(ctx) + except MissingCredential as e: + click.echo(str(e), err=True) + sys.exit(1) + + try: + resumo = refinar(tracker, run_id, alvo, queixa, gemini_key) + except RefinamentoConcorrente as e: + if as_json: + click.echo(json.dumps({"error": str(e)}, ensure_ascii=False), err=True) + else: + click.echo(f"✗ {e}", err=True) + sys.exit(2) + except RefinamentoError as e: + if as_json: + click.echo(json.dumps({"error": str(e)}, ensure_ascii=False), err=True) + else: + click.echo(f"✗ {e}", err=True) + sys.exit(1) + + if as_json: + click.echo(json.dumps(resumo, ensure_ascii=False)) + else: + click.echo( + f"✓ Run #{run_id}: {alvo} refinado " + f"({resumo['o_que_mudou'] or 'sem resumo do motor'})" + ) + + # --- cleanup -------------------------------------------------------------- diff --git a/lib/reason_engines.py b/lib/reason_engines.py index d0ccaea..629aef1 100644 --- a/lib/reason_engines.py +++ b/lib/reason_engines.py @@ -56,6 +56,22 @@ MOTORES = frozenset(FUNCOES_POR_MOTOR) +# Lista fechada do verbo `refine` (`lib/refinamento.py`), espelhando o padrão de +# `FUNCOES_POR_MOTOR`: em Python, não no manifesto YAML, porque a edição de um +# YAML não pode alcançar refinamento. O par alvo->função é 1:1 (`prompt_video` +# só pode chamar `refine_video_prompt`, `cena` só `regenerate_scene`), então não +# há `--funcao` para escolher — o alvo decide. +# +# Não entra no manifesto como bloco `conversa:` por uma razão medida e +# documentada no RECON da frente: o interpolador do runner é cego para índice +# de lista (`{{ steps.pacote.prompts_de_cena.0... }}` devolve `None` em +# silêncio), então não há como o YAML declarar "a cena N" sem essa lacuna. Não +# é para arrumar agora — é decisão já tomada. +FUNCOES_POR_CONVERSA: dict[str, str] = { + "prompt_video": "refine_video_prompt", + "cena": "regenerate_scene", +} + class ReasonError(RuntimeError): """Argumento, função ou motor recusado antes de qualquer chamada de rede.""" diff --git a/lib/refinamento.py b/lib/refinamento.py new file mode 100644 index 0000000..1dfcac7 --- /dev/null +++ b/lib/refinamento.py @@ -0,0 +1,240 @@ +"""Verbo `refine`: ajusta um prompt já pronto (o de vídeo, ou uma cena) fora +do pipeline declarado, a partir de uma queixa em texto livre. + +Não é passo de manifesto. `FUNCOES_POR_MOTOR` (`reason_engines.py`) exclui de +propósito as funções de refinamento do motor: um passo declarado não pode +disparar refinamento de algo que ainda não existe, porque o refinamento entra +pela CONVERSA, depois que o pacote já foi gerado. O par alvo->função é 1:1 +(`FUNCOES_POR_CONVERSA`), então não há `--funcao` para escolher — o `--alvo` +decide. + +Escreve sempre no lugar CANÔNICO que as `ui.vistas` do manifesto do Crystal +Ball leem (`step_outputs.pacote.resultado...`): escrever em outro lugar deixa +a tela mostrando o valor velho. Nunca recalcula a nota — `_criticar` é privada +e cara — só marca `nota_desatualizada` e deixa a régua para quem decidir +reprocessar o pacote inteiro. +""" + +from __future__ import annotations + +import json +from datetime import datetime, timezone +from typing import Any + +from .reason_engines import FUNCOES_POR_CONVERSA +from .tracker import Tracker +from .workflow_runner import marcar_custo_nao_apurado + +_STATUS_ACEITOS = ("done", "paused") + + +class RefinamentoError(RuntimeError): + """Recusa antes de qualquer chamada de rede ou escrita no state.""" + + +class RefinamentoConcorrente(RefinamentoError): + """O status ou o state da run mudou entre a leitura e a escrita do refine. + + Cobre tanto a mudança de status (a run terminou de pausar, falhou, ou foi + retomada) quanto duas chamadas concorrentes com o MESMO status (dois + refines na mesma run, ao mesmo tempo, cada um editando uma cópia do state + lida antes de qualquer escrita). + + Nada foi perdido: a chamada de rede já aconteceu (o `generation` fica + registrado), mas o `state` novo não foi gravado por cima de uma mudança + concorrente. Rodar `refine` de novo parte do estado atual. + """ + + +def _parse_alvo(alvo: str) -> tuple[str, int | None]: + if alvo == "prompt_video": + return "prompt_video", None + if alvo.startswith("cena:"): + _, _, resto = alvo.partition(":") + try: + return "cena", int(resto) + except ValueError: + raise RefinamentoError( + f"`--alvo cena:{resto}` não é um número de cena válido." + ) from None + raise RefinamentoError( + f"`--alvo` '{alvo}' não é reconhecido. Use `prompt_video` ou `cena:N`." + ) + + +def _cena_por_numero(cenas: list[Any], numero: int) -> tuple[int, dict[str, Any]]: + for i, c in enumerate(cenas): + if isinstance(c, dict) and c.get("cena") == numero: + return i, c + disponiveis = [c.get("cena") for c in cenas if isinstance(c, dict)] + raise RefinamentoError( + f"cena {numero} não existe neste pacote. Disponíveis: {disponiveis}." + ) + + +def refinar( + tracker: Tracker, + run_id: int, + alvo: str, + queixa: str, + gemini_key: str, +) -> dict[str, Any]: + """Executa o refine e devolve um resumo: `run_id`, `alvo`, `funcao`, + `generation_id`, `o_que_mudou`, `status` (da run, depois do refine), + `resultado` (o valor novo escrito — o prompt de vídeo inteiro, ou a ficha + da cena), `avisos` (revalidado quando `alvo=prompt_video`; `None` para + `cena:N`, que não recalcula) e `nota_desatualizada` (sempre `True` depois + de um refine). + + Levanta `RefinamentoError` (ou a subclasse `RefinamentoConcorrente`, para + a escrita perdendo a corrida) para toda recusa. Nenhuma recusa tem chamada + de rede nem escrita de estado antes de si: status errado, run sem sessão, + alvo desconhecido ou cena inexistente são todos conferidos ANTES de abrir + a generation. + """ + if not (queixa or "").strip(): + raise RefinamentoError("`--queixa` não pode ser vazia.") + tipo, numero_cena = _parse_alvo(alvo) + funcao = FUNCOES_POR_CONVERSA[tipo] + + rows = tracker.query("SELECT * FROM runs WHERE id = ?", (run_id,)) + if not rows: + raise RefinamentoError(f"Run #{run_id} não existe.") + row = rows[0] + status = row["status"] + if status not in _STATUS_ACEITOS: + raise RefinamentoError( + f"Run #{run_id} está '{status}'. Refine só corre em run " + f"{list(_STATUS_ACEITOS)} — 'running' ainda está gravando o " + "pacote, e 'failed' não tem resultado para ajustar." + ) + session_id = row["session_id"] + if session_id is None: + raise RefinamentoError( + f"Run #{run_id} não tem sessão associada. `generations.session_id` " + "é obrigatório (NOT NULL), e sem sessão não há onde debitar o " + "custo do refine." + ) + + state_bruto = row["state"] + state = json.loads(state_bruto or "{}") + try: + resultado = state["step_outputs"]["pacote"]["resultado"] + except (KeyError, TypeError) as e: + raise RefinamentoError( + f"Run #{run_id} não tem `step_outputs.pacote.resultado` — o " + "pacote ainda não foi gerado, então não há o que refinar." + ) from e + + if tipo == "cena": + cenas = resultado.get("prompts_de_cena") or [] + idx, cena_atual = _cena_por_numero(cenas, numero_cena) + antes: Any = cena_atual + else: + cenas = idx = cena_atual = None + antes = resultado.get("prompt_video_final") + + modelo = f"crystalball/{funcao}" + gen_id = tracker.create_generation( + project_id=row["project_id"], + session_id=session_id, + model=modelo, + kind="reason", + prompt=None, + params={"run_id": run_id, "alvo": alvo, "queixa": queixa}, + run_id=run_id, + provider="crystalball", + ) + + from .motores.carga import chave_ligada + + try: + with chave_ligada(gemini_key) as cb: + if tipo == "prompt_video": + saida = cb.refine_video_prompt(antes, queixa, resultado.get("biblia_visual")) + resultado["prompt_video_final"] = saida.get("prompt_video_final") + # `refine_video_prompt` não revalida sozinho: sem isto, + # `_avisos` continuaria mostrando a checagem do prompt VELHO. + resultado["_avisos"] = cb.validar_prompt_video(resultado) + o_que_mudou = saida.get("o_que_mudou", "") + else: + inputs = json.loads(row["inputs"] or "{}") + briefing_original = inputs.get("briefing", "") + # `regenerate_scene` (vendorado) não tem parâmetro de queixa + # no contrato original — não existe onde encaixá-la como + # argumento nomeado. O único lugar que a leva de verdade até + # o Gemini é dentro de `briefingText`: o motor a copia para + # `ctx["briefing_resumo"]` e serializa no CONTEXTO que monta o + # prompt (`crystalball_llm.py`, `regenerate_scene`, por volta + # da linha 1318 — não é metadado solto, entra no `_call` de + # verdade). Sem isto, a queixa ficava só no histórico + # (`state["refinamentos"]`) e nunca chegava ao modelo: o JSON + # da cena mudava de hash, mas a descrição continuava a mesma. + briefing_com_queixa = ( + f"{briefing_original}\n\nAjuste pedido para esta cena " + f"especificamente (corrija isso, mantenha o resto): {queixa}" + if briefing_original + else ( + "Ajuste pedido para esta cena especificamente " + f"(corrija isso, mantenha o resto): {queixa}" + ) + ) + contexto = { + "briefingText": briefing_com_queixa, + "direcao": resultado.get("direcao_escolhida", ""), + "biblia_visual": resultado.get("biblia_visual"), + } + nova_cena = cb.regenerate_scene(contexto, cena_atual) + cenas[idx] = nova_cena + o_que_mudou = "" + except Exception as e: + tracker.finish_generation(gen_id, status="failed", error=str(e)) + raise + + marcar_custo_nao_apurado(tracker, gen_id, "crystalball") + + state["nota_desatualizada"] = True + historico = state.setdefault("refinamentos", []) + historico.append( + { + "n": len(historico) + 1, + "em": datetime.now(timezone.utc).isoformat(), + "alvo": alvo, + "funcao": funcao, + "queixa": queixa, + "generation_id": gen_id, + "o_que_mudou": o_que_mudou, + "antes": antes, + } + ) + + if not tracker.update_run_state_if_status(run_id, state, status, state_bruto): + raise RefinamentoConcorrente( + f"Run #{run_id} mudou de status ou de state entre a leitura e a " + f"escrita do refine (status era '{status}'). A geração #{gen_id} " + "foi registrada; o state novo não foi gravado. Rode `refine` de " + "novo sobre o estado atual." + ) + + if tipo == "prompt_video": + # o valor novo inteiro, e o `_avisos` já revalidado contra ele. + resultado_novo: Any = resultado["prompt_video_final"] + avisos: Any = resultado.get("_avisos") + else: + # a ficha da cena nova. `_avisos` é sobre o prompt de vídeo inteiro — + # refinar uma cena não o recalcula, então expor o valor velho aqui + # sugeriria uma revalidação que não aconteceu. + resultado_novo = cenas[idx] + avisos = None + + return { + "run_id": run_id, + "alvo": alvo, + "funcao": funcao, + "generation_id": gen_id, + "o_que_mudou": o_que_mudou, + "status": status, + "resultado": resultado_novo, + "avisos": avisos, + "nota_desatualizada": state["nota_desatualizada"], + } diff --git a/lib/tracker.py b/lib/tracker.py index 1e72ee9..67cc55b 100644 --- a/lib/tracker.py +++ b/lib/tracker.py @@ -418,6 +418,38 @@ def create_workflow( ) return cur.lastrowid + def upsert_workflow( + self, + slug: str, + name: str, + yaml_path: str, + source_session_id: int | None = None, + description: str | None = None, + ) -> int: + """Registra um workflow, ou atualiza o que já existe com este `slug`. + + `create_workflow` é um INSERT puro contra `slug UNIQUE` + (`migrations/001_initial.sql:38`), então registrar duas vezes (`workflow + save` rodado de novo, `workflow install` reinstalando o catálogo) + estourava `IntegrityError` cru. `source_session_id` não entra no + `DO UPDATE`: reinstalar um manifesto de app não tem sessão de origem, e + sobrescrever a que já existia apagaria a proveniência de um workflow + salvo por `workflow save`. + """ + with self.transaction() as conn: + conn.execute( + "INSERT INTO workflows (slug, name, yaml_path, source_session_id, description) " + "VALUES (?,?,?,?,?) " + "ON CONFLICT(slug) DO UPDATE SET " + "name=excluded.name, yaml_path=excluded.yaml_path, " + "description=excluded.description", + (slug, name, yaml_path, source_session_id, description), + ) + row = conn.execute( + "SELECT id FROM workflows WHERE slug = ?", (slug,) + ).fetchone() + return row["id"] + def get_workflow_by_slug(self, slug: str) -> sqlite3.Row | None: return self._conn.execute( "SELECT * FROM workflows WHERE slug = ?", (slug,) @@ -465,6 +497,43 @@ def update_run( with self.transaction() as conn: conn.execute(f"UPDATE runs SET {', '.join(sets)} WHERE id = ?", vals) + def update_run_state_if_status( + self, + run_id: int, + state: dict[str, Any], + status_esperado: str, + state_bruto_esperado: str | None = None, + ) -> bool: + """Escreve `state` só se `status` ainda for `status_esperado`. Devolve + se escreveu. + + `update_run` incondicional entre a leitura e a escrita do verbo + `refine` perderia em silêncio uma mudança concorrente de status (a run + terminou de pausar de novo, falhou, ou foi retomada enquanto a chamada + de rede do refine estava em voo): o `state` novo sobrescreveria o que + essa mudança gravou, sem que nada acusasse. O `WHERE status = ?` faz o + `UPDATE` não afetar nenhuma linha nesse caso, em vez de gravar por + cima — quem chama decide o que fazer com zero linhas afetadas. + + Isso sozinho não pega DUAS chamadas concorrentes com o MESMO status + (ex: refine de `prompt_video` e refine de `cena:1` na mesma run, ao + mesmo tempo): nenhuma das duas muda `status`, então as duas passariam + pela guarda antiga e a segunda escrita apagaria a primeira por + completo, sem erro. `state_bruto_esperado`, quando passado, estende o + `WHERE` para comparar também a string exata de `state` lida no início + da chamada (o valor bruto de `row["state"]`, antes do parse) — se + qualquer coisa escreveu no state entre a leitura e esta escrita, + combinando status ou não, a comparação falha e `rowcount` vem zero. + """ + sql = "UPDATE runs SET state = ? WHERE id = ? AND status = ?" + params: list[Any] = [json.dumps(state), run_id, status_esperado] + if state_bruto_esperado is not None: + sql += " AND state = ?" + params.append(state_bruto_esperado) + with self.transaction() as conn: + cur = conn.execute(sql, params) + return cur.rowcount > 0 + # --- queries de leitura (suporte a reports) ------------------------------- def query(self, sql: str, params: tuple = ()) -> list[sqlite3.Row]: diff --git a/lib/workflow_runner.py b/lib/workflow_runner.py index 9eae39a..0e564af 100644 --- a/lib/workflow_runner.py +++ b/lib/workflow_runner.py @@ -145,6 +145,33 @@ def _contrato_da_pausa( return aceita, campo +def contrato_persistido(state: dict[str, Any]) -> dict[str, Any] | None: + """O contrato da pausa gravada em `state`, ou `None` se a run não está pausada. + + Existe porque `state["pending_aceita"]`/`state["pending_campo_livre"]` + (`:382-383`) ficavam presos no disco: nada os publicava no JSON que a CLI + imprime (`run --json`, `runs --json`, `runs show`). Espelha o fallback que + `_resolver_escolha` aplica a `pending_aceita` ausente (uma pausa gravada + antes de esta chave existir no state): `[item_da_lista]` quando há opções, + em vez de `aceita=None`, que faria os três lugares acima publicarem `null` + onde uma run de verdade está esperando resposta. + """ + pending_step = state.get("pending_step") + if not pending_step: + return None + opcoes = state.get("pending_options") or [] + aceita = state.get("pending_aceita") + if aceita is None: + aceita = [ACEITA_ITEM] if opcoes else [] + return { + "step_id": pending_step, + "prompt_to_user": state.get("pending_prompt"), + "options": list(opcoes), + "aceita": list(aceita), + "campo_livre": state.get("pending_campo_livre"), + } + + def _interpolate(value: Any, ctx: dict[str, Any]) -> Any: """Substitui {{ a.b.c }} dentro de strings (recursivo em dicts/lists).""" if isinstance(value, str): @@ -185,6 +212,29 @@ class StepResult: cost_brl: float = 0.0 +def marcar_custo_nao_apurado(tracker: Tracker, generation_id: int, motor: str) -> None: + """Fecha uma generation de raciocínio com custo zero por IGNORÂNCIA, não + por gratuidade: o motor não mede token, então "estimado pelo catálogo" + seria mentira. + + Fonte única do texto `cost_source='nao-apurado'`, chamada tanto por um + passo `reason` do runner quanto pelo verbo `refine` (`lib/refinamento.py`), + para o rótulo não divergir entre os dois lugares. + """ + tracker.finish_generation(generation_id, status="done", cost_brl=0.0) + tracker.merge_generation_params( + generation_id, + { + "cost_source": "nao-apurado", + "cost_gaps": [ + f"passo de raciocínio ({motor}): tokens não medidos", + "o motor não lê usageMetadata e models.yaml não tem preço " + "por milhão de token: R$ 0,00 aqui é ignorância, não gratuidade", + ], + }, + ) + + class WorkflowRunner: def __init__( self, @@ -563,20 +613,7 @@ def _run_reason_step( except Exception as e: self.tracker.finish_generation(gen_id, status="failed", error=str(e)) raise - self.tracker.finish_generation(gen_id, status="done", cost_brl=0.0) - # O zero precisa dizer que é zero por ignorância, não por gratuidade: o - # motor não mede token, então "estimado pelo catálogo" seria mentira. - self.tracker.merge_generation_params( - gen_id, - { - "cost_source": "nao-apurado", - "cost_gaps": [ - f"passo de raciocínio ({motor}): tokens não medidos", - "o motor não lê usageMetadata e models.yaml não tem preço " - "por milhão de token: R$ 0,00 aqui é ignorância, não gratuidade", - ], - }, - ) + marcar_custo_nao_apurado(self.tracker, gen_id, motor) return StepResult(step_id=step["id"], outputs=outputs, cost_brl=0.0) def _finalize( diff --git a/tests/test_motor_vendorado.py b/tests/test_motor_vendorado.py new file mode 100644 index 0000000..ab6f5ce --- /dev/null +++ b/tests/test_motor_vendorado.py @@ -0,0 +1,189 @@ +"""O motor vendorado de verdade, carregado offline (PROD-2128, item 4). + +Todo outro teste do Crystal Ball troca o `ReasonRouter` por um duplo +(`MotorFalso`) — nenhum chega a carregar `lib/motores/crystalball_llm.py`. O +motor só importa stdlib mais `requests` (confira o topo do arquivo), então +carregá-lo aqui é offline e barato: nenhuma chamada de rede, só import + +inspeção. O que este arquivo prova, e que nenhum outro prova: + +1. o blob do arquivo vendorado bate com o hash declarado em `PROVENIENCIA.md` + — a checagem por blob que pega o dia em que o arquivo derivar do ref fixado + (`13a55d5`), mesmo por um `ruff format` que não muda comportamento nenhum; +2. `carga.motor()` registra `_crystalball_motor` em `sys.modules` (não + `crystalball_llm`, de propósito) e não deixa `sys.path` diferente de como + encontrou; +3. as assinaturas reais das seis funções alcançáveis (as quatro de + `FUNCOES_POR_MOTOR["crystalball"]` mais as duas de `FUNCOES_POR_CONVERSA`) + aceitam os argumentos que os adaptadores (`reason_engines.py`, + `refinamento.py`) de fato passam; +4. `PESOS` do motor tem 7 chaves, soma 1.0, e é IDÊNTICO ao `PESOS_REAIS` + copiado à mão em `tests/test_reason_engines.py:29-37` — pega o dia em que + os dois derivarem; +5. `chave_ligada()` injeta e restaura `GEMINI_API_KEY`, sem deixá-la vazando + para o processo depois do `with`, inclusive quando o bloco levanta. +""" + +from __future__ import annotations + +import hashlib +import inspect +import os +import re +import sys +from pathlib import Path + +import pytest + +from lib.motores import carga +from lib.reason_engines import FUNCOES_POR_CONVERSA, FUNCOES_POR_MOTOR +from tests.test_reason_engines import PESOS_REAIS + +MOTOR_PATH = Path(__file__).resolve().parents[1] / "lib" / "motores" / "crystalball_llm.py" +PROVENIENCIA_PATH = Path(__file__).resolve().parents[1] / "lib" / "motores" / "PROVENIENCIA.md" + + +@pytest.fixture(autouse=True) +def motor_sem_cache_entre_testes(): + """`motor()` cacheia em `sys.modules` por processo. Sem este teardown, o + teste que prova o registro dependeria de rodar antes de qualquer outro que + já tenha carregado o motor — acoplamento de ordem que o próprio item pede + para evitar.""" + sys.modules.pop("_crystalball_motor", None) + sys.modules.pop("store", None) + yield + sys.modules.pop("_crystalball_motor", None) + sys.modules.pop("store", None) + + +def _blob_esperado() -> str: + texto = PROVENIENCIA_PATH.read_text() + m = re.search(r"\|\s*blob git\s*\|\s*`([0-9a-f]{40})`\s*\|", texto) + assert m, "PROVENIENCIA.md não tem a linha `blob git` no formato esperado" + return m.group(1) + + +# --- 1. blob bate com a proveniência --------------------------------------- + + +def test_blob_do_motor_bate_com_o_hash_declarado_na_proveniencia() -> None: + conteudo = MOTOR_PATH.read_bytes() + # Mesmo método que `git hash-object`: sha1("blob \0" + conteúdo). + calculado = hashlib.sha1(b"blob %d\0" % len(conteudo) + conteudo).hexdigest() + assert calculado == _blob_esperado() + + +# --- 2. carga: sys.modules e sys.path --------------------------------------- + + +def test_motor_registra_em_sys_modules_com_nome_interno_e_nao_toca_sys_path() -> None: + assert "_crystalball_motor" not in sys.modules + assert "crystalball_llm" not in sys.modules + path_antes = list(sys.path) + + mod = carga.motor() + + assert sys.modules.get("_crystalball_motor") is mod + # Nome interno, de propósito: um `import crystalball_llm` acidental em + # outro lugar do processo não pode pegar este módulo por acaso. + assert "crystalball_llm" not in sys.modules + assert sys.path == path_antes + + +def test_motor_carregado_duas_vezes_devolve_o_mesmo_modulo() -> None: + m1 = carga.motor() + m2 = carga.motor() + assert m1 is m2 + + +def test_motor_registra_o_store_sintetico() -> None: + carga.motor() + store = sys.modules.get("store") + assert store is not None + assert store.SUGESTOES_NARRATIVE and store.SUGESTOES_EMOTION and store.SUGESTOES_AUDIO + + +# --- 3. assinaturas batem com o que os adaptadores passam ------------------- + +# (posicionais, kwargs) EXATAMENTE como `reason_engines.py` e +# `lib/refinamento.py` chamam cada função hoje. Ler a assinatura real e +# tentar `Signature.bind` com isto é o que pega o dia em que um nome de +# parâmetro mudar no upstream. +_CHAMADAS_REAIS: dict[str, tuple[tuple, dict]] = { + "pesquisar_referencias": (({"cliente": "x"},), {}), + "caminhos_criativos": ( + ({"cliente": "x"},), + {"images": None, "ja_apresentados": None, "dossie": {}}, + ), + "diagnose": (({"cliente": "x"},), {"images": None, "caminho": "ideia"}), + "generate_video": ( + ({"cliente": "x"},), + {"images": None, "diagnostico": {"d": 1}, "direcao": "d"}, + ), + "refine_video_prompt": (("prompt atual", "queixa", {"paleta": []}), {}), + "regenerate_scene": (({"briefingText": ""}, {"cena": 1}), {}), +} + + +def test_lista_fechada_bate_exatamente_com_as_chamadas_conferidas() -> None: + """Controle: se alguém acrescentar função à lista fechada sem atualizar + este teste, ele acusa a lacuna em vez de passar em silêncio.""" + todas = set(FUNCOES_POR_MOTOR["crystalball"]) | set(FUNCOES_POR_CONVERSA.values()) + assert todas == set(_CHAMADAS_REAIS) + + +@pytest.mark.parametrize("nome", sorted(_CHAMADAS_REAIS)) +def test_assinatura_real_aceita_o_que_o_adaptador_passa(nome: str) -> None: + mod = carga.motor() + fn = getattr(mod, nome) + args, kwargs = _CHAMADAS_REAIS[nome] + sig = inspect.signature(fn) + sig.bind(*args, **kwargs) # levanta TypeError se a assinatura não aceitar + + +# --- 4. PESOS: 7 chaves, soma 1.0, idêntico ao gabarito --------------------- + + +def test_pesos_tem_sete_chaves_e_soma_um() -> None: + mod = carga.motor() + assert len(mod.PESOS) == 7 + assert sum(mod.PESOS.values()) == pytest.approx(1.0, abs=1e-9) + + +def test_pesos_e_identico_ao_gabarito_copiado_a_mao() -> None: + """Pega o dia em que `PESOS` do motor e `PESOS_REAIS` de + `test_reason_engines.py` derivarem um do outro.""" + mod = carga.motor() + assert mod.PESOS == PESOS_REAIS + + +# --- 5. chave_ligada: injeta e restaura ------------------------------------- + + +def test_chave_ligada_injeta_e_restaura_quando_nao_havia_chave( + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.delenv("GEMINI_API_KEY", raising=False) + with carga.chave_ligada("chave-de-teste") as mod: + assert os.environ.get("GEMINI_API_KEY") == "chave-de-teste" + assert mod.has_key() + assert "GEMINI_API_KEY" not in os.environ + + +def test_chave_ligada_restaura_o_valor_anterior_em_vez_de_apagar( + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setenv("GEMINI_API_KEY", "valor-anterior-do-processo") + with carga.chave_ligada("chave-de-teste"): + assert os.environ["GEMINI_API_KEY"] == "chave-de-teste" + assert os.environ["GEMINI_API_KEY"] == "valor-anterior-do-processo" + + +def test_chave_ligada_restaura_mesmo_quando_o_bloco_levanta( + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.delenv("GEMINI_API_KEY", raising=False) + with pytest.raises(RuntimeError, match="boom"): + with carga.chave_ligada("chave-de-teste"): + assert os.environ.get("GEMINI_API_KEY") == "chave-de-teste" + raise RuntimeError("boom") + assert "GEMINI_API_KEY" not in os.environ diff --git a/tests/test_pausa_publicada.py b/tests/test_pausa_publicada.py new file mode 100644 index 0000000..cd1dfb9 --- /dev/null +++ b/tests/test_pausa_publicada.py @@ -0,0 +1,194 @@ +"""O contrato da pausa chega ao JSON que a CLI imprime (PROD-2128, item 1). + +O runner já persistia `pending_aceita`/`pending_campo_livre` no state +(`workflow_runner.py:382-383`), mas nada os publicava: nem `run --json`, nem +`runs --json`, e não existia como consultar UMA run pausada por fora do stdout +de quem a disparou. `contrato_persistido`, extraída de `_resolver_escolha`, é o +que os três lugares agora leem. Este teste prova dois pontos que ela precisa +acertar: + +- o contrato de uma pausa REAL do Crystal Ball (`escolha_do_caminho`), com + `aceita`/`campo_livre` vindos do YAML, não inventados; +- o fallback de uma pausa gravada por uma versão anterior, sem + `pending_aceita` no state, que não pode virar `aceita: null` — o valor que + faria a tela achar que a pausa não aceita nada. +""" + +from __future__ import annotations + +import json +from pathlib import Path + +import pytest +from click.testing import CliRunner + +from lib.asset_store import AssetStore +from lib.cli import main +from lib.models_registry import ModelsRegistry +from lib.tracker import Tracker +from lib.workflow_runner import ( + ACEITA_ITEM, + WorkflowPaused, + WorkflowRunner, + WorkflowSpec, + contrato_persistido, +) + +MANIFESTO = Path(__file__).resolve().parents[1] / "templates/apps/crystal-ball.yaml" + +BRIEFING = { + "cliente": "Rider", + "segmento": "Calçados", + "objetivo": "Lançar a linha R10", + "plataforma": "reels-tiktok", + "duracao": 30, + "objetivo_primario": "retencao", + "briefing": "Sandália de borracha, público jovem, tom irreverente.", + "referencias": [], +} + + +class RouterFalso: + """Duplo mínimo: só o suficiente para o Crystal Ball chegar à 1a pausa.""" + + def valida(self, motor: str, funcao: str) -> None: + pass + + def chamar(self, motor: str, funcao: str, args: dict) -> dict: + if funcao == "pesquisar_referencias": + return {"resultado": {}, "pesquisa_ok": False, "mecanicas": 0} + if funcao == "caminhos_criativos": + caminhos = [{"numero": f"{i:02d}", "nome": f"Caminho {i}"} for i in range(1, 6)] + return { + "resultado": {"caminhos": caminhos, "_pesquisa_ok": False}, + "opcoes": [ + {"id": c["numero"], "rotulo": c["nome"], "caminho": c} for c in caminhos + ], + "pesquisa_ok": False, + } + raise AssertionError(f"funcao inesperada para este teste: {funcao}") + + +def _env(root: Path) -> dict[str, str | None]: + return {"STUDIOLOCAL_ROOT": str(root), "STUDIO_DATA_ROOT": None} + + +@pytest.fixture +def root(tmp_path: Path) -> Path: + r = tmp_path / "data-root" / ".studiolocal" + r.mkdir(parents=True) + return r + + +def _root_pronto(root: Path) -> None: + """Só o suficiente para `root.is_ready()` aceitar (tracker.db existe).""" + tracker = Tracker(root / "tracker.db", create=True) + tracker.apply_migrations() + tracker.close() + + +def _pausar_crystal_ball(root: Path) -> tuple[int, WorkflowPaused]: + """Roda o manifesto real do Crystal Ball, com motor falso, até a primeira + pausa (`escolha_do_caminho`). Devolve (run_id, a exceção de pausa) e fecha + o próprio tracker, para a CLI abrir o mesmo arquivo livre depois.""" + tracker = Tracker(root / "tracker.db", create=True) + tracker.apply_migrations() + runner = WorkflowRunner(tracker, ModelsRegistry(), None, AssetStore(root)) + runner._reason_router = lambda: RouterFalso() + pid = tracker.create_project("rider", "Rider", [], None) + sid = tracker.open_session(pid) + wid = tracker.create_workflow("crystal-ball", "Crystal Ball", str(MANIFESTO)) + spec = WorkflowSpec.from_yaml(MANIFESTO) + with pytest.raises(WorkflowPaused) as e: + runner.start(spec, pid, "rider", sid, dict(BRIEFING), wid) + tracker.close() + return e.value.run_id, e.value + + +# --- contrato de uma pausa real -------------------------------------------- + + +def test_escolha_do_caminho_publica_aceita_e_campo_livre_do_yaml(root: Path) -> None: + run_id, pausa = _pausar_crystal_ball(root) + assert pausa.step_id == "escolha_do_caminho" + assert pausa.aceita == ["item_da_lista", "texto_livre"] + assert pausa.campo_livre == "caminho" + + tracker = Tracker(root / "tracker.db") + estado = json.loads(tracker.query("SELECT state FROM runs WHERE id = ?", (run_id,))[0]["state"]) + tracker.close() + + contrato = contrato_persistido(estado) + assert contrato is not None + assert contrato["step_id"] == "escolha_do_caminho" + assert contrato["aceita"] == ["item_da_lista", "texto_livre"] + assert contrato["campo_livre"] == "caminho" + assert len(contrato["options"]) == 5 + + +# --- fallback de pausa gravada por versão anterior ------------------------- + + +def test_state_antigo_sem_pending_aceita_cai_no_fallback_e_nao_publica_null() -> None: + """Pausa gravada antes de `pending_aceita` existir no state: sem opções + aqui seria `aceita=[]`, mas COM opções ela vira `[item_da_lista]`, nunca + `None` — que é o valor que faria a tela achar a pausa sem contrato.""" + legado = { + "pending_step": "escolha", + "pending_options": [{"id": "01"}, {"id": "02"}], + "pending_prompt": "Qual?", + } + contrato = contrato_persistido(legado) + assert contrato is not None + assert contrato["aceita"] == [ACEITA_ITEM] + assert contrato["campo_livre"] is None + assert contrato["step_id"] == "escolha" + + +def test_state_sem_pausa_devolve_none() -> None: + assert contrato_persistido({"step_outputs": {}}) is None + assert contrato_persistido({"pending_step": None}) is None + + +# --- `runs show` pela CLI --------------------------------------------------- + + +def test_runs_show_de_id_inexistente_sai_com_exit_1(root: Path) -> None: + _root_pronto(root) + runner = CliRunner() + result = runner.invoke(main, ["runs", "show", "999", "--json"], env=_env(root)) + assert result.exit_code == 1 + assert "não existe" in result.output + + +def test_runs_show_de_run_pausada_mostra_pausa_com_as_chaves_do_run_json(root: Path) -> None: + run_id, pausa = _pausar_crystal_ball(root) + runner = CliRunner() + result = runner.invoke(main, ["runs", "show", str(run_id), "--json"], env=_env(root)) + assert result.exit_code == 0, result.output + saida = json.loads(result.output) + + # Mesmas chaves que o ramo `except WorkflowPaused` de `run --json` publica. + assert saida["pausa"] is not None + assert set(saida["pausa"]) == { + "step_id", "prompt_to_user", "options", "aceita", "campo_livre" + } + assert saida["pausa"]["step_id"] == pausa.step_id + assert saida["pausa"]["aceita"] == pausa.aceita + assert saida["pausa"]["campo_livre"] == pausa.campo_livre + assert saida["pausa"]["options"] == pausa.options + assert saida["status"] == "paused" + + +def test_runs_json_traz_as_quatro_chaves_pending_extras(root: Path) -> None: + run_id, pausa = _pausar_crystal_ball(root) + runner = CliRunner() + result = runner.invoke(main, ["runs", "--json"], env=_env(root)) + assert result.exit_code == 0, result.output + linhas = json.loads(result.output) + row = next(r for r in linhas if r["run_id"] == run_id) + assert row["pending_step"] == "escolha_do_caminho" + assert row["pending_aceita"] == pausa.aceita + assert row["pending_campo_livre"] == pausa.campo_livre + assert row["pending_prompt"] == pausa.prompt_to_user + assert row["pending_opcoes"] == len(pausa.options) diff --git a/tests/test_refine.py b/tests/test_refine.py new file mode 100644 index 0000000..9d12717 --- /dev/null +++ b/tests/test_refine.py @@ -0,0 +1,519 @@ +"""`studiolocal refine`: ajusta um prompt já pronto fora do pipeline (PROD-2128, +item 3). + +Motor substituído por um duplo, no padrão de `MotorFalso` de +`test_reason_engines.py:42`: zero chamada de rede. O que estes testes seguram: + +- recusa ANTES de qualquer chamada de rede ou escrita: status errado, cena que + não existe, run sem sessão; +- escreve no lugar CANÔNICO que as `ui.vistas` do manifesto leem + (`step_outputs.pacote.resultado...`), nunca em outro lugar; +- `_avisos` é recalculado (`refine_video_prompt` não revalida sozinho); +- a run nunca sai de `done`/`paused` por causa do refine; +- o ledger de custo é `cost_source='nao-apurado'`, a mesma fonte que o runner + usa para passo `reason` (`marcar_custo_nao_apurado`); +- a escrita do state é atômica: uma mudança de status entre a leitura e a + escrita não pode ser sobrescrita em silêncio. +""" + +from __future__ import annotations + +import contextlib +import json +from pathlib import Path + +import pytest +from click.testing import CliRunner + +from lib.cli import main +from lib.refinamento import RefinamentoConcorrente, RefinamentoError, refinar +from lib.tracker import Tracker + +PACOTE = { + "titulo": "T", + "direcao_escolhida": "a direção afiada", + "biblia_visual": {"paleta": ["#000"]}, + "prompt_video_final": "prompt original", + "prompts_de_cena": [ + {"cena": i, "prompt_imagem": f"cena {i}", "ficha_tecnica": {}} for i in range(1, 6) + ], + "analise_preditiva": {"gancho_3s": {"nota": 4.0}}, + "nota_final_ponderada": 4.0, + "_avisos": ["aviso antigo, de antes do refine"], +} + + +class MotorFalso: + """Duplo do motor. Registra cada chamada, zero rede.""" + + def __init__(self) -> None: + self.chamadas: list[tuple[str, dict]] = [] + + def has_key(self) -> bool: + return True + + def refine_video_prompt(self, prompt, complaint, biblia=None): + self.chamadas.append( + ("refine_video_prompt", {"prompt": prompt, "complaint": complaint, "biblia": biblia}) + ) + return {"prompt_video_final": "NOVO PROMPT REFINADO", "o_que_mudou": "trocou o gancho"} + + def regenerate_scene(self, contexto, cena): + self.chamadas.append(("regenerate_scene", {"contexto": dict(contexto), "cena": dict(cena)})) + return {"cena": cena.get("cena"), "prompt_imagem": "cena regenerada", "ficha_tecnica": {}} + + def validar_prompt_video(self, conceito): + self.chamadas.append(("validar_prompt_video", {"conceito": dict(conceito)})) + return ["aviso recalculado depois do refine"] + + +@pytest.fixture +def motor(monkeypatch: pytest.MonkeyPatch) -> MotorFalso: + m = MotorFalso() + + @contextlib.contextmanager + def chave_falsa(key, modelo=None): + yield m + + monkeypatch.setattr("lib.motores.carga.chave_ligada", chave_falsa, raising=True) + return m + + +def _tracker(tmp_path: Path) -> Tracker: + t = Tracker(tmp_path / "tracker.db", create=True) + t.apply_migrations() + return t + + +def _run_pronta( + tracker: Tracker, + status: str = "done", + com_sessao: bool = True, + briefing: str = "Sandália de borracha, público jovem.", +) -> int: + """Uma run com o pacote do Crystal Ball já no state, pronta para refinar.""" + pid = tracker.create_project("rider", "Rider", [], None) + sid = tracker.open_session(pid) if com_sessao else None + wid = tracker.create_workflow("crystal-ball", "Crystal Ball", "crystal-ball.yaml") + run_id = tracker.create_run(wid, pid, sid, {"briefing": briefing}) + state = { + "step_outputs": {"pacote": {"resultado": json.loads(json.dumps(PACOTE))}}, + "cost_total": 0.0, + } + tracker.update_run(run_id, status=status, state=state) + return run_id + + +def _estado(tracker: Tracker, run_id: int) -> dict: + row = tracker.query("SELECT state FROM runs WHERE id = ?", (run_id,))[0] + return json.loads(row["state"]) + + +# --- recusas ANTES de qualquer chamada de rede ------------------------------ + + +def test_refine_numa_run_running_e_recusado(tmp_path: Path, motor: MotorFalso) -> None: + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker, status="running") + with pytest.raises(RefinamentoError, match="running"): + refinar(tracker, run_id, "prompt_video", "o gancho está fraco", "chave-falsa") + assert motor.chamadas == [] + assert tracker.query("SELECT COUNT(*) AS n FROM generations")[0]["n"] == 0 + + +def test_refine_numa_run_failed_e_recusado(tmp_path: Path, motor: MotorFalso) -> None: + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker, status="failed") + with pytest.raises(RefinamentoError, match="failed"): + refinar(tracker, run_id, "prompt_video", "o gancho está fraco", "chave-falsa") + assert motor.chamadas == [] + + +def test_refine_numa_run_paused_e_aceito(tmp_path: Path, motor: MotorFalso) -> None: + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker, status="paused") + refinar(tracker, run_id, "prompt_video", "o gancho está fraco", "chave-falsa") + row = tracker.query("SELECT status FROM runs WHERE id = ?", (run_id,))[0] + assert row["status"] == "paused", "refine não pode tocar o status da run" + + +def test_cena_9_numa_run_com_5_cenas_falha_listando_as_cinco( + tmp_path: Path, motor: MotorFalso +) -> None: + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker) + with pytest.raises(RefinamentoError, match=r"\[1, 2, 3, 4, 5\]"): + refinar(tracker, run_id, "cena:9", "queixa", "chave-falsa") + assert motor.chamadas == [] + + +def test_refine_numa_run_sem_sessao_falha_com_mensagem_clara( + tmp_path: Path, motor: MotorFalso +) -> None: + """`generations.session_id` é NOT NULL (migrations/001_initial.sql:62): sem + sessão, a recusa tem de vir ANTES do INSERT, não deixar o banco estourar.""" + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker, com_sessao=False) + with pytest.raises(RefinamentoError, match="sessão"): + refinar(tracker, run_id, "prompt_video", "queixa", "chave-falsa") + assert motor.chamadas == [] + assert tracker.query("SELECT COUNT(*) AS n FROM generations")[0]["n"] == 0 + + +def test_alvo_desconhecido_e_recusado(tmp_path: Path, motor: MotorFalso) -> None: + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker) + with pytest.raises(RefinamentoError, match="alvo"): + refinar(tracker, run_id, "storyboard", "queixa", "chave-falsa") + assert motor.chamadas == [] + + +def test_queixa_vazia_e_recusada(tmp_path: Path, motor: MotorFalso) -> None: + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker) + with pytest.raises(RefinamentoError, match="queixa"): + refinar(tracker, run_id, "prompt_video", " ", "chave-falsa") + assert motor.chamadas == [] + + +# --- escreve no lugar canônico ---------------------------------------------- + + +def test_refine_de_prompt_video_troca_o_valor_e_mantem_status_done( + tmp_path: Path, motor: MotorFalso +) -> None: + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker, status="done") + resumo = refinar(tracker, run_id, "prompt_video", "o gancho está fraco", "chave-falsa") + + row = tracker.query("SELECT status, finished_at FROM runs WHERE id = ?", (run_id,))[0] + assert row["status"] == "done" + + estado = _estado(tracker, run_id) + resultado = estado["step_outputs"]["pacote"]["resultado"] + assert resultado["prompt_video_final"] == "NOVO PROMPT REFINADO" + assert estado["nota_desatualizada"] is True + assert resumo["o_que_mudou"] == "trocou o gancho" + assert resumo["funcao"] == "refine_video_prompt" + + +def test_refine_de_cena_troca_so_o_elemento_correspondente( + tmp_path: Path, motor: MotorFalso +) -> None: + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker) + refinar(tracker, run_id, "cena:3", "a cena 3 ficou escura", "chave-falsa") + + estado = _estado(tracker, run_id) + cenas = estado["step_outputs"]["pacote"]["resultado"]["prompts_de_cena"] + assert cenas[2]["prompt_imagem"] == "cena regenerada" + # as outras quatro cenas não foram tocadas + assert cenas[0]["prompt_imagem"] == "cena 1" + assert cenas[3]["prompt_imagem"] == "cena 4" + assert motor.chamadas[0][1]["cena"]["cena"] == 3 + + +def test_refine_de_cena_leva_a_queixa_para_o_motor( + tmp_path: Path, motor: MotorFalso +) -> None: + """Achado real: `refine --alvo cena:N` devolvia sucesso e trocava o hash + da cena, mas a queixa nunca chegava ao motor — `regenerate_scene` não tem + parâmetro de queixa no contrato vendorado, e o wrapper montava `contexto` + só com briefing geral, direção e bíblia visual. + + A queixa tem que entrar em `briefingText`, que é o único campo que o + motor de fato copia para o prompt que vai pro Gemini (`ctx["briefing_ + resumo"]`, em `regenerate_scene`). E o briefing original não pode + desaparecer — a correção é um acréscimo, não uma substituição. + """ + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker, briefing="Sandália de borracha, público jovem.") + refinar(tracker, run_id, "cena:3", "rua com grafite", "chave-falsa") + + ((_, args),) = [c for c in motor.chamadas if c[0] == "regenerate_scene"] + contexto = args["contexto"] + assert "rua com grafite" in contexto["briefingText"] + assert "Sandália de borracha, público jovem." in contexto["briefingText"] + + +# --- o resumo devolvido (o JSON de `refine --json`) ------------------------- + + +def test_resumo_de_prompt_video_traz_status_resultado_avisos_e_nota_desatualizada( + tmp_path: Path, motor: MotorFalso +) -> None: + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker, status="done") + resumo = refinar(tracker, run_id, "prompt_video", "o gancho está fraco", "chave-falsa") + + estado = _estado(tracker, run_id) + resultado = estado["step_outputs"]["pacote"]["resultado"] + + assert resumo["status"] == "done" + assert resumo["resultado"] == resultado["prompt_video_final"] == "NOVO PROMPT REFINADO" + assert resumo["avisos"] == resultado["_avisos"] == ["aviso recalculado depois do refine"] + assert resumo["nota_desatualizada"] is True + + +def test_resumo_de_cena_traz_status_resultado_avisos_none_e_nota_desatualizada( + tmp_path: Path, motor: MotorFalso +) -> None: + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker, status="paused") + resumo = refinar(tracker, run_id, "cena:3", "a cena 3 ficou escura", "chave-falsa") + + estado = _estado(tracker, run_id) + cena_nova = estado["step_outputs"]["pacote"]["resultado"]["prompts_de_cena"][2] + + assert resumo["status"] == "paused" + assert resumo["resultado"] == cena_nova + assert resumo["resultado"]["prompt_imagem"] == "cena regenerada" + # refinar uma cena não recalcula `_avisos` (que é sobre o prompt de vídeo + # inteiro) — expor o valor velho aqui sugeriria uma revalidação que não + # aconteceu, então o campo vem `None`. + assert resumo["avisos"] is None + assert resumo["nota_desatualizada"] is True + + +def test_avisos_e_recalculado_depois_do_refine(tmp_path: Path, motor: MotorFalso) -> None: + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker) + refinar(tracker, run_id, "prompt_video", "queixa", "chave-falsa") + + estado = _estado(tracker, run_id) + resultado = estado["step_outputs"]["pacote"]["resultado"] + assert resultado["_avisos"] == ["aviso recalculado depois do refine"] + assert any(nome == "validar_prompt_video" for nome, _ in motor.chamadas) + + +def test_refine_de_cena_nao_recalcula_avisos(tmp_path: Path, motor: MotorFalso) -> None: + """Só `prompt_video` rechama `validar_prompt_video` — regenerar uma cena + não dispara a validação do prompt de vídeo inteiro.""" + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker) + refinar(tracker, run_id, "cena:1", "queixa", "chave-falsa") + assert not any(nome == "validar_prompt_video" for nome, _ in motor.chamadas) + estado = _estado(tracker, run_id) + assert estado["step_outputs"]["pacote"]["resultado"]["_avisos"] == PACOTE["_avisos"] + + +# --- `nota_desatualizada` chega na CLI, não só no state --------------------- + + +def test_runs_show_json_traz_nota_desatualizada_depois_de_um_refine( + tmp_path: Path, motor: MotorFalso +) -> None: + """Achado real: depois de um refine, `state["nota_desatualizada"]` fica + `true` no banco (conferido direto via sqlite), mas `studiolocal runs show + --json` não incluía o campo — a lista fixa de chaves da saída não o + citava, então o aviso de nota possivelmente desatualizada nunca aparecia + na tela, mesmo com o dado presente. + + Passa pela CLI de verdade (`click.testing.CliRunner` sobre `lib.cli.main`) + porque o bug estava exatamente na função que monta esse JSON + (`runs_show`), não em `refinar()` — o resumo devolvido por `refinar()` já + trazia `nota_desatualizada` corretamente antes desta correção. + """ + root = tmp_path / "data-root" / ".studiolocal" + root.mkdir(parents=True) + tracker = Tracker(root / "tracker.db", create=True) + tracker.apply_migrations() + run_id = _run_pronta(tracker, status="done") + refinar(tracker, run_id, "prompt_video", "o gancho está fraco", "chave-falsa") + tracker.close() # solta o handle antes do CLI abrir o mesmo arquivo + + runner = CliRunner() + result = runner.invoke( + main, + ["runs", "show", str(run_id), "--json"], + env={"STUDIOLOCAL_ROOT": str(root), "STUDIO_DATA_ROOT": None}, + ) + assert result.exit_code == 0, result.output + saida = json.loads(result.output) + assert saida["nota_desatualizada"] is True + + +# --- histórico append-only --------------------------------------------------- + + +def test_refinamentos_cresce_um_item_por_chamada(tmp_path: Path, motor: MotorFalso) -> None: + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker) + refinar(tracker, run_id, "prompt_video", "queixa 1", "chave-falsa") + refinar(tracker, run_id, "cena:2", "queixa 2", "chave-falsa") + + estado = _estado(tracker, run_id) + hist = estado["refinamentos"] + assert len(hist) == 2 + assert [h["n"] for h in hist] == [1, 2] + assert hist[0]["alvo"] == "prompt_video" + assert hist[0]["antes"] == "prompt original" + assert hist[1]["alvo"] == "cena:2" + assert hist[1]["antes"]["cena"] == 2 + assert all(h["generation_id"] for h in hist) + assert all(h["queixa"] for h in hist) + assert all("em" in h for h in hist) + + +# --- ledger de custo --------------------------------------------------------- + + +def test_generation_criada_tem_cost_source_nao_apurado( + tmp_path: Path, motor: MotorFalso +) -> None: + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker) + resumo = refinar(tracker, run_id, "prompt_video", "queixa", "chave-falsa") + + gen = tracker.query( + "SELECT * FROM generations WHERE id = ?", (resumo["generation_id"],) + )[0] + assert gen["model"] == "crystalball/refine_video_prompt" + assert gen["provider"] == "crystalball" + assert gen["kind"] == "reason" + assert gen["status"] == "done" + params = json.loads(gen["params"]) + assert params["cost_source"] == "nao-apurado" + + +def test_generation_de_regenerate_scene_usa_o_nome_da_funcao( + tmp_path: Path, motor: MotorFalso +) -> None: + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker) + resumo = refinar(tracker, run_id, "cena:1", "queixa", "chave-falsa") + gen = tracker.query( + "SELECT model FROM generations WHERE id = ?", (resumo["generation_id"],) + )[0] + assert gen["model"] == "crystalball/regenerate_scene" + + +# --- escrita atômica --------------------------------------------------------- + + +def test_mudanca_de_status_concorrente_levanta_e_nao_perde_a_generation( + tmp_path: Path, motor: MotorFalso, monkeypatch: pytest.MonkeyPatch +) -> None: + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker) + monkeypatch.setattr(tracker, "update_run_state_if_status", lambda *a, **k: False) + + with pytest.raises(RefinamentoConcorrente): + refinar(tracker, run_id, "prompt_video", "queixa", "chave-falsa") + + # a chamada de rede já tinha acontecido: a generation existe, mesmo que o + # state novo não tenha sido gravado por cima da mudança concorrente. + gens = tracker.query("SELECT * FROM generations WHERE run_id = ?", (run_id,)) + assert len(gens) == 1 + assert gens[0]["status"] == "done" + + +def test_update_run_state_if_status_so_grava_se_o_status_bater(tmp_path: Path) -> None: + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker, status="done") + + assert tracker.update_run_state_if_status(run_id, {"x": 1}, "done") is True + assert _estado(tracker, run_id) == {"x": 1} + + assert tracker.update_run_state_if_status(run_id, {"y": 2}, "paused") is False + assert _estado(tracker, run_id) == {"x": 1}, "escrita rejeitada não pode ter passado" + + +def test_update_run_state_if_status_com_state_esperado_recusa_state_obsoleto( + tmp_path: Path, +) -> None: + """A guarda por `status` sozinha não pega duas escritas com o MESMO + status: `state_bruto_esperado` estende o `WHERE` para também comparar o + state bruto lido antes da escrita, e uma segunda escrita contra um state + já obsoleto tem que ser recusada mesmo com o status intacto.""" + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker, status="done") + state_bruto_original = tracker.query( + "SELECT state FROM runs WHERE id = ?", (run_id,) + )[0]["state"] + + assert ( + tracker.update_run_state_if_status( + run_id, {"x": 1}, "done", state_bruto_original + ) + is True + ) + assert _estado(tracker, run_id) == {"x": 1} + + # Mesmo state bruto esperado, agora obsoleto (a escrita acima já mudou o + # state) — mesmo com o status ainda 'done', a escrita não pode passar. + assert ( + tracker.update_run_state_if_status( + run_id, {"y": 2}, "done", state_bruto_original + ) + is False + ) + assert _estado(tracker, run_id) == {"x": 1}, "escrita contra state obsoleto não pode ter passado" + + +def test_duas_refines_concorrentes_com_o_mesmo_status_a_segunda_perde_e_nao_apaga_a_primeira( + tmp_path: Path, motor: MotorFalso, monkeypatch: pytest.MonkeyPatch +) -> None: + """Reproduz o achado: `refine` nunca muda `status`, então duas chamadas + concorrentes sobre a MESMA run com o MESMO status (A refina + `prompt_video`, B refina `cena:1`) passavam as duas pela guarda antiga — + a segunda escrita sobrescrevia a primeira por completo, em silêncio. + + Simula a concorrência lendo o MESMO (status, state) para as duas + chamadas: A lê e escreve normalmente; B é forçada a ler o snapshot + ANTIGO (o que ela teria lido se tivesse consultado antes de A escrever, + como aconteceria com duas conexões concorrentes). A segunda escrita tem + que falhar com `RefinamentoConcorrente`, e o state final tem que ser o + de A, intacto. + """ + tracker = _tracker(tmp_path) + run_id = _run_pronta(tracker, status="done") + + # O snapshot que as DUAS chamadas concorrentes teriam lido: mesmo + # (status, state), antes de qualquer uma escrever. + snapshot_pre_concorrencia = tracker.query( + "SELECT * FROM runs WHERE id = ?", (run_id,) + )[0] + + # Chamada A: refina o prompt de vídeo. Lê e escreve normalmente — nada + # ainda mudou no banco, então ela lê exatamente o mesmo snapshot. + refinar(tracker, run_id, "prompt_video", "o gancho está fraco", "chave-falsa") + estado_apos_a = _estado(tracker, run_id) + assert ( + estado_apos_a["step_outputs"]["pacote"]["resultado"]["prompt_video_final"] + == "NOVO PROMPT REFINADO" + ) + + # Chamada B: refina a cena 1, mas a leitura do state é forçada a devolver + # o snapshot de ANTES de A escrever — exatamente o que uma segunda + # conexão concorrente teria lido, porque a leitura dela aconteceu antes + # da escrita de A ser commitada. + query_original = tracker.query + + def query_com_leitura_concorrente(sql: str, params: tuple = ()) -> list: + if sql.strip().startswith("SELECT * FROM runs WHERE id"): + return [snapshot_pre_concorrencia] + return query_original(sql, params) + + monkeypatch.setattr(tracker, "query", query_com_leitura_concorrente) + + with pytest.raises(RefinamentoConcorrente): + refinar(tracker, run_id, "cena:1", "a cena 1 ficou escura", "chave-falsa") + + monkeypatch.setattr(tracker, "query", query_original) + + # a chamada de rede de B já aconteceu (registrou generation), mas o + # state em disco continua sendo o de A, intacto — B não apagou nada. + gens = tracker.query("SELECT * FROM generations WHERE run_id = ?", (run_id,)) + assert len(gens) == 2 + estado_final = _estado(tracker, run_id) + assert estado_final == estado_apos_a + assert ( + estado_final["step_outputs"]["pacote"]["resultado"]["prompt_video_final"] + == "NOVO PROMPT REFINADO" + ) + # a edição de B (a cena 1 regenerada) não está no state final: foi + # perdida na chamada de rede, não no state — exatamente o comportamento + # esperado quando a escrita otimista recusa. + cenas_finais = estado_final["step_outputs"]["pacote"]["resultado"]["prompts_de_cena"] + assert cenas_finais[0]["prompt_imagem"] == "cena 1", "B não pode ter tocado o state" diff --git a/tests/test_workflow_install.py b/tests/test_workflow_install.py new file mode 100644 index 0000000..447c07c --- /dev/null +++ b/tests/test_workflow_install.py @@ -0,0 +1,215 @@ +"""`workflow install`: o manifesto que ninguém registrava (PROD-2128, item 2). + +`templates/apps/crystal-ball.yaml` não era citado por nenhum instalador: quem +queria o app no catálogo tinha que rodar `workflow save` na mão, e rodar duas +vezes estourava `IntegrityError` cru (`create_workflow` é um INSERT puro contra +o `slug UNIQUE` de `migrations/001_initial.sql:38`) — o bug que o David batia +ao registrar na mão. Este teste prova o instalador: idempotente via +`upsert_workflow`, defensivo com YAML editado à mão no root, e que rejeita +manifesto quebrado sem gravar nada (nem arquivo, nem banco). +""" + +from __future__ import annotations + +import json +from pathlib import Path + +import pytest +import yaml +from click.testing import CliRunner + +import lib.cli as cli +from lib.cli import main +from lib.tracker import Tracker + +REPO_APPS = Path(__file__).resolve().parents[1] / "templates" / "apps" + + +def _env(root: Path) -> dict[str, str | None]: + return {"STUDIOLOCAL_ROOT": str(root), "STUDIO_DATA_ROOT": None} + + +@pytest.fixture +def root(tmp_path: Path) -> Path: + """Root pronto (migrado, com as pastas), mas sem o catálogo instalado.""" + r = tmp_path / "data-root" / ".studiolocal" + tracker = Tracker(r / "tracker.db", create=True) + tracker.apply_migrations() + tracker.close() + for sub in ("projects", "workflows", "_tmp", "archive", "logs"): + (r / sub).mkdir(exist_ok=True) + return r + + +def _manifestos_reais() -> list[str]: + return sorted(p.stem for p in REPO_APPS.glob("*.yaml")) + + +# --- instala o catálogo real ------------------------------------------------ + + +def test_install_registra_crystal_ball_e_os_outros_apps_do_catalogo(root: Path) -> None: + slugs = _manifestos_reais() + assert "crystal-ball" in slugs # controle: é o app que este item prova + + runner = CliRunner() + result = runner.invoke(main, ["workflow", "install"], env=_env(root)) + assert result.exit_code == 0, result.output + + tracker = Tracker(root / "tracker.db") + try: + for slug in slugs: + assert tracker.get_workflow_by_slug(slug) is not None, slug + finally: + tracker.close() + + +def test_rodar_duas_vezes_nao_duplica_nem_estoura(root: Path) -> None: + runner = CliRunner() + r1 = runner.invoke(main, ["workflow", "install"], env=_env(root)) + assert r1.exit_code == 0, r1.output + r2 = runner.invoke(main, ["workflow", "install"], env=_env(root)) + assert r2.exit_code == 0, r2.output + + tracker = Tracker(root / "tracker.db") + try: + n = tracker.query( + "SELECT COUNT(*) AS n FROM workflows WHERE slug = 'crystal-ball'" + )[0]["n"] + finally: + tracker.close() + assert n == 1 + + +def test_workflow_list_apps_devolve_crystal_ball_com_bloco_ui(root: Path) -> None: + runner = CliRunner() + runner.invoke(main, ["workflow", "install"], env=_env(root)) + result = runner.invoke(main, ["workflow", "list", "--json", "--apps"], env=_env(root)) + assert result.exit_code == 0, result.output + apps = json.loads(result.output) + cb = next((a for a in apps if a["slug"] == "crystal-ball"), None) + assert cb is not None, apps + assert isinstance(cb["ui"], dict) and cb["ui"] + + +# --- postura defensiva com YAML editado no root ----------------------------- + + +def test_yaml_editado_a_mao_e_preservado_sem_force_e_sobrescrito_com_force( + root: Path, +) -> None: + runner = CliRunner() + runner.invoke(main, ["workflow", "install", "crystal-ball"], env=_env(root)) + alvo = root / "workflows" / "crystal-ball.yaml" + editado = alvo.read_text() + "\n# edição manual do David\n" + alvo.write_text(editado) + + sem_force = runner.invoke( + main, ["workflow", "install", "crystal-ball"], env=_env(root) + ) + assert sem_force.exit_code == 0, sem_force.output + assert "preservado" in sem_force.output + assert alvo.read_text() == editado, "sem --force não pode sobrescrever" + + com_force = runner.invoke( + main, ["workflow", "install", "crystal-ball", "--force"], env=_env(root) + ) + assert com_force.exit_code == 0, com_force.output + assert "# edição manual" not in alvo.read_text(), "--force tem de sobrescrever" + + +# --- manifesto quebrado ------------------------------------------------------ + + +def test_yaml_quebrado_e_rejeitado_sem_gravar_nada( + root: Path, tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + """`aceita: [texto_livre]` sem `campo_livre` é exatamente o que + `_contrato_da_pausa` recusa em runtime — a instalação tem que recusar + ANTES, não deixar a pessoa descobrir na primeira run que pausa.""" + apps_fake = tmp_path / "apps-fake" + apps_fake.mkdir() + quebrado = { + "schema": "studiolocal/workflow/v1", + "slug": "quebrado", + "name": "Quebrado", + "inputs": {}, + "steps": [ + { + "id": "pausa", + "kind": "human_pick", + "prompt_to_user": "Qual?", + "aceita": ["texto_livre"], + # campo_livre ausente de propósito + } + ], + } + (apps_fake / "quebrado.yaml").write_text(yaml.safe_dump(quebrado)) + monkeypatch.setattr(cli, "_apps_dir", lambda: apps_fake) + + runner = CliRunner() + result = runner.invoke(main, ["workflow", "install"], env=_env(root)) + assert result.exit_code == 1 + assert "campo_livre" in result.output + + tracker = Tracker(root / "tracker.db") + try: + assert tracker.get_workflow_by_slug("quebrado") is None + finally: + tracker.close() + assert not (root / "workflows" / "quebrado.yaml").exists() + + +def test_manifesto_valido_ao_lado_do_quebrado_ainda_e_instalado( + root: Path, tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + """Um app quebrado não pode impedir os outros do catálogo de entrarem.""" + apps_fake = tmp_path / "apps-fake" + apps_fake.mkdir() + bom = { + "schema": "studiolocal/workflow/v1", + "slug": "bom", + "name": "Bom", + "inputs": {}, + "steps": [{"id": "s", "kind": "reason", "motor": "template", "outputs": {"x": 1}}], + } + quebrado = { + "schema": "studiolocal/workflow/v1", + "slug": "quebrado", + "name": "Quebrado", + "inputs": {}, + "steps": [ + {"id": "pausa", "kind": "human_pick", "prompt_to_user": "Qual?", + "aceita": ["texto_livre"]}, + ], + } + (apps_fake / "bom.yaml").write_text(yaml.safe_dump(bom)) + (apps_fake / "quebrado.yaml").write_text(yaml.safe_dump(quebrado)) + monkeypatch.setattr(cli, "_apps_dir", lambda: apps_fake) + + runner = CliRunner() + result = runner.invoke(main, ["workflow", "install"], env=_env(root)) + assert result.exit_code == 1 + + tracker = Tracker(root / "tracker.db") + try: + assert tracker.get_workflow_by_slug("bom") is not None + assert tracker.get_workflow_by_slug("quebrado") is None + finally: + tracker.close() + + +# --- o bootstrap `install` também semeia o catálogo ------------------------- + + +def test_install_bootstrap_tambem_registra_o_catalogo(tmp_path: Path) -> None: + root = tmp_path / "fresh" / ".studiolocal" + runner = CliRunner() + result = runner.invoke(main, ["install"], env=_env(root)) + assert result.exit_code == 0, result.output + + tracker = Tracker(root / "tracker.db") + try: + assert tracker.get_workflow_by_slug("crystal-ball") is not None + finally: + tracker.close()