Lab Notes
Research

Mission DAG

El DAG de fases del research framework: qué fases se ejecutan, en qué orden, qué produce cada una y de qué depende.

Propósito

Esta página es la referencia canónica del grafo de fases. El DAG es la estructura que ejecuta el orquestador. El framework de alto nivel está en Framework; el ciclo de vida de las fases (input, output, retry, failure) está en Mission Lifecycle.

Topología

Loading diagram…
Mission DAG: F0 siembra F1, F1 reparte a F2/F4/Y1 y opcionalmente P3, Y1 alimenta Y2, todo alimenta F6.

Tres cosas a notar:

  • F1 es el punto de fan-out. Todas las demás fases dependen de los candidatos producidos por F1. Esto es deliberado: fuerza al framework a operar sobre nombres de candidatos reales y normalizados en lugar de sobre la query cruda.
  • F2, F4, Y1 se ejecutan en paralelo (parallel_group: 1). No comparten estado y pueden ejecutarse concurrentemente.
  • P3 es opcional. Solo se ejecuta si el flag enable_perplexity_synthesis está activo en la configuración del run.

Referencia de fases

F0 — Init

AspectoDetalle
Ficheroscout/phases/f0_init.py
Statussolo succeeded. Sin retry.
NecesitaNada (fase raíz).
Producetask.json (tipado: Task)
InputsLa query cruda del usuario (string) y el mission_id (string).
OutputsUn directorio de run bajo runs_dir/<mission_id>/ con task.json y un checkpoint phase.json.
ParalelismoSequential (siempre corre primero).
Regla de skipNunca se omite.

F0 parsea la query cruda en un ParsedQuery (topic, location, period, budget, constraints) usando un parser basado en keywords. La estructura parseada es lo que F1 usa para construir las queries de búsqueda de candidatos.

F1 — Candidates

AspectoDetalle
Ficheroscout/phases/f1_candidates.py
Statussucceeded, succeeded_partial o failed_terminal.
Necesitatask.json de F0.
Producecandidates.json (tipado: CandidateList)
Inputstask.json. La fase usa el topic y location parseados para construir queries de búsqueda.
OutputsUna lista de registros Candidate, cada uno con un Source. Deduplicados y normalizados.
ParalelismoSequential (tras F0).
Regla de skipNunca se omite.

F1 es la fase más importante. Su salida es sobre la que opera cada fase downstream. La fase llama a los adaptadores de búsqueda configurados en secuencia (Tavily → DuckDuckGo → Google), fusiona resultados, deduplica por URL y título normalizado, valida cada candidato contra el esquema Candidate y escribe el candidates.json tipado en el directorio del run.

F2 — Reviews

AspectoDetalle
Ficheroscout/phases/f2_reviews.py
Statussucceeded, succeeded_partial o failed_retryable.
Necesitacandidates.json de F1.
Producereviews.json (tipado: ReviewSummary)
InputsLa URL de cada candidato y su identificador primario.
OutputsUn ReviewSummary por candidato con agregados de sentimiento y temas.
ParalelismoParallel (parallel_group: 1).
Retry2 intentos.
Regla de skipSe omite si candidates.json está vacío o tiene menos de 1 candidato.

F4 — Price

AspectoDetalle
Ficheroscout/phases/f4_price.py
Statussucceeded, succeeded_partial o failed_retryable.
Necesitacandidates.json de F1.
Produceprices.json (tipado: PriceSeries)
InputsEl identificador de cada candidato y el rango de fechas del run.
OutputsUn PriceSeries por candidato con snapshots y estadísticas.
ParalelismoParallel (parallel_group: 1).
Retry2 intentos.
Regla de skipSe omite si la query no tiene restricción budget_max.
AspectoDetalle
Ficheroscout/phases/y1_youtube.py
Statussucceeded, succeeded_partial o failed_retryable.
Necesitacandidates.json de F1.
Produceyoutube_search.json (tipado: lista de metadatos de vídeo)
InputsEl nombre de cada candidato y el topic de la query.
OutputsUna lista de registros de vídeo de YouTube por candidato.
ParalelismoParallel (parallel_group: 1).
Retry1 intento (los rate limits de YouTube son estrictos).
Regla de skipSe omite si enable_youtube_research es false.

Y2 — YouTube Extract

AspectoDetalle
Ficheroscout/phases/y2_extract.py
Statussucceeded, succeeded_partial o failed_retryable.
Necesitayoutube_search.json de Y1.
Producetranscripts.json (tipado: lista de Transcript)
InputsLas URLs de YouTube de Y1.
OutputsTranscripciones y segmentos chunked por vídeo.
ParalelismoSequential (tras Y1).
Retry1 intento.
Regla de skipSe omite si Y1 no produjo resultados.

P3 — Perplexity Synthesis (Opcional)

AspectoDetalle
Ficheroscout/phases/p3_perplexity_synthesis.py
Statussucceeded o failed_terminal.
Necesitacandidates.json de F1.
Producesynthesis.json (tipado: lista de resúmenes narrativos)
InputsLos top N candidatos (default 5).
OutputsUn resumen narrativo por candidato escrito por Perplexity.
ParalelismoParallel (parallel_group: 1) si está habilitado.
Retry2 intentos.
Regla de skipSe omite si enable_perplexity_synthesis es false.

