From 538267ed3916674f298d0567fa359c52c81b6863 Mon Sep 17 00:00:00 2001 From: davidbenal <144815978+davidbenal@users.noreply.github.com> Date: Sat, 22 Aug 2026 18:45:33 -0300 Subject: [PATCH 1/2] =?UTF-8?q?Um=20app=20declarado=20precisa=20guardar=20?= =?UTF-8?q?o=20que=20n=C3=A3o=20=C3=A9=20m=C3=ADdia,=20e=20a=20retomada=20?= =?UTF-8?q?precisa=20ser=20confi=C3=A1vel?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A camada de Apps do Workbench (PROD-2127) tem passos cujo resultado é estrutura, não arquivo: um dossiê de pesquisa, uma lista de caminhos criativos, um diagnóstico. O runner não sabia representar isso. `ProviderOutput` só aceita `url` ou `data` (provider_base.py:25-32), e `_run_step` assume que todo passo salva arquivo. Passo de raciocínio não podia virar generation com asset falso só para caber nesse caminho: asset draft acende "precisa de você" no kanban e infla a contagem de entrega. Entra o kind `reason`: grava a estrutura em `runs.state`, grava uma linha em `generations` e nenhum asset. A linha entra no ledger por três motivos medidos: enquanto está `pending` o card acende "em geração", o que cobre espera longa de motor; `created_at`/`finished_at` dão latência por passo, que não existia em lugar nenhum; e `cost_brl` ganha onde morar quando o custo de token passar a ser medido. O zero vem marcado `cost_source: nao-apurado`, porque zero sem procedência é o erro que a apuração de junho já cometeu. Três consertos na retomada, todos encontrados ao tentar provar o mecanismo: 1. `resume` herdava `session_id` do parâmetro, e o CLI passava `None` (cli.py:939). Com `session_id or 0` isso violava a FK de `generations.session_id` e a run inteira virava `failed` na primeira geração depois da pausa. Reproduzido: `FOREIGN KEY constraint failed`. Agora o `session_id` vem da linha da run, que é quem sabe. 2. A escolha do humano era gravada crua, sem nenhuma checagem. Qualquer `--selected` era aceito, inclusive um que nunca esteve nas opções, e o passo seguinte consumia um valor que nunca esteve no estado. A pausa passa a persistir as opções, e a retomada resolve a resposta contra elas: `'inventada' não está entre as opções da pausa em 'escolha'`. Opção escalar continua devolvendo o dicionário do humano, para `{{ steps.pick.selected }}` seguir resolvendo. 3. Falha de step apagava o ponteiro e o trabalho já pago. Um 429 do provider depois de uma escolha humana perdia a escolha. `step_outputs` continua preservado e a run falhada passa a ser retomável sem repetir a pergunta. Mais: a pausa acende `sessions.awaiting_input`, que é o que faz o card ir para "precisa de você" no Workbench; `from` string vira uma opção, não uma por caractere; passo de mídia que declara output de estrutura falha alto em vez de descartar em silêncio; `workflow save` respeita o `slug` do YAML em vez de derivá-lo do `--name`; e `run`, `workflow list` e `workflow show` ganham `--json`, que é como o Workbench lê o mesmo arquivo sem ponte. `studiolocal runs [--paused]` é novo e existe porque o `run_id` só vivia no stdout de quem disparou: sem ele, run pausada era invisível depois do terminal fechar. CI: primeiro workflow deste repositório. Gate é pytest mais o subconjunto do ruff que pega defeito (E9, F). Os 23 achados de estilo herdados ficam declarados em voz alta num passo não bloqueante, porque seis deles pedem renomear exceção pública. Provado no data root de desenvolvimento do Workbench, com geração real: run 1 pausou em `escolha`, o processo morreu, a retomada de processo frio recusou `{"id":"inventada"}` e aceitou `{"id":"chuva"}`, e o prompt gravado saiu "uma foto de jabuticaba sob chuva, tom neon (eco de jabuticaba)", com as duas metades vindas de `runs.state`. Asset 1024x1024, R$ 0,68, sessão aberta. Co-Authored-By: Claude Opus 5 (1M context) --- .github/workflows/testes.yml | 45 +++++ lib/cli.py | 174 +++++++++++++++- lib/workflow_runner.py | 184 +++++++++++++++-- templates/apps/eco.yaml | 86 ++++++++ tests/test_workflow_app_steps.py | 327 +++++++++++++++++++++++++++++++ 5 files changed, 788 insertions(+), 28 deletions(-) create mode 100644 .github/workflows/testes.yml create mode 100644 templates/apps/eco.yaml create mode 100644 tests/test_workflow_app_steps.py diff --git a/.github/workflows/testes.yml b/.github/workflows/testes.yml new file mode 100644 index 0000000..31cd43f --- /dev/null +++ b/.github/workflows/testes.yml @@ -0,0 +1,45 @@ +# Primeiro CI deste repositório. +# +# Até 2026-08-22 não havia nenhum: `.github/` não existia, e os 194 testes só +# rodavam se alguém lembrasse de rodar. Isso passou a doer porque a metade do +# app do Workbench que pausa, retoma e guarda estado mora aqui, e "está verde" +# era afirmação de quem escreveu, nunca fato observável por quem revisa. +# +# O gate é o pytest mais o subconjunto do ruff que pega defeito de verdade +# (E9 sintaxe, F nome indefinido e import morto). O ruff completo roda ao lado, +# sem reprovar, porque o repositório carrega 23 achados de estilo herdados e +# seis deles pedem renomear exceção pública, que é mudança de API e não cabe +# nesta frente. A dívida está declarada em voz alta na saída do job, com issue +# aberta: preferível a um `continue-on-error` mudo ou a um verde por omissão. + +name: Testes + +on: + push: + branches: [main] + pull_request: + workflow_dispatch: + +jobs: + pytest: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + + - uses: astral-sh/setup-uv@v5 + with: + enable-cache: true + + - name: Instala com os extras de dev + run: uv pip install --system -e ".[dev]" + + - name: Gate do ruff (sintaxe, nome indefinido, import morto) + run: ruff check . --select E9,F --output-format concise + + - name: Dívida de estilo, declarada e não bloqueante + run: | + echo "Achados de estilo herdados (não reprovam este job):" + ruff check . --statistics --exit-zero + + - name: Testes + run: pytest -q diff --git a/lib/cli.py b/lib/cli.py index 89bbb20..fbfbcfd 100644 --- a/lib/cli.py +++ b/lib/cli.py @@ -11,8 +11,10 @@ import subprocess import sys from pathlib import Path +from typing import Any import click +import yaml from rich.console import Console from rich.table import Table @@ -867,9 +869,22 @@ def workflow() -> None: @workflow.command("list") -def workflow_list() -> None: - tracker, _, _, _ = _ctx() +@click.option("--json", "as_json", is_flag=True, help="Uma linha JSON, para consumo por programa.") +@click.option("--apps", "only_apps", is_flag=True, help="Só os que têm bloco `ui` (apps).") +def workflow_list(as_json: bool, only_apps: bool) -> None: + tracker, store, _, _ = _ctx() rows = tracker.query("SELECT * FROM workflows ORDER BY created_at DESC") + if as_json: + saida = [] + for r in rows: + ui = _ui_do_yaml(store.root / r["yaml_path"]) + if only_apps and not ui: + continue + saida.append( + {"slug": r["slug"], "name": r["name"], "yaml_path": r["yaml_path"], "ui": ui} + ) + click.echo(json.dumps(saida, ensure_ascii=False)) + return if not rows: click.echo("Nenhum workflow salvo ainda.") return @@ -882,6 +897,42 @@ def workflow_list() -> None: console.print(table) +def _ui_do_yaml(path: Path) -> dict[str, Any] | None: + """Devolve o bloco `ui` de um manifesto, ou None se não houver. + + Nunca levanta: um arquivo ruim não pode derrubar a lista inteira, pela mesma + razão que ESTADO.md ausente devolve vazio em vez de erro. + """ + try: + data = yaml.safe_load(path.read_text()) or {} + except (OSError, yaml.YAMLError): + return None + ui = data.get("ui") + return ui if isinstance(ui, dict) else None + + +@workflow.command("show") +@click.argument("workflow_slug") +@click.option("--json", "as_json", is_flag=True) +def workflow_show(workflow_slug: str, as_json: bool) -> None: + """Despeja o manifesto inteiro. É por aqui que o Workbench lê o app.""" + tracker, store, _, _ = _ctx() + wrow = tracker.get_workflow_by_slug(workflow_slug) + if not wrow: + click.echo(f"Workflow `{workflow_slug}` não encontrado.", err=True) + sys.exit(1) + path = store.root / wrow["yaml_path"] + try: + data = yaml.safe_load(path.read_text()) or {} + except (OSError, yaml.YAMLError) as e: + click.echo(f"Manifesto ilegível em {path}: {e}", err=True) + sys.exit(1) + if as_json: + click.echo(json.dumps(data, ensure_ascii=False)) + else: + click.echo(path.read_text()) + + @workflow.command("save") @click.option("--name", required=True, help="Nome legível: 'Hero Fashion Light'") @click.option("--from-yaml", "from_yaml", required=True, help="Path do YAML pronto") @@ -892,7 +943,15 @@ def workflow_save(name: str, from_yaml: str, description: str) -> None: if not src.exists(): click.echo(f"Arquivo {src} não existe.", err=True) sys.exit(1) - slug = slugify(name) + # O slug do arquivo vence o do nome: senão `--name "Eco Bonito"` grava + # `eco-bonito` no banco enquanto o YAML continua dizendo `slug: eco`, e o + # lookup por slug passa a apontar para nada. + try: + declarado = (yaml.safe_load(src.read_text()) or {}).get("slug") + except yaml.YAMLError as e: + click.echo(f"YAML inválido em {src}: {e}", err=True) + sys.exit(1) + slug = str(declarado) if declarado else slugify(name) target = store.root / "workflows" / f"{slug}.yaml" target.parent.mkdir(parents=True, exist_ok=True) target.write_text(src.read_text()) @@ -913,13 +972,15 @@ def workflow_save(name: str, from_yaml: str, description: str) -> None: @click.option("--project", "project_slug", required=True) @click.option("--inputs", help='JSON: \'{"prompt":"X"}\'') @click.option("--resume", "resume_run_id", type=int, help="Resume run pausada") -@click.option("--selected", help='Para human_pick: \'{"selected": }\'') +@click.option("--selected", help='Para human_pick: \'{"id": ""}\'') +@click.option("--json", "as_json", is_flag=True, help="Uma linha JSON, para consumo por programa.") def run_cmd( workflow_slug: str, project_slug: str, inputs: str | None, resume_run_id: int | None, selected: str | None, + as_json: bool, ) -> None: tracker, store, registry, sm = _ctx() proj = tracker.get_project_by_slug(project_slug) @@ -933,32 +994,127 @@ 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]: + row = tracker.query("SELECT * FROM runs WHERE id = ?", (run_id,))[0] + estado = json.loads(row["state"] or "{}") + return { + "run_id": run_id, + "status": row["status"], + "session_id": row["session_id"], + "cost_brl": row["cost_brl"], + "step_outputs": estado.get("step_outputs", {}), + } + try: if resume_run_id: user_input = json.loads(selected) if selected else {} runner.resume(spec, resume_run_id, proj["id"], proj["slug"], None, user_input) - click.echo(f"✓ Run #{resume_run_id} retomada e concluída") + if as_json: + click.echo(json.dumps(_resumo(resume_run_id), ensure_ascii=False)) + else: + click.echo(f"✓ Run #{resume_run_id} retomada e concluída") else: session_id = sm.ensure_session(proj["id"]) inp = json.loads(inputs) if inputs else {} run_id = runner.start( spec, proj["id"], proj["slug"], session_id, inp, wrow["id"] ) - click.echo(f"✓ Run #{run_id} concluída") + if as_json: + click.echo(json.dumps(_resumo(run_id), ensure_ascii=False)) + else: + click.echo(f"✓ Run #{run_id} concluída") except WorkflowPaused as p: + if as_json: + saida = _resumo(p.run_id) + saida.update( + { + "step_id": p.step_id, + "prompt_to_user": p.prompt_to_user, + "options": p.options, + } + ) + click.echo(json.dumps(saida, ensure_ascii=False)) + return click.echo(f"\n⏸ Run #{p.run_id} pausada no step '{p.step_id}'") click.echo(f" {p.prompt_to_user}") - click.echo(f" Opções (asset IDs): {p.options}") + click.echo(f" Opções: {p.options}") click.echo( f"\n Para continuar:\n" f" studiolocal run {workflow_slug} --project {project_slug} " - f"--resume {p.run_id} --selected '{{\"selected\": }}'" + f"--resume {p.run_id} --selected '{{\"id\": \"\"}}'" ) except WorkflowError as e: - click.echo(f"✗ {e}", err=True) + 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) +@main.command("runs") +@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.""" + tracker, _, _, _ = _ctx() + sql = ( + "SELECT r.id, r.status, r.session_id, r.cost_brl, r.started_at, " + "w.slug AS workflow, p.slug AS project, r.state " + "FROM runs r JOIN workflows w ON w.id = r.workflow_id " + "JOIN projects p ON p.id = r.project_id" + ) + where, args = [], [] + if project_slug: + where.append("p.slug = ?") + args.append(project_slug) + if paused: + where.append("r.status = 'paused'") + if where: + sql += " WHERE " + " AND ".join(where) + 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( + [ + { + "run_id": r["id"], + "status": r["status"], + "workflow": r["workflow"], + "project": r["project"], + "session_id": r["session_id"], + "cost_brl": r["cost_brl"], + "started_at": r["started_at"], + "pending_step": _pendente(r), + } + for r in rows + ], + ensure_ascii=False, + ) + ) + return + if not rows: + click.echo("Nenhuma run." if not paused else "Nenhuma run pausada.") + return + table = Table(title="Runs") + for col in ("run", "status", "workflow", "project", "sessão", "passo pendente", "início"): + table.add_column(col) + for r in rows: + table.add_row( + str(r["id"]), r["status"], r["workflow"], r["project"], + str(r["session_id"] or "-"), _pendente(r) or "-", r["started_at"], + ) + console.print(table) + + # --- cleanup -------------------------------------------------------------- diff --git a/lib/workflow_runner.py b/lib/workflow_runner.py index c28c57b..0bbaa63 100644 --- a/lib/workflow_runner.py +++ b/lib/workflow_runner.py @@ -1,7 +1,14 @@ """Parser + executor de Workflow YAMLs. -Suporta pipelines lineares com 4 kinds de step: image, video, upscale, human_pick. -Steps human_pick PAUSAM a Run e persistem state. Resume com user input. +Suporta pipelines lineares com 5 kinds de step: image, video, upscale, human_pick +e reason. Steps human_pick PAUSAM a Run e persistem state. Resume com user input. + +`reason` é o passo cujo resultado não é mídia: ele grava estrutura em +`runs.state` e uma linha em `generations` sem nenhum asset. Existe porque um app +declarado tem passos de raciocínio (um dossiê, uma lista de caminhos, um +diagnóstico) que o caminho de mídia não sabe representar: `ProviderOutput` só +aceita `url` ou `data` (provider_base.py:25-32) e `_run_step` assume arquivo em +todo passo. Não suporta branching/condicionais — composição complexa = Claude orquestrando múltiplas Workflows conversacionalmente. @@ -30,7 +37,7 @@ class WorkflowError(RuntimeError): class WorkflowPaused(Exception): """Sinaliza que a Run pausou em um human_pick. State já persistido.""" - def __init__(self, run_id: int, step_id: str, prompt_to_user: str, options: list[str]): + def __init__(self, run_id: int, step_id: str, prompt_to_user: str, options: list[Any]): self.run_id = run_id self.step_id = step_id self.prompt_to_user = prompt_to_user @@ -152,16 +159,73 @@ def resume( row = self.tracker.query("SELECT * FROM runs WHERE id = ?", (run_id,))[0] state = json.loads(row["state"] or "{}") inputs = json.loads(row["inputs"] or "{}") - # injeta user_input no step pendente + # O session_id da run é o da linha, nunca o do parâmetro: quem chama pelo + # CLI não tem como saber qual era, e `session_id or 0` violava a FK de + # generations.session_id, matando a run inteira na primeira geração + # depois da pausa. + session_id = row["session_id"] if row["session_id"] is not None else session_id + state.setdefault("step_outputs", {}) + pending_step = state.get("pending_step") - if not pending_step: + if pending_step: + state["step_outputs"][pending_step] = self._resolver_escolha( + state, pending_step, user_input + ) + state["pending_step"] = None + state.pop("pending_options", None) + state.pop("pending_prompt", None) + elif row["status"] == "failed": + # Retomada de falha: os passos já concluídos continuam em + # step_outputs e são pulados. Existe porque um 429 do provider + # depois de uma escolha humana não pode custar a escolha. + pass + else: raise WorkflowError(f"Run {run_id} não está pausada") - state["step_outputs"][pending_step] = user_input - state["pending_step"] = None + + if session_id is not None: + self.tracker.set_session_awaiting(session_id, False) self._execute( workflow, run_id, project_id, project_slug, session_id, inputs, state ) + @staticmethod + def _resolver_escolha( + state: dict[str, Any], pending_step: str, user_input: dict[str, Any] + ) -> Any: + """Resolve a resposta do humano contra as opções que ficaram no disco. + + Sem isto, `--selected` aceita qualquer coisa e o passo seguinte consome + um valor que nunca esteve em `runs.state`: a retomada passaria a provar + que o processo aceita entrada, não que o estado sobreviveu. + """ + opcoes = state.get("pending_options") or [] + if not opcoes: + # Passo sem `from` declarado: o valor do humano é o próprio output. + return user_input + escolhido = user_input.get("id", user_input.get("selected")) + if escolhido is None: + raise WorkflowError( + f"Step '{pending_step}' espera uma escolha: informe " + '{"id": }.' + ) + for op in opcoes: + if isinstance(op, dict) and op.get("id") == escolhido: + # Opção estruturada: o output do passo é o objeto inteiro, que + # veio do disco. É o que permite `{{ steps.escolha.prompt }}`. + return op + if op == escolhido: + # Opção escalar (o caso histórico: uma lista de asset_id). O + # output continua sendo o dicionário que o humano mandou, para + # `{{ steps.pick.selected }}` seguir resolvendo. + return user_input + disponiveis = [ + op.get("id") if isinstance(op, dict) else op for op in opcoes + ] + raise WorkflowError( + f"'{escolhido}' não está entre as opções da pausa em " + f"'{pending_step}': {disponiveis}" + ) + # --- core execution ------------------------------------------------------- def _execute( @@ -185,16 +249,35 @@ def _execute( if step["kind"] == "human_pick": # persiste pausa, levanta exceção - state["pending_step"] = step_id - self.tracker.update_run(run_id, status="paused", state=state) prompt = _interpolate(step.get("prompt_to_user", "Escolha"), ctx) from_list = _interpolate(step.get("from"), ctx) or [] + if isinstance(from_list, (str, bytes)) or not isinstance( + from_list, (list, tuple) + ): + # `list()` de uma string devolveria uma opção por caractere. + from_list = [from_list] + state["pending_step"] = step_id + state["pending_options"] = list(from_list) + state["pending_prompt"] = str(prompt) + self.tracker.update_run(run_id, status="paused", state=state) + if session_id is not None: + # Acende `precisa_voce` no kanban do Workbench: a coluna é + # derivada de sessions.awaiting_input (queries.ts:312). + self.tracker.set_session_awaiting(session_id, True, str(prompt)) raise WorkflowPaused(run_id, step_id, str(prompt), list(from_list)) params = _interpolate(step.get("params", {}), ctx) try: - result = self._run_step(step, params, project_id, project_slug, session_id, run_id, idx) + if step["kind"] == "reason": + result = self._run_reason_step( + step, ctx, project_id, session_id, run_id, idx + ) + else: + result = self._run_step(step, params, project_id, project_slug, session_id, run_id, idx) except Exception as e: + # `pending_step` fica NULL, mas `step_outputs` é preservado e o + # resume aceita run falhada: uma falha transitória do provider + # não pode apagar o trabalho já pago nem a escolha do humano. self.tracker.update_run( run_id, status="failed", state=state, finished=True ) @@ -269,17 +352,80 @@ def _run_step( ) outputs_decl = step.get("outputs", {}) or {} - outputs: dict[str, Any] = {} - # decl tipo "assets: candidates.assets" — mapeamos primary - for key in outputs_decl: - if kind == "image" and key in ("assets", "asset"): - outputs[key] = asset_ids if key == "assets" else asset_ids[0] - elif kind in ("video", "upscale") and key in ("asset", "assets"): - outputs[key] = asset_ids[0] if key == "asset" else asset_ids - if not outputs: - outputs = {"assets": asset_ids, "asset": asset_ids[0] if asset_ids else None} + desconhecidas = [k for k in outputs_decl if k not in ("asset", "assets")] + if desconhecidas: + # Antes estas chaves eram descartadas em silêncio, e o passo + # seguinte lia None sem que nada acusasse. Passo de mídia só sabe + # produzir asset: quem precisa de estrutura usa `kind: reason`. + raise WorkflowError( + f"Step '{step['id']}' ({kind}) declara outputs que um passo de " + f"mídia não produz: {desconhecidas}. Use kind: reason." + ) + outputs: dict[str, Any] = { + "assets": asset_ids, + "asset": asset_ids[0] if asset_ids else None, + } return StepResult(step_id=step["id"], outputs=outputs, cost_brl=cost) + def _run_reason_step( + self, + step: dict[str, Any], + ctx: dict[str, Any], + project_id: int, + session_id: int | None, + run_id: int, + step_index: int, + ) -> StepResult: + """Passo cujo resultado é estrutura, não arquivo. + + Grava uma linha em `generations` e NENHUM asset. A linha entra no ledger + por três motivos medidos: enquanto ela está `pending` o card acende + `em_geracao` no kanban (queries.ts:281), o que cobre esperas longas de + motor; `created_at`/`finished_at` dão latência por passo, que hoje não + existe em lugar nenhum; e `cost_brl` ganha onde morar quando o custo de + token passar a ser medido. Nenhum asset é criado porque asset falso + acende `precisa_voce` e infla a contagem de entrega (queries.ts:272,313). + """ + motor = step.get("motor", "template") + if motor != "template": + raise WorkflowError( + f"Step '{step['id']}': motor '{motor}' não existe neste runner. " + "O v1 só conhece `template`." + ) + model = f"{motor}/{step.get('funcao') or step['id']}" + gen_id = self.tracker.create_generation( + project_id=project_id, + session_id=session_id, + model=model, + kind="reason", + prompt=None, + params={"motor": motor, "step": step["id"]}, + run_id=run_id, + step_index=step_index, + provider=motor, + ) + try: + outputs = _interpolate(step.get("outputs", {}) or {}, ctx) + if not isinstance(outputs, dict): + raise WorkflowError( + f"Step '{step['id']}': `outputs` de um passo reason tem de " + "ser um mapa de chaves." + ) + 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": ["passo de raciocínio: tokens não medidos"], + }, + ) + return StepResult(step_id=step["id"], outputs=outputs, cost_brl=0.0) + def _finalize( self, workflow: WorkflowSpec, project_slug: str, state: dict[str, Any] ) -> None: diff --git a/templates/apps/eco.yaml b/templates/apps/eco.yaml new file mode 100644 index 0000000..ba885ca --- /dev/null +++ b/templates/apps/eco.yaml @@ -0,0 +1,86 @@ +# Eco: o menor app que fecha o contrato do Workbench. +# +# Você dá uma palavra e um tom, o Eco monta três frases, para e pergunta qual +# delas vira imagem, e desenha a escolhida. É bobo de propósito: nenhuma +# inteligência, nenhuma chamada de modelo no meio, as três frases são template +# com substituição. O que ele faz de sério é atravessar o contrato inteiro, +# antes do Crystal Ball entrar em cima. +# +# O bloco `ui` é lido só pelo Workbench. WorkflowSpec.from_yaml só pega as +# chaves que conhece, então o runner ignora esse bloco sem nenhuma mudança: é +# um arquivo, dois leitores, nenhuma ponte. + +schema: studiolocal/workflow/v1 +slug: eco +name: Eco +description: Uma palavra vira três frases, você escolhe uma, ela vira imagem. + +ui: + versao: 1 + categoria: imagens + resumo: O menor app possível que exercita o contrato de ponta a ponta. + escopo: cliente-campanha + campos: + - { id: palavra, tipo: texto, rotulo: Palavra, obrigatorio: true, placeholder: girassol } + - { id: tom, tipo: escolha, apresentacao: select, rotulo: Tom, obrigatorio: true, + opcoes: [ensolarado, sombrio, neon], default: ensolarado } + - { id: quantos, tipo: numero, rotulo: Quantos desenhos, obrigatorio: true, min: 1, max: 2, default: 1 } + - { id: observacao, tipo: texto, linhas: 3, rotulo: Observação, obrigatorio: false } + pausa: + step: escolha + apresentacao: cartoes + rotulo_por_opcao: rotulo + # O custo NÃO é declarado aqui: um manifesto que declara o próprio preço pode + # mentir na tela de revisão. A tela lê o catálogo de modelos. + custo: + unidade: brl + modelo_do_passo: desenho + resultado: + tipo: galeria_da_sessao + +inputs: + palavra: { type: string, required: true } + tom: { type: string, required: true } + quantos: { type: int, required: true } + observacao: { type: string, required: false } + +steps: + # Passo de raciocínio: o resultado é estrutura, não arquivo. Grava em + # runs.state e uma linha em generations sem nenhum asset. + - id: eco + kind: reason + motor: template + outputs: + palavra: "{{ inputs.palavra }}" + opcoes: + - id: amanhecer + rotulo: "{{ inputs.palavra }} ao amanhecer" + prompt: "uma foto de {{ inputs.palavra }} ao amanhecer, tom {{ inputs.tom }}" + - id: chuva + rotulo: "{{ inputs.palavra }} sob chuva" + prompt: "uma foto de {{ inputs.palavra }} sob chuva, tom {{ inputs.tom }}" + - id: letreiro + rotulo: "{{ inputs.palavra }} em letreiro" + prompt: "uma foto de {{ inputs.palavra }} em letreiro, tom {{ inputs.tom }}" + + - id: escolha + kind: human_pick + prompt_to_user: "Qual eco de '{{ inputs.palavra }}' vira imagem?" + from: "{{ steps.eco.opcoes }}" + + # O passo que custa dinheiro lê DUAS chaves que só existem em runs.state: + # `steps.escolha.prompt`, que o runner resolveu contra a lista persistida, e + # `steps.eco.palavra`, que nenhum humano digita na retomada. Sem isso, a prova + # de reinício provaria que o processo aceita entrada, não que o estado voltou. + - id: desenho + kind: image + model: nano-banana-pro-google + params: + prompt: "{{ steps.escolha.prompt }} (eco de {{ steps.eco.palavra }})" + num_outputs: "{{ inputs.quantos }}" + image_size: "1K" + outputs: + assets: desenho.assets + +finalize: + promote_to_library: [] diff --git a/tests/test_workflow_app_steps.py b/tests/test_workflow_app_steps.py new file mode 100644 index 0000000..63bf0c0 --- /dev/null +++ b/tests/test_workflow_app_steps.py @@ -0,0 +1,327 @@ +"""Testes dos passos que um app declarado precisa, e da retomada. + +Cobre o que a Etapa 1 do Crystal Ball (PROD-2127) exige do runner: +- `reason`: passo cujo resultado é estrutura, grava em runs.state, zero assets +- a pausa acende `sessions.awaiting_input`, que é o que faz o card ir para + "precisa de você" no kanban do Workbench +- a retomada resolve a escolha contra as opções PERSISTIDAS, e recusa o que não + estava lá. Sem isso, o passo seguinte consome um valor que nunca esteve no + estado, e a prova de reinício provaria só que o processo aceita entrada +- a retomada herda `session_id` da linha da run. Passar None violava a FK de + generations.session_id e matava a run inteira +- falha de step não apaga o trabalho já pago: a run falhada é retomável +""" + +from __future__ import annotations + +import json +from pathlib import Path + +import pytest +import yaml + +from lib.asset_store import AssetStore +from lib.models_registry import ModelsRegistry +from lib.provider_base import ProviderError, ProviderOutput, ProviderResult +from lib.tracker import Tracker +from lib.workflow_runner import ( + WorkflowError, + WorkflowPaused, + WorkflowRunner, + WorkflowSpec, +) + + +class FakeProvider: + def __init__(self, calls: list[dict], falhar: bool = False): + self.calls = calls + self.falhar = falhar + + def call(self, model_name: str, params: dict) -> ProviderResult: + self.calls.append({"model": model_name, "params": params}) + if self.falhar: + raise ProviderError("429 do provider, transitório") + n = int(params.get("num_outputs", 1)) + outs = [ + ProviderOutput(url=f"https://fake/{model_name}/{len(self.calls)}/{i}.bin") + for i in range(n) + ] + return ProviderResult(job_id="req-1", outputs=outs, raw_response={}) + + +class FakeRouter: + def __init__(self, falhar: bool = False): + self.calls: list[dict] = [] + self._provider = FakeProvider(self.calls, falhar) + + def for_model(self, model_name: str): + return self._provider + + +def _fake_save_output(output: ProviderOutput, dest: Path) -> int: + dest.parent.mkdir(parents=True, exist_ok=True) + dest.write_bytes(b"fake") + return 4 + + +@pytest.fixture +def amb(tmp_path: Path, monkeypatch): + import lib.workflow_runner as wr + + monkeypatch.setattr(wr, "save_output", _fake_save_output) + tracker = Tracker(tmp_path / "tracker.db", create=True) + tracker.apply_migrations() + store = AssetStore(tmp_path) + router = FakeRouter() + runner = WorkflowRunner(tracker, ModelsRegistry(), router, store) + project_id = tracker.create_project("teste-app", "Teste App", [], None) + session_id = tracker.open_session(project_id) + return { + "runner": runner, + "tracker": tracker, + "store": store, + "router": router, + "project_id": project_id, + "project_slug": "teste-app", + "session_id": session_id, + "tmp": tmp_path, + } + + +def _spec_eco(tmp_path: Path, prompt_do_desenho: str | None = None) -> WorkflowSpec: + """O manifesto do Eco, na forma mínima que os testes precisam.""" + p = tmp_path / "eco.yaml" + p.write_text( + yaml.safe_dump( + { + "schema": "studiolocal/workflow/v1", + "slug": "eco", + "name": "Eco", + # bloco que só o Workbench lê: o runner tem de ignorar + "ui": {"versao": 1, "campos": [{"id": "palavra", "tipo": "texto"}]}, + "inputs": { + "palavra": {"type": "string", "required": True}, + "tom": {"type": "string", "required": True}, + }, + "steps": [ + { + "id": "eco", + "kind": "reason", + "motor": "template", + "outputs": { + "palavra": "{{ inputs.palavra }}", + "opcoes": [ + { + "id": "amanhecer", + "rotulo": "{{ inputs.palavra }} ao amanhecer", + "prompt": "foto de {{ inputs.palavra }} ao amanhecer, tom {{ inputs.tom }}", + }, + { + "id": "chuva", + "rotulo": "{{ inputs.palavra }} sob chuva", + "prompt": "foto de {{ inputs.palavra }} sob chuva, tom {{ inputs.tom }}", + }, + ], + }, + }, + { + "id": "escolha", + "kind": "human_pick", + "prompt_to_user": "Qual eco de '{{ inputs.palavra }}'?", + "from": "{{ steps.eco.opcoes }}", + }, + { + "id": "desenho", + "kind": "image", + "model": "nano-banana-pro-google", + "params": { + "prompt": prompt_do_desenho + or "{{ steps.escolha.prompt }} (eco de {{ steps.eco.palavra }})", + "num_outputs": 1, + "image_size": "1K", + }, + "outputs": {"assets": "desenho.assets"}, + }, + ], + "finalize": {"promote_to_library": []}, + } + ) + ) + return WorkflowSpec.from_yaml(p) + + +def _pausar(amb) -> tuple[WorkflowSpec, int, WorkflowPaused]: + spec = _spec_eco(amb["tmp"]) + wid = amb["tracker"].create_workflow("eco", "Eco", "eco.yaml") + with pytest.raises(WorkflowPaused) as exc: + amb["runner"].start( + spec, + project_id=amb["project_id"], + project_slug=amb["project_slug"], + session_id=amb["session_id"], + inputs={"palavra": "girassol", "tom": "neon"}, + workflow_id=wid, + ) + run_id = exc.value.run_id + return spec, run_id, exc.value + + +def test_reason_grava_estrutura_em_runs_state_e_nao_cria_asset(amb): + _, run_id, _ = _pausar(amb) + estado = json.loads( + amb["tracker"].query("SELECT state FROM runs WHERE id=?", (run_id,))[0]["state"] + ) + eco = estado["step_outputs"]["eco"] + assert eco["palavra"] == "girassol" + assert [o["id"] for o in eco["opcoes"]] == ["amanhecer", "chuva"] + assert eco["opcoes"][1]["prompt"] == "foto de girassol sob chuva, tom neon" + + # o passo entrou no ledger, e entrou sem asset + gens = amb["tracker"].query("SELECT * FROM generations WHERE run_id=?", (run_id,)) + assert len(gens) == 1 + assert gens[0]["kind"] == "reason" + assert gens[0]["model"] == "template/eco" + assert gens[0]["cost_brl"] == 0.0 + assert gens[0]["session_id"] == amb["session_id"] + assert json.loads(gens[0]["params"])["cost_source"] == "nao-apurado" + assert amb["tracker"].query("SELECT * FROM assets") == [] + + +def test_pausa_acende_awaiting_input_e_persiste_as_opcoes(amb): + _, run_id, pausa = _pausar(amb) + sess = amb["tracker"].query( + "SELECT * FROM sessions WHERE id=?", (amb["session_id"],) + )[0] + assert sess["awaiting_input"] == 1 + assert sess["awaiting_reason"] == "Qual eco de 'girassol'?" + assert sess["awaiting_at"] is not None + + estado = json.loads( + amb["tracker"].query("SELECT state FROM runs WHERE id=?", (run_id,))[0]["state"] + ) + assert estado["pending_step"] == "escolha" + assert [o["id"] for o in estado["pending_options"]] == ["amanhecer", "chuva"] + assert len(pausa.options) == 2 + + +def test_retomada_recusa_escolha_que_nao_estava_no_estado(amb): + spec, run_id, _ = _pausar(amb) + with pytest.raises(WorkflowError) as exc: + amb["runner"].resume( + spec, run_id, amb["project_id"], amb["project_slug"], None, + {"id": "inventada"}, + ) + assert "inventada" in str(exc.value) + # e nada foi gerado por causa da tentativa + assert amb["router"].calls == [] + + +def test_retomada_le_a_escolha_do_disco_e_herda_a_sessao(amb): + spec, run_id, _ = _pausar(amb) + # session_id=None é exatamente o que o CLI passa: o runner tem de herdar da + # linha da run, senão a FK de generations.session_id derruba a run. + amb["runner"].resume( + spec, run_id, amb["project_id"], amb["project_slug"], None, {"id": "chuva"} + ) + run = amb["tracker"].query("SELECT * FROM runs WHERE id=?", (run_id,))[0] + assert run["status"] == "done" + + gen = amb["tracker"].query( + "SELECT * FROM generations WHERE run_id=? AND kind='image'", (run_id,) + )[0] + assert gen["session_id"] == amb["session_id"] + # o prompt do passo que custa dinheiro veio de duas chaves que só existiam + # em runs.state: a opção resolvida e a palavra do passo de raciocínio + assert gen["prompt"] == "foto de girassol sob chuva, tom neon (eco de girassol)" + assert gen["cost_brl"] > 0 + + sess = amb["tracker"].query( + "SELECT * FROM sessions WHERE id=?", (amb["session_id"],) + )[0] + assert sess["awaiting_input"] == 0 + + +def test_falha_de_step_nao_apaga_o_trabalho_ja_pago(amb, monkeypatch): + spec, run_id, _ = _pausar(amb) + amb["runner"].providers = FakeRouter(falhar=True) + with pytest.raises(WorkflowError): + amb["runner"].resume( + spec, run_id, amb["project_id"], amb["project_slug"], None, {"id": "chuva"} + ) + run = amb["tracker"].query("SELECT * FROM runs WHERE id=?", (run_id,))[0] + assert run["status"] == "failed" + estado = json.loads(run["state"]) + # a escolha do humano e o raciocínio continuam lá + assert estado["step_outputs"]["escolha"]["id"] == "chuva" + assert estado["step_outputs"]["eco"]["palavra"] == "girassol" + + # e a run falhada é retomável sem repetir a pergunta ao humano + amb["runner"].providers = FakeRouter() + amb["runner"].resume( + spec, run_id, amb["project_id"], amb["project_slug"], None, {} + ) + run = amb["tracker"].query("SELECT * FROM runs WHERE id=?", (run_id,))[0] + assert run["status"] == "done" + + +def test_passo_de_midia_que_declara_output_de_estrutura_falha_alto(amb): + p = amb["tmp"] / "ruim.yaml" + p.write_text( + yaml.safe_dump( + { + "schema": "studiolocal/workflow/v1", + "slug": "ruim", + "name": "Ruim", + "inputs": {}, + "steps": [ + { + "id": "gen", + "kind": "image", + "model": "nano-banana-pro-google", + "params": {"prompt": "x", "num_outputs": 1}, + # antes isto era descartado em silêncio e o passo + # seguinte lia None sem que nada acusasse + "outputs": {"dossie": "gen.dossie"}, + } + ], + } + ) + ) + spec = WorkflowSpec.from_yaml(p) + wid = amb["tracker"].create_workflow("ruim", "Ruim", "ruim.yaml") + with pytest.raises(WorkflowError) as exc: + amb["runner"].start( + spec, amb["project_id"], amb["project_slug"], amb["session_id"], {}, wid + ) + assert "dossie" in str(exc.value) + assert "kind: reason" in str(exc.value) + + +def test_from_string_nao_vira_uma_opcao_por_caractere(amb): + p = amb["tmp"] / "str.yaml" + p.write_text( + yaml.safe_dump( + { + "schema": "studiolocal/workflow/v1", + "slug": "str", + "name": "Str", + "inputs": {"d": {"type": "string"}}, + "steps": [ + { + "id": "pick", + "kind": "human_pick", + "prompt_to_user": "Confirma?", + "from": "{{ inputs.d }}", + } + ], + } + ) + ) + spec = WorkflowSpec.from_yaml(p) + wid = amb["tracker"].create_workflow("str", "Str", "str.yaml") + with pytest.raises(WorkflowPaused) as exc: + amb["runner"].start( + spec, amb["project_id"], amb["project_slug"], amb["session_id"], + {"d": "Abrir com o plano geral"}, wid, + ) + assert exc.value.options == ["Abrir com o plano geral"] From 28503cd6d0dfd54c9e41c51f40a2a7fd3eaf8fe2 Mon Sep 17 00:00:00 2001 From: davidbenal <144815978+davidbenal@users.noreply.github.com> Date: Sat, 22 Aug 2026 19:02:39 -0300 Subject: [PATCH 2/2] CI: instalar em venv, porque --system falha no runner MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit O Python do ubuntu-latest é externally managed, e `uv pip install --system` sai com exit 2 antes de instalar qualquer coisa. Passa a criar venv e rodar os três passos com `uv run`, da raiz do repositório, que é onde o editable resolve `from lib.X`. Co-Authored-By: Claude Opus 5 (1M context) --- .github/workflows/testes.yml | 14 ++++++++++---- 1 file changed, 10 insertions(+), 4 deletions(-) diff --git a/.github/workflows/testes.yml b/.github/workflows/testes.yml index 31cd43f..c3b725f 100644 --- a/.github/workflows/testes.yml +++ b/.github/workflows/testes.yml @@ -30,16 +30,22 @@ jobs: with: enable-cache: true + # Em venv, não --system: o Python do runner é externally managed e o + # --system falha com exit 2 antes de instalar qualquer coisa. - name: Instala com os extras de dev - run: uv pip install --system -e ".[dev]" + run: | + uv venv + uv pip install -e ".[dev]" - name: Gate do ruff (sintaxe, nome indefinido, import morto) - run: ruff check . --select E9,F --output-format concise + run: uv run ruff check . --select E9,F --output-format concise - name: Dívida de estilo, declarada e não bloqueante run: | echo "Achados de estilo herdados (não reprovam este job):" - ruff check . --statistics --exit-zero + uv run ruff check . --statistics --exit-zero + # Da raiz do repo: os testes importam `from lib.X`, e o editable depende + # do cwd para resolver o pacote. - name: Testes - run: pytest -q + run: uv run pytest -q