P3 es la única fase que usa un servicio de LLM de pago (Perplexity) directamente dentro del DAG. Es opt-in porque cuesta créditos de API. La salida es prosa narrativa que el informe final (F6) puede incluir verbatim o parafrasear. P3 está documentado por separado en External Providers porque el patrón de adaptador, el token pool y el request journal viven todos allí.

F6 — Report

AspectoDetalle
Ficheroscout/phases/f6_report.py
Statussolo succeeded. Sin retry.
NecesitaTodos los artefactos de fase (F1, F2, F4, Y2, P3).
Producereport.md (markdown, no tipado)
InputsTodos los artefactos bajo el directorio del run.
OutputsEl informe markdown final.
ParalelismoSequential (siempre corre el último).
Regla de skipNunca se omite.

F6 agrega todos los artefactos en un único informe markdown con una estructura estable: resumen ejecutivo, candidatos, reseñas, pricing, cobertura de vídeo, síntesis de fuentes (si P3 corrió) y una lista deduplicada de fuentes. El formato completo está en Artifacts → report.md.

Estados de fallo

StatusSignificadoAcción
pendingLa fase no ha empezado todavía.Esperar al orquestador.
runningLa fase se está ejecutando.Esperar.
succeededLa fase se completó y produjo todas las salidas esperadas.Marcar las fases downstream como listas.
succeeded_partialLa fase se completó pero faltan algunas salidas.Marcar downstream como listas; marcar en el informe.
failed_retryableLa fase falló pero se puede reintentar.El orquestador reintenta.
failed_terminalLa fase falló y no se puede reintentar.El orquestador detiene la misión; notifica al usuario.
skippedLa fase no se ejecutó intencionadamente.Marcar downstream como blocked_missing_input.
blocked_missing_inputLa fase no pudo empezar porque falta un input.Omitir; marcar downstream.
staleEl input de la fase se actualizó después de que la fase corriera.Re-ejecutar.

Ejecución paralela

Las fases paralelas (F2, F4, Y1) se lanzan con un ThreadPoolExecutor. El número de workers por defecto es SCOUTE_PARALLEL_WORKERS (default 3). El runner coordina de forma que todas las fases paralelas arrancan cuando F1 termina, cada fase paralela escribe su propio artefacto de forma independiente, y el runner espera a que todas las fases paralelas alcancen un estado terminal antes de arrancar Y2 (si Y1 fue paralela) y F6.

El flag allow_partial de cada fase controla si se acepta una completion parcial:

Faseallow_partialJustificación
F0falseSin F0, no hay run.
F1trueAlgunos candidatos son mejor que ninguno.
F2trueLas reseñas son nice-to-have, no críticas.
F4trueEl pricing es crítico pero parcial es mejor que nada.
Y1trueLa cobertura de YouTube es opcional.
Y2trueY2 no puede succeed sin Y1.
P3trueLa síntesis es opcional.
F6falseEl informe siempre se produce si la misión corre.

Checkpointing

El DAG es totalmente resumible. Tras cada fase, el StateStore escribe el estado actual a <run_dir>/phase.json. Si el orquestador crashea, la siguiente llamada a run_mission cargará el estado, identificará la siguiente fase a ejecutar según el DAG y el estado, y continuará desde ahí. El modelo completo de checkpointing está en Mission Lifecycle → Checkpointing.

Construcción del DAG

El DAG se construye en scout/orchestrator/dag.py. El DAG por defecto se construye en tiempo de import. Se pueden construir DAGs custom pasando una lista de PhaseConfig a PhaseDAG(...).

from scout.orchestrator.dag import PhaseConfig, PhaseDAG, PhaseKind
 
dag = PhaseDAG([
    PhaseConfig(name="F0", needs=[], produces=["task.json"]),
    PhaseConfig(name="F1", needs=["task.json"], produces=["candidates.json"]),
    PhaseConfig(name="F2", needs=["candidates.json"], produces=["reviews.json"],
                parallel_group=1, retry=2, allow_partial=True),
    PhaseConfig(name="F4", needs=["candidates.json"], produces=["prices.json"],
                parallel_group=1, retry=2, allow_partial=True),
    PhaseConfig(name="Y1", needs=["candidates.json"], produces=["youtube_search.json"],
                parallel_group=1, retry=1, allow_partial=True),
    PhaseConfig(name="Y2", needs=["youtube_search.json"], produces=["transcripts.json"],
                retry=1, allow_partial=True),
    PhaseConfig(name="P3", needs=["candidates.json"], produces=["synthesis.json"],
                parallel_group=1, retry=2, allow_partial=True),
    PhaseConfig(name="F6", needs=["reviews.json", "prices.json", "transcripts.json",
                                  "synthesis.json"], produces=["report.md"]),
])

Los DAGs custom no se usan en producción pero son útiles para testear fases individuales o para flujos de investigación que necesiten una forma distinta.

Ver también

On this page