Flujo de trabajo asistido por IA para cambios en bases de datos (parte 1 de 2)
Este manual describe un flujo de trabajo de cambio de base de datos asistido por IA de principio a fin, utilizando el SDK de OpenAI Agents.
Demuestra cómo el ecosistema de herramientas de OpenAI se puede aplicar para orquestar flujos de trabajo complejos y con gran cantidad de datos en infraestructuras empresariales modernas. Si bien la implementación actual se centra en un caso de uso de cambio de esquema y análisis de impacto orientado al comercio minorista, los patrones arquitectónicos subyacentes son agnósticos al dominio y extensibles. El mismo diseño de flujo de trabajo se puede adaptar a industrias como la manufactura, productos farmacéuticos, atención médica, logística, finanzas y operaciones de la cadena de suministro, donde se requieran flujos de trabajo de datos estructurados, razonamiento operativo, análisis aumentado por recuperación y validación automatizada.
El ejemplo en ejecución es un cambio de nivel de lealtad minorista, pero el mismo patrón se aplica a muchas solicitudes de cambio de base de datos donde los equipos necesitan un análisis de impacto rastreable y un resultado de implementación revisable.
El flujo de trabajo comienza con una solicitud de cambio de base de datos en lenguaje natural, la convierte en JSON estructurado, opcionalmente fundamenta el análisis de impacto con contexto de búsqueda de archivos basado en PDF, genera un plan de implementación seguro, redacta SQL en todas las capas de la plataforma de datos, valida la salida con barreras de seguridad deterministas, guarda un artefacto reutilizable y, opcionalmente, evalúa el flujo con Promptfoo.
El notebook es intencionalmente autocontenido: toda la lógica central del flujo de trabajo, los prompts, las barreras de seguridad, la generación de artefactos y los archivos de tiempo de ejecución de evaluación se crean a partir de las celdas del notebook.
Descripción general
Los cambios de esquema son engañosamente simples. Una solicitud como "agregar una columna anulable y rellenarla" puede afectar las tablas de destino, los modelos de preparación, las tablas dimensionales, los marts, la lógica de informes, las suposiciones de linaje, las comprobaciones de validación, los procedimientos de reversión y la secuencia de lanzamiento.
Los ejemplos utilizan datos de clientes minoristas porque las dependencias son fáciles de ver, pero los mismos tipos de traspasos aparecen en muchos equipos de análisis y plataforma.
Este manual demuestra un patrón práctico para usar agentes como asistente de análisis de cambios e implementación para el trabajo de ingeniería de bases de datos. En lugar de pedirle a un modelo que produzca un script SQL final directamente, el flujo de trabajo divide la tarea en etapas explícitas:
- Analiza la solicitud en lenguaje natural en JSON estructurado.
- Analiza los objetos impactados y los riesgos operativos.
- Crea un plan de implementación con verificaciones previas, posteriores y orientación para la reversión.
- Genera SQL en todas las capas de la plataforma.
- Ejecuta verificaciones de cordura deterministas.
- Guarda un artefacto legible por máquina.
- Opcionalmente, ejecuta evaluaciones de Promptfoo contra el flujo actual.
El resultado no es solo un script SQL generado. Es un paquete auditable que contiene la solicitud de cambio interpretada, el análisis de impacto, el plan, el SQL, los resultados de validación, los resúmenes opcionales de evidencia RAG y los resultados de la evaluación.
Por qué esto es importante
Las solicitudes de cambio de base de datos a menudo pasan por varias transferencias: los propietarios de productos describen la necesidad, los ingenieros de datos la interpretan, los equipos de plataforma evalúan el riesgo, los ingenieros de análisis propagan el campo aguas abajo y los revisores verifican si el cambio es seguro. El contexto importante puede perderse en cada paso.
SchemaFlow aborda esto convirtiendo una solicitud de cambio de formato libre en un flujo de trabajo estructurado e inspeccionable.
Esto es importante porque los cambios en la base de datos pueden crear modos de falla ocultos:
- Una columna agregada a ODS puede no propagarse a staging, core o marts.
- Un campo anulable puede generarse accidentalmente como
NOT NULL. - La lógica de relleno puede omitirse aunque la solicitud pida una población histórica.
- Los requisitos de índice pueden pasarse por alto.
- Las dependencias de informes posteriores pueden ser desconocidas a menos que se consulte la documentación de referencia.
- El SQL generado puede parecer plausible pero fallar en las comprobaciones básicas de coherencia.
Este manual muestra un patrón para reducir esos riesgos con razonamiento de agente por etapas, salidas tipadas, contexto de recuperación opcional, barreras de seguridad deterministas, artefactos guardados y evaluaciones repetibles.
Beneficios clave
- Interpretación estructurada – Convierte las solicitudes de bases de datos en lenguaje natural en un contrato
change_jsonnormalizado. - Separación de responsabilidades – Utiliza agentes especializados para el análisis, el análisis de impacto, la planificación de la implementación y la generación de SQL.
- Fundamentación RAG opcional – Permite que el agente de análisis de impacto use la búsqueda de archivos sobre un PDF cargado, como un IFD, una especificación de esquema o un documento de linaje.
- Salidas de etapa tipadas – Utiliza modelos Pydantic y esquemas de salida del SDK de Agents para las etapas de análisis, impacto y planificación.
- Flujo de trabajo primero con barreras de seguridad – Agrega verificaciones deterministas entre etapas para que las fallas obvias se detecten antes de que los pasos posteriores consuman un estado incorrecto.
- Trazabilidad – Emite trazas y spans del SDK de OpenAI Agents para ejecuciones de agentes, barreras de seguridad, generación de artefactos y ejecución de evaluaciones.
- Artefactos portátiles – Guarda el paquete de flujo de trabajo final como JSON en
artifacts/notebook_runs/. - Diseño listo para evaluación – Genera archivos de proveedor, aserción, configuración y resultados de Promptfoo a partir del estado del notebook en vivo.
- Sin efectos secundarios en la base de datos – Produce un borrador de SQL y una salida de validación sin ejecutar contra una base de datos en vivo.
Lo que construirás
Al final de este notebook, tendrás una pipeline de SchemaFlow funcional que produce:
Una solicitud de cambio de base de datos analizada:
- título
- dominio
- esquema de destino
- tabla de destino
- operaciones normalizadas
- notas
Un informe de análisis de impacto:
- tablas, columnas, índices, vistas o relaciones impactadas
- riesgos
- suposiciones
- resúmenes opcionales de evidencia de búsqueda de archivos
Un plan de implementación:
- pasos de implementación
- verificaciones previas
- verificaciones posteriores
- acciones de reversión
Un script SQL borrador con cuatro secciones requeridas:
LANDING (ODS)STAGING (STG)CORE (DIM/FACT/VIEW)MARTS (SERVING)
Un resultado de validación:
- verificaciones de tabla esperadas
- verificaciones de columna esperadas
- verificaciones de palabras clave requeridas como
ALTER TABLE,UPDATEoCREATE INDEX
Un artefacto JSON guardado:
- solicitud de cambio
- análisis de impacto
- plan
- SQL
- validación
- metadatos RAG opcionales
Un arnés de evaluación de Promptfoo:
- proveedor de Python
- archivo de aserción de Python
- configuración de Promptfoo generada
- caso de evaluación solo de análisis
- caso de evaluación de flujo completo
- informes de evaluación JSON y HTML con marca de tiempo
Introducción: Caso de uso y solución
Este manual se centra en un escenario común de ingeniería de datos empresariales: un interesado solicita un cambio de esquema de base de datos en lenguaje natural, y el equipo de datos necesita convertir esa solicitud en un plan listo para la implementación.
Aquí, el dominio minorista es solo una forma concreta de hacer que el flujo de trabajo sea tangible. El mismo enfoque por etapas se puede adaptar a otros sistemas de origen, productos de datos y procesos de revisión.
La solicitud predeterminada en este notebook es:
Add LOYALTY_TIER VARCHAR(20) to ODS.ODS_CUSTOMER_PROFILE as nullable.
Backfill from CORE.DIM_CUSTOMER on CUSTOMER_ID where IS_CURRENT=true.
Add a non-unique index on (CUSTOMER_ID, LOYALTY_TIER).
Un ingeniero de datos humano normalmente necesitaría responder varias preguntas antes de escribir SQL de producción:
- ¿Qué tabla y esquema se están cambiando?
- ¿Qué columna, tipo y anulabilidad exactos se solicitaron?
- ¿Se requiere un relleno histórico?
- ¿La solicitud implica un índice?
- ¿Qué capas posteriores necesitan que se propague el campo?
- ¿Qué riesgos deben buscar los revisores?
- ¿Qué comprobaciones deben ejecutarse antes y después de la implementación?
- ¿Qué pasos de reversión son razonables?
- ¿El SQL generado incluye los elementos requeridos?
SchemaFlow implementa esto como un flujo de trabajo de agente por etapas. Cada etapa crea una salida intermedia tipada que la siguiente etapa consume. Las comprobaciones deterministas validan las salidas antes de que el notebook guarde el paquete final y, opcionalmente, ejecute las evaluaciones.
Descripción general del flujo de trabajo
A grandes rasgos, SchemaFlow sigue esta secuencia:

El notebook está organizado para que los lectores puedan ejecutar primero el flujo de trabajo central y luego decidir si quieren ejecutar la sección de evaluación opcional de Promptfoo.
Tabla de Contenidos
Guía Conceptual
- Descripción general
- Por qué esto es importante
- Beneficios clave
- Lo que construirás
- Introducción: Caso de uso y solución
- Descripción general del flujo de trabajo
- Arquitectura - Patrones de diseño
- Diseño del sistema
- Flujo de trabajo de ejecución
Implementación del Notebook
- Configuración del entorno
- Entrada
- Contexto RAG de PDF opcional
- Etapas 1-2: Análisis de la solicitud de cambio + Análisis de impacto
- Etapas 3-4: Plan de ejecución + Generación de SQL
- Etapa 5: Comprobaciones ligeras de cordura de SQL
- Paquete final
- Guardar artefacto
- Limpieza opcional
- Evaluar el flujo con Promptfoo
- Configuración del directorio de tiempo de ejecución de Promptfoo
- Comprobación de tiempo de ejecución de Node.js y npm
- Publicar el tiempo de ejecución central de SchemaFlow
- Tiempo de ejecución del proveedor de Promptfoo
- Tiempo de ejecución de aserción de Promptfoo
- Construir casos de prueba y configuración de Promptfoo
- Ejecutar evaluación de Promptfoo
- Revisar los últimos resultados de Promptfoo
Referencia
Arquitectura - Patrones de diseño
SchemaFlow utiliza una arquitectura de agente por etapas y basada en contratos. El objetivo es evitar tratar el modelo como un único generador de SQL de caja negra. En cambio, cada etapa tiene una responsabilidad limitada y produce una salida que puede ser inspeccionada, validada, rastreada y reutilizada.
1. Especialización del agente
Cada agente realiza una tarea principal:
| Agente | Responsabilidad | Salida principal |
|---|---|---|
| Agente de análisis | Extrae campos estructurados de la solicitud en lenguaje natural | change_json |
| Agente de impacto | Identifica objetos afectados, suposiciones y riesgos | impact_json |
| Agente de plan | Convierte el cambio y el impacto en pasos de implementación | plan_json |
| Agente SQL | Redacta SQL en todas las capas de la plataforma de datos | sql_text |
Esta especialización facilita la depuración del flujo de trabajo. Si a SQL le falta una columna, puedes inspeccionar si el problema comenzó en el análisis, el análisis de impacto, la planificación o la generación de SQL.
2. Contratos de salida tipados
El notebook define modelos Pydantic para las etapas estructuradas:
ChangeRequestModelImpactModelPlanModel
Esos modelos se envuelven con AgentOutputSchema para que el SDK de Agents conozca la forma de salida esperada. El flujo de trabajo también normaliza las salidas después de las llamadas al modelo para asegurar que existan las claves esperadas antes de que se ejecuten las etapas posteriores.
3. Análisis de impacto aumentado por recuperación
La sección RAG de PDF es opcional. Cuando se establece PDF_PATH, el notebook:
- Crea un almacén de vectores de OpenAI.
- Carga el PDF.
- Permite que OpenAI lo analice, fragmente, incruste e indexe.
- Le da al Agente de Impacto un
FileSearchTool. - Captura un resumen de los resultados de la búsqueda de archivos devueltos.
Esto es útil cuando la solicitud de cambio necesita fundamentarse en un IFD, un documento de esquema, un archivo de linaje, un contrato de datos o una referencia de arquitectura.
4. Puertas de seguridad entre etapas
El notebook añade comprobaciones deterministas después de las etapas principales:
- Las barreras de seguridad de las etapas 1-2 validan las salidas de análisis e impacto.
- Las barreras de seguridad de las etapas 3-4 validan la completitud del plan, la propagación del tipo de datos y el manejo de la anulabilidad.
- Las comprobaciones SQL de la etapa 5 validan la presencia de la tabla, columna y palabras clave SQL esperadas.
- Las comprobaciones posteriores al artefacto verifican que el artefacto JSON guardado existe y se completa el ciclo.
- Las comprobaciones previas a Promptfoo verifican que el estado del notebook esté listo para las evaluaciones.
Estas comprobaciones no reemplazan la revisión humana, pero detectan fallas silenciosas comunes a tiempo.
5. Ejecución centrada en artefactos
El paquete final es el artefacto principal del flujo de trabajo. Captura el estado necesario para revisar o depurar la ejecución:
bundle = {
"summary": ...,
"rag": ...,
"change_json": ...,
"impact_json": ...,
"plan": ...,
"sql": ...,
"validation": ...
}
El notebook guarda este paquete en artifacts/notebook_runs/.
6. Tiempo de ejecución de evaluación generado a partir del estado del notebook
Promptfoo se ejecuta en un proceso separado, por lo que no puede leer directamente las variables del kernel activo del notebook. Para resolver esto, la Sección 10 escribe un pequeño módulo Python reutilizable y archivos de tiempo de ejecución de Promptfoo a partir del estado actual del notebook.
Esto asegura que las ediciones de prompts, las ediciones de CHANGE_TEXT y los cambios de configuración del modelo se reflejen cuando se regeneran los archivos de evaluación.
Diseño del sistema
Arquitectura de componentes

Objetos principales en tiempo de ejecución
| Objeto | Creado en | Propósito |
|---|---|---|
CHANGE_TEXT |
Sección de entrada | La solicitud de cambio de base de datos en lenguaje natural |
change_json |
Etapa 1 | Interpretación estructurada de la solicitud |
rag_vector_store_id |
Sección opcional de RAG de PDF | ID del almacén de vectores alojado para el contexto PDF cargado |
rag_file_search_results |
Etapa 2 | Resumen de los resultados de la búsqueda de archivos devueltos al Agente de Impacto |
impact_json |
Etapa 2 | Objetos impactados, riesgos y suposiciones |
plan_json |
Etapa 3 | Plan de implementación, verificaciones y guía de reversión |
sql_text |
Etapa 4 | Borrador de script SQL |
validation |
Etapa 5 | Resultado de la verificación de cordura determinista de SQL |
bundle |
Sección del paquete final | Salida consolidada del flujo de trabajo |
out_path |
Sección Guardar artefacto | Ruta del artefacto JSON guardado |
promptfoo_config |
Sección Promptfoo | Configuración de evaluación generada |
Límite importante
SchemaFlow genera borradores de artefactos de implementación. No ejecuta SQL contra una base de datos, aplica migraciones, abre solicitudes de extracción ni modifica sistemas de producción.
Flujo de trabajo de ejecución
Ejecuta el notebook en orden.
Flujo de trabajo central
Configuración del entorno
- Importa las dependencias.
- Verifica la versión del SDK de OpenAI Agents.
- Lee
OPENAI_API_KEY. - Configura el rastreo y la selección del modelo.
Entrada
- Define
CHANGE_TEXT. - Esta es la única entrada de negocio requerida para el flujo de trabajo central.
- Define
Contexto RAG de PDF opcional
- Deja
PDF_PATH = Nonepara ejecutar sin recuperación. - Establece
PDF_PATHen un PDF local para habilitar el contexto de búsqueda de archivos para el análisis de impacto.
- Deja
Etapas 1-2
- Analiza la solicitud de cambio.
- Analiza el impacto.
- Opcionalmente, usa la búsqueda de archivos durante el análisis de impacto.
Barreras de seguridad de las etapas 1-2
- Confirma que la salida del análisis está bien formada.
- Confirma que la salida del impacto incluye el objetivo.
- Confirma que los objetos impactados contienen los campos requeridos.
Etapas 3-4
- Genera un plan de ejecución.
- Genera SQL en las capas de aterrizaje, preparación, núcleo y mart.
Barreras de seguridad de las etapas 3-4
- Confirma que las secciones del plan están pobladas.
- Confirma la propagación del tipo de datos.
- Confirma que el comportamiento de anulabilidad coincide con la solicitud.
Etapa 5 Comprobaciones de cordura de SQL ligeras
- Verifica si el SQL está vacío.
- Verifica la tabla y columnas de destino esperadas.
- Verifica las acciones SQL requeridas implícitas en la solicitud.
Paquete y artefacto final
- Ensambla el paquete de salida completo.
- Guárdalo como JSON.
- Verifica que el artefacto se completa con éxito.
Flujo de trabajo de evaluación opcional
Comprobaciones previas a Promptfoo
- Confirma que el estado del notebook está listo para las evaluaciones.
Generación del tiempo de ejecución de Promptfoo
- Crea un módulo central de SchemaFlow reutilizable.
- Escribe un proveedor de Promptfoo.
- Escribe un archivo de aserción de Promptfoo.
- Genera casos de prueba y configuración de Promptfoo.
Ejecución de la evaluación de Promptfoo
- Ejecuta evaluaciones de solo análisis y de flujo completo.
- Guarda informes JSON y HTML con marca de tiempo.
- Actualiza los alias de
schemaflow_cookbook_eval_latest.*.
1) Configuración del entorno
Esta sección prepara el entorno de ejecución para el flujo de trabajo de SchemaFlow.
La celda de configuración hace lo siguiente:
- Importa las utilidades estándar de Python utilizadas en todo el notebook.
- Importa el cliente de OpenAI.
- Importa las primitivas del SDK de OpenAI Agents:
AgentRunnerRunConfigAgentOutputSchemaFileSearchTool- ayudantes de rastreo y span
- Verifica que el paquete
openai-agentsinstalado cumpla con la versión mínima requerida. - Lee
OPENAI_API_KEYdel entorno o lo solicita. - Establece el modelo con
OPENAI_MODEL, con un valor predeterminado degpt-5.5. - Crea un ID de grupo de rastreo para que todas las ejecuciones de agentes y spans de barreras de seguridad relacionadas puedan agruparse.
El flujo de trabajo habilita intencionalmente cargas útiles de rastreo sensibles para esta demostración, de modo que los prompts, las salidas, los paquetes de evaluación y los datos de las herramientas sean visibles en los rastreos. Para uso en producción, revisa esta configuración antes de manejar datos privados.
%pip install --quiet -U "openai" "openai-agents>=0.17.0"
import os
import json
import re
import uuid
from datetime import datetime, timezone
from getpass import getpass
from importlib.metadata import PackageNotFoundError, version
try:
from openai import OpenAI
except Exception as e:
raise RuntimeError("Install dependency first: pip install -U openai") from e
MIN_AGENTS_SDK_VERSION = "0.17.0"
try:
from agents import (
Agent,
AgentOutputSchema,
FileSearchTool,
Runner,
RunConfig,
custom_span,
flush_traces,
function_span,
guardrail_span,
trace,
)
except Exception as e:
raise RuntimeError(
'Install or upgrade the OpenAI Agents SDK first: pip install -U "openai-agents>=0.17.0"'
) from e
def _version_tuple(value):
match = re.match(r"^(\d+)\.(\d+)\.(\d+)", str(value or ""))
return tuple(int(part) for part in match.groups()) if match else (0, 0, 0)
try:
AGENTS_SDK_VERSION = version("openai-agents")
except PackageNotFoundError as e:
raise RuntimeError('Install the OpenAI Agents SDK first: pip install -U "openai-agents>=0.17.0"') from e
if _version_tuple(AGENTS_SDK_VERSION) < _version_tuple(MIN_AGENTS_SDK_VERSION):
raise RuntimeError(
f'OpenAI Agents SDK {MIN_AGENTS_SDK_VERSION}+ is required; found {AGENTS_SDK_VERSION}. '
'Upgrade with: pip install -U "openai-agents>=0.17.0"'
)
def _clean_openai_api_key(value):
key = (value or "").strip()
if not key:
raise RuntimeError("OPENAI_API_KEY is required.")
return key
if not os.getenv("OPENAI_API_KEY", "").strip():
os.environ["OPENAI_API_KEY"] = getpass("Enter your OpenAI API key: ")
os.environ["OPENAI_API_KEY"] = _clean_openai_api_key(os.getenv("OPENAI_API_KEY"))
OPENAI_ORG_ID = os.getenv("OPENAI_ORG_ID", "").strip()
if OPENAI_ORG_ID:
os.environ["OPENAI_ORG_ID"] = OPENAI_ORG_ID
MODEL = os.getenv("OPENAI_MODEL", "gpt-5.5")
TRACE_INCLUDE_SENSITIVE_DATA = os.getenv("OPENAI_AGENTS_TRACE_INCLUDE_SENSITIVE_DATA", "false").lower() in {"1", "true", "yes", "on"}
os.environ["OPENAI_AGENTS_TRACE_INCLUDE_SENSITIVE_DATA"] = "true" if TRACE_INCLUDE_SENSITIVE_DATA else "false"
SCHEMAFLOW_TRACE_GROUP_ID = os.getenv("SCHEMAFLOW_TRACE_GROUP_ID") or (
"schemaflow-cookbook-" + datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%SZ") + "-" + uuid.uuid4().hex[:8]
)
os.environ["SCHEMAFLOW_TRACE_GROUP_ID"] = SCHEMAFLOW_TRACE_GROUP_ID
client = OpenAI(api_key=os.environ["OPENAI_API_KEY"])
print("Using model:", MODEL)
print("OpenAI Agents SDK:", AGENTS_SDK_VERSION)
print("OpenAI organization:", os.getenv("OPENAI_ORG_ID") or "(default for API key)")
print("Trace group:", SCHEMAFLOW_TRACE_GROUP_ID)
print("Trace payloads include prompts/outputs:", TRACE_INCLUDE_SENSITIVE_DATA)
from concurrent.futures import ThreadPoolExecutor
from pydantic import BaseModel, ConfigDict, Field
class SchemaFlowBaseModel(BaseModel):
model_config = ConfigDict(extra="allow")
class OperationModel(SchemaFlowBaseModel):
op: str
details: dict = Field(default_factory=dict)
class ChangeRequestModel(SchemaFlowBaseModel):
title: str | None = None
domain: str | None = None
target_schema: str | None = None
target_table: str | None = None
operations: list[OperationModel] = Field(default_factory=list)
notes: list = Field(default_factory=list)
class ImpactObjectModel(SchemaFlowBaseModel):
type: str
name: str
reason: str
source: str
class ImpactModel(SchemaFlowBaseModel):
impacted_objects: list[ImpactObjectModel] = Field(default_factory=list)
risks: list[str] = Field(default_factory=list)
assumptions: list[str] = Field(default_factory=list)
class PlanStepModel(SchemaFlowBaseModel):
id: str
description: str
class PlanModel(SchemaFlowBaseModel):
plan_steps: list[PlanStepModel] = Field(default_factory=list)
prechecks: list[str] = Field(default_factory=list)
postchecks: list[str] = Field(default_factory=list)
rollback: list[str] = Field(default_factory=list)
CHANGE_OUTPUT_SCHEMA = AgentOutputSchema(ChangeRequestModel, strict_json_schema=False)
IMPACT_OUTPUT_SCHEMA = AgentOutputSchema(ImpactModel, strict_json_schema=False)
PLAN_OUTPUT_SCHEMA = AgentOutputSchema(PlanModel, strict_json_schema=False)
def _parse_json_text(text: str):
text = (text or "{}").strip()
if text.startswith("```"):
text = re.sub(r"^```(?:json)?\s*", "", text)
text = re.sub(r"\s*```$", "", text).strip()
try:
return json.loads(text)
except json.JSONDecodeError:
match = re.search(r"\{.*\}", text, flags=re.DOTALL)
if not match:
raise
return json.loads(match.group(0))
def _model_dump(value):
if value is None or isinstance(value, (str, int, float, bool, bytes)):
return value
if isinstance(value, type):
return value
if hasattr(value, "model_dump"):
try:
return value.model_dump()
except TypeError:
pass
if hasattr(value, "to_dict"):
try:
return value.to_dict()
except TypeError:
pass
if hasattr(value, "__dict__"):
try:
return {k: v for k, v in vars(value).items() if not k.startswith("_")}
except TypeError:
pass
return value
def _agent_output_to_json(value):
value = _model_dump(value)
if isinstance(value, dict):
return value
if isinstance(value, str):
return _parse_json_text(value)
return json.loads(json.dumps(value, default=str))
def _agent_output_to_text(value):
value = _model_dump(value)
if isinstance(value, str):
return value.strip()
return json.dumps(value, ensure_ascii=False)
def _trace_metadata(metadata: dict | None = None):
cleaned = {}
for key, value in (metadata or {}).items():
if value is None:
cleaned[str(key)] = ""
elif isinstance(value, bool):
cleaned[str(key)] = "true" if value else "false"
elif isinstance(value, (dict, list, tuple, set)):
cleaned[str(key)] = json.dumps(value, ensure_ascii=False, default=str)
else:
cleaned[str(key)] = str(value)
return cleaned
def _schemaflow_run_config(workflow_name: str, metadata: dict | None = None):
return RunConfig(
workflow_name=workflow_name,
group_id=SCHEMAFLOW_TRACE_GROUP_ID,
trace_include_sensitive_data=TRACE_INCLUDE_SENSITIVE_DATA,
trace_metadata=_trace_metadata({"notebook": "schemaflow_cookbook", **(metadata or {})}),
)
def _runner_run_sync(agent, prompt: str, *, workflow_name: str, metadata: dict | None = None, max_turns: int = 4):
kwargs = {"run_config": _schemaflow_run_config(workflow_name, metadata), "max_turns": max_turns}
try:
return Runner.run_sync(agent, prompt, **kwargs)
except RuntimeError as exc:
if "event loop" not in str(exc).lower():
raise
with ThreadPoolExecutor(max_workers=1) as pool:
return pool.submit(lambda: Runner.run_sync(agent, prompt, **kwargs)).result()
def run_schemaflow_json_agent(*, name, instructions, prompt, output_schema, model=MODEL, tools=None, workflow_name=None, metadata=None):
agent = Agent(name=name, instructions=instructions, model=model, output_type=output_schema, tools=tools or [])
result = _runner_run_sync(agent, prompt, workflow_name=workflow_name or name, metadata={"agent": name, **(metadata or {})})
return _agent_output_to_json(result.final_output), result
def run_schemaflow_text_agent(*, name, instructions, prompt, model=MODEL, tools=None, workflow_name=None, metadata=None):
agent = Agent(name=name, instructions=instructions, model=model, tools=tools or [])
result = _runner_run_sync(agent, prompt, workflow_name=workflow_name or name, metadata={"agent": name, **(metadata or {})})
return _agent_output_to_text(result.final_output), result
def _collect_file_search_results(value):
results = []
seen = set()
def visit(node):
if node is None or isinstance(node, (str, int, float, bool, bytes)):
return
if isinstance(node, type) or callable(node):
return
node_id = id(node)
if node_id in seen:
return
seen.add(node_id)
node = _model_dump(node)
if node is None or isinstance(node, (str, int, float, bool, bytes)):
return
if isinstance(node, type) or callable(node):
return
if isinstance(node, dict):
if node.get("type") == "file_search_call":
for result in node.get("results", []) or []:
result = _model_dump(result)
if isinstance(result, dict):
text = result.get("text") or result.get("content") or ""
if isinstance(text, list):
text = "\n".join(str(x) for x in text)
results.append({"file_id": result.get("file_id"), "filename": result.get("filename") or result.get("file_name") or result.get("title"), "score": result.get("score"), "text_preview": str(text)[:1200]})
for child in node.values():
visit(child)
elif isinstance(node, (list, tuple, set)):
for child in node:
visit(child)
visit(value)
return results
def agent_file_search_results(run_result):
return _collect_file_search_results(run_result)
def trace_function_result(name: str, *, input_obj=None, output_obj=None):
with function_span(
name,
input=json.dumps(input_obj, ensure_ascii=False, default=str) if input_obj is not None else None,
output=json.dumps(output_obj, ensure_ascii=False, default=str) if output_obj is not None else None,
):
pass
def pretty(obj):
print(json.dumps(obj, indent=2, ensure_ascii=False))
2) Entrada
Esta sección define la solicitud de cambio de base de datos que procesará SchemaFlow. Piensa en ella como el ticket, problema o mensaje compacto que un equipo de datos podría recibir antes de convertir la solicitud en detalles de implementación.
La solicitud predeterminada pide al flujo de trabajo que:
- Agregue
LOYALTY_TIER VARCHAR(20)aODS.ODS_CUSTOMER_PROFILE. - Trate la nueva columna como anulable.
- Rellene desde
CORE.DIM_CUSTOMER. - Una en
CUSTOMER_ID. - Filtre la fuente a registros actuales con
IS_CURRENT=true. - Agregue un índice no único en
(CUSTOMER_ID, LOYALTY_TIER).
Esta entrada es intencionalmente compacta pero lo suficientemente rica como para ejercitar el flujo de trabajo completo:
- análisis del esquema y la tabla de destino
- extracción del nombre de la columna, tipo y anulabilidad
- reconocimiento de los requisitos de relleno
- reconocimiento de los requisitos de índice
- generación de SQL multicapa
- ejecución de comprobaciones de validación para la tabla, columna y acciones SQL esperadas
CHANGE_TEXT = """Add LOYALTY_TIER VARCHAR(20) to ODS.ODS_CUSTOMER_PROFILE as nullable.
Backfill from CORE.DIM_CUSTOMER on CUSTOMER_ID where IS_CURRENT=true.
Add a non-unique index on (CUSTOMER_ID, LOYALTY_TIER)."""
print(CHANGE_TEXT)
3) Contexto RAG de PDF opcional
SchemaFlow puede ejecutarse con o sin contexto de recuperación, por lo que los lectores pueden comenzar solo con la solicitud y agregar documentos de referencia solo cuando el cambio los necesite.
La ruta de PDF de ejemplo en la celda de código a continuación apunta a un archivo incluido en la carpeta del manual en data/, no a bytes incrustados dentro del notebook. Deja PDF_PATH = None para vistas previas de artículos estáticos o ejecuciones genéricas.
Con el PDF_PATH = None predeterminado, el notebook utiliza solo la solicitud de cambio en lenguaje natural. Esto es suficiente para demostrar el flujo de trabajo central por etapas.
Establece PDF_PATH en un PDF local cuando quieras que el Agente de Impacto base su análisis en material de referencia, como:
- documentos de diseño de interfaz
- especificaciones de esquema
- documentación de linaje
- contratos de datos
- notas de arquitectura de plataforma
- documentación de dependencias posteriores
Cuando se configura un PDF, esta sección:
- Valida que el archivo existe y es un PDF.
- Crea un almacén de vectores de OpenAI con una política de caducidad de un día.
- Carga el PDF al almacén de vectores.
- Permite que OpenAI maneje el análisis, la fragmentación, la incrustación y la recuperación.
- Almacena el ID del almacén de vectores para el Agente de Impacto.
- Más tarde resume cualquier resultado de búsqueda de archivos devuelto durante el análisis de impacto.
Esto mantiene el manual ligero porque no requiere modelos de incrustación locales, Chroma, Neo4j, LangGraph o módulos Python específicos del proyecto.
from pathlib import Path
# Optional PDF RAG example.
# The GitHub repo includes this sample PDF under schemaflow_cookbook/data.
# When running from a repo checkout in the cookbook folder, uncomment the path below to upload it to File Search.
# For meaningful retrieval hits, pair it with the LOYALTY_TIER change request used in this notebook.
PDF_PATH = None
# PDF_PATH = "data/sample_customer_loyalty_ifd.pdf"
RAG_MAX_RESULTS = 6
rag_vector_store = None
rag_vector_store_id = None
rag_vector_store_file = None
rag_file_search_results = []
impact_response = None
def create_pdf_vector_store(pdf_path):
pdf_path = Path(pdf_path).expanduser().resolve()
if not pdf_path.exists():
raise FileNotFoundError(f"PDF not found: {pdf_path}")
if pdf_path.suffix.lower() != ".pdf":
raise ValueError(f"Expected a PDF file, got: {pdf_path}")
with trace("SchemaFlow PDF Vector Store", group_id=SCHEMAFLOW_TRACE_GROUP_ID, metadata={"step": "pdf_vector_store", "pdf_path": str(pdf_path)}):
with custom_span("Create vector store", {"pdf_path": str(pdf_path)}):
vector_store = client.vector_stores.create(name=f"schemaflow-cookbook-{datetime.now(timezone.utc).strftime('%Y%m%dT%H%M%SZ')}", expires_after={"anchor": "last_active_at", "days": 1})
with custom_span("Upload PDF to vector store", {"vector_store_id": vector_store.id, "pdf_path": str(pdf_path)}):
with pdf_path.open("rb") as handle:
vector_store_file = client.vector_stores.files.upload_and_poll(vector_store_id=vector_store.id, file=handle)
trace_function_result("PDF vector store ready", input_obj={"pdf_path": str(pdf_path)}, output_obj={"vector_store_id": vector_store.id, "status": getattr(vector_store_file, "status", "unknown")})
flush_traces()
return vector_store, vector_store_file
if PDF_PATH:
rag_vector_store, rag_vector_store_file = create_pdf_vector_store(PDF_PATH)
rag_vector_store_id = rag_vector_store.id
print("Created vector store:", rag_vector_store_id)
print("Uploaded PDF status:", getattr(rag_vector_store_file, "status", "unknown"))
else:
print("No PDF configured. Leave PDF_PATH as None to run without RAG, or set it to a local PDF path.")
4) Etapas 1-2 - Análisis de la solicitud de cambio + Análisis de impacto
Esta sección ejecuta las dos primeras etapas del agente consecutivamente. Juntas, responden dos preguntas prácticas: ¿qué se solicitó exactamente y qué más podría verse afectado?
Etapa 1: Analizar la solicitud de cambio
El Agente de Análisis convierte CHANGE_TEXT en un objeto estructurado change_json.
Los campos esperados incluyen:
titledomaintarget_schematarget_tableoperationsnotes
Esta etapa crea el contrato normalizado que consume cada etapa posterior. Si el paso de análisis omite la tabla de destino, la columna, el tipo de datos, la anulabilidad, el relleno o la intención del índice, las etapas posteriores pueden producir una salida incompleta. Por eso el notebook valida esta etapa inmediatamente después.
Etapa 2: Análisis de impacto
El Agente de Impacto consume change_json y produce impact_json.
Los campos esperados incluyen:
impacted_objectsrisksassumptions
Si PDF_PATH se configuró anteriormente, el Agente de Impacto también recibe un FileSearchTool conectado al almacén de vectores PDF cargado. Esto permite que el modelo busque documentación de referencia antes de devolver las afirmaciones de impacto.
La salida es intencionalmente conservadora. Cuando el agente no está seguro, debe señalar suposiciones y riesgos en lugar de inventar certezas indocumentadas.
Vista previa del panel de impacto
La etapa de análisis de impacto produce impact_json estructurados que se pueden visualizar como un gráfico de objetos y relaciones afectados.
La vista previa a continuación muestra el tipo de gráfico de linaje de lealtad del cliente construido en la sección opcional del panel de Neo4j más adelante en el notebook. Ejecuta esa sección para generar la interfaz de usuario del gráfico local a partir de la semilla del gráfico de conocimiento de muestra e inspeccionar los objetos impactados de forma interactiva.

# =============================================================
# Stage 1 - Parse Change Request
# =============================================================
print("=" * 60)
print("Stage 1 - Parse Change Request")
print("=" * 60)
PARSE_SYSTEM = """
You are a precise information extraction system for database change requests.
Return STRICT JSON only (no prose, no code fences, no comments).
Required keys:
{
"title": str,
"domain": str|null,
"target_schema": str|null,
"target_table": str|null,
"operations": [{"op": str, "details": object}],
"notes": []
}
Rules:
- Use lowercase op names.
- If schema/table unknown, set null.
- Keep details explicit and typed where possible.
""".strip()
parse_user = "Change Request:\n\n" + CHANGE_TEXT
change_json, parse_agent_result = run_schemaflow_json_agent(name="SchemaFlow Parse Agent", instructions=PARSE_SYSTEM, prompt=parse_user, output_schema=CHANGE_OUTPUT_SCHEMA, workflow_name="SchemaFlow Stage 1 Parse", metadata={"stage": "parse_change_request"})
if isinstance(change_json, dict):
change_json.setdefault("title", None)
change_json.setdefault("domain", None)
change_json.setdefault("target_schema", None)
change_json.setdefault("target_table", None)
if not isinstance(change_json.get("operations"), list):
change_json["operations"] = [change_json.get("operations")] if change_json.get("operations") else []
if not isinstance(change_json.get("notes"), list):
change_json["notes"] = []
pretty(change_json)
# =============================================================
# Stage 2 - Impact Analysis
# =============================================================
print("\n" + "=" * 60)
print("Stage 2 - Impact Analysis")
print("=" * 60)
IMPACT_SYSTEM = """
You are a cautious impact analysis assistant.
Inputs:
- change_json: normalized change request.
- optional File Search context from an uploaded IFD/reference PDF.
Task:
Return JSON exactly as:
{
"impacted_objects": [
{"type":"table|column|fk|index|view","name":str,"reason":str,"source":"file_search|ifd|inference"}
],
"risks": [str],
"assumptions": [str]
}
Rules:
- Be conservative when uncertain.
- Call out data quality/backfill risks explicitly.
- If File Search context is available, use it to ground table, column, and downstream-impact claims.
""".strip()
impact_user_parts = ["CHANGE_JSON:\n" + json.dumps(change_json, ensure_ascii=False)]
impact_tools = []
if rag_vector_store_id:
impact_tools.append(FileSearchTool(vector_store_ids=[rag_vector_store_id], max_num_results=RAG_MAX_RESULTS, include_search_results=True))
impact_user_parts.append("Use the file_search tool against the uploaded PDF to look for relevant IFD, schema, table, column, lineage, and downstream dependency context before returning JSON.")
impact_json, impact_agent_result = run_schemaflow_json_agent(name="SchemaFlow Impact Agent", instructions=IMPACT_SYSTEM, prompt="\n\n".join(impact_user_parts), output_schema=IMPACT_OUTPUT_SCHEMA, tools=impact_tools, workflow_name="SchemaFlow Stage 2 Impact Analysis", metadata={"stage": "impact_analysis", "rag_enabled": bool(rag_vector_store_id)})
impact_response = impact_agent_result
try:
rag_file_search_results = agent_file_search_results(impact_agent_result)
except Exception as exc:
rag_file_search_results = []
print(f"File Search result summary skipped: {type(exc).__name__}: {exc}")
if rag_vector_store_id:
print("File Search results returned:", len(rag_file_search_results))
for i, result in enumerate(rag_file_search_results, start=1):
print(f"{i}. {result.get('filename') or result.get('file_id')} score={result.get('score')}")
if isinstance(impact_json, dict):
impact_json.setdefault("impacted_objects", [])
impact_json.setdefault("risks", [])
impact_json.setdefault("assumptions", [])
pretty(impact_json)
flush_traces()
Barreras de protección de salida de las etapas 1-2
Esta celda de barrera de protección realiza verificaciones deterministas en las salidas de Parse e Impact antes de que el flujo de trabajo continúe.
Las verificaciones comprueban que:
change_jsoncontiene un esquema de destino.change_jsoncontiene una tabla de destino.change_json.operationses una lista no vacía.impact_json.impacted_objectscontiene al menos un objeto.- La salida de impacto hace referencia a la tabla de destino analizada.
- Cada objeto impactado tiene campos básicos requeridos como tipo, nombre y razón.
Estas verificaciones son deliberadamente ligeras. No prueban que el análisis esté completo, pero detectan modos de falla obvios antes de que el Agente de Planificación o el Agente de SQL consuman un estado malformado o incompleto.
# Stages 1-2 Output Guardrails - inspects change_json (Parse) and impact_json (Impact).
stages_1_2_guardrails = []
with trace("SchemaFlow Stages 1-2 Guardrails", group_id=SCHEMAFLOW_TRACE_GROUP_ID, metadata={"stage": "stages_1_2_guardrails"}):
def _check(name, ok, detail=""):
ok = bool(ok)
stages_1_2_guardrails.append({"name": name, "ok": ok, "detail": detail})
with guardrail_span(name, triggered=not ok):
trace_function_result(name + " detail", output_obj={"ok": ok, "detail": detail})
_target_schema = (change_json.get("target_schema") or "").strip() if isinstance(change_json, dict) else ""
_target_table = (change_json.get("target_table") or "").strip() if isinstance(change_json, dict) else ""
_ops = change_json.get("operations") if isinstance(change_json, dict) else None
_check("parse_output_well_formed", bool(_target_schema) and bool(_target_table) and isinstance(_ops, list) and len(_ops) > 0, f"target={_target_schema}.{_target_table}, ops={len(_ops or [])}")
_impacted = impact_json.get("impacted_objects") if isinstance(impact_json, dict) else []
_target_fqn = f"{_target_schema}.{_target_table}" if (_target_schema and _target_table) else ""
_target_in_impact = any(isinstance(o, dict) and (o.get("name", "").upper() == _target_fqn.upper() or (_target_table and _target_table.upper() in o.get("name", "").upper())) for o in (_impacted or []))
_check("impact_includes_target", bool(_impacted) and _target_in_impact, f"{len(_impacted or [])} impacted object(s), target_match={_target_in_impact}")
_malformed = [i for i, o in enumerate(_impacted or []) if not (isinstance(o, dict) and o.get("type") and o.get("name") and o.get("reason"))]
_check("impacted_objects_well_formed", not _malformed, "all populated" if not _malformed else f"missing fields at indices {_malformed[:5]}")
stages_1_2_guardrails_passed = all(c["ok"] for c in stages_1_2_guardrails)
trace_function_result("Stages 1-2 guardrails summary", output_obj={"passed": stages_1_2_guardrails_passed, "checks": stages_1_2_guardrails})
flush_traces()
print(f"Stages 1-2 Output Guardrails: {'PASS' if stages_1_2_guardrails_passed else 'FAIL'}")
for _c in stages_1_2_guardrails:
_flag = "OK " if _c["ok"] else "FAIL"
print(f" [{_flag}] {_c['name']:35s} {_c['detail']}")
5) Etapas 3-4 - Plan de ejecución + Generación de SQL
Esta sección ejecuta las etapas de planificación de la implementación y generación de SQL. En este punto, el flujo de trabajo pasa de comprender la solicitud a redactar una entrega de implementación.
Etapa 3: Plan de ejecución
El Agente de Planificación consume:
change_jsonimpact_json
Devuelve plan_json con cuatro secciones:
plan_stepsprecheckspostchecksrollback
El objetivo es hacer explícita la estrategia de implementación antes de generar SQL. Esto ayuda a separar "lo que se debe hacer" de "qué SQL exacto se debe redactar".
Etapa 4: Generación de SQL
El Agente de SQL consume:
change_jsonplan_json
Devuelve un único script SQL de texto plano. El prompt requiere cuatro secciones en orden:
-- === LANDING (ODS) ===-- === STAGING (STG) ===-- === CORE (DIM/FACT/VIEW) ===-- === MARTS (SERVING) ===
El SQL generado está pensado como un borrador revisable. Debe ser revisado por ingenieros antes de cualquier uso en producción.
# =============================================================
# Stage 3 - Execution Plan
# =============================================================
print("=" * 60)
print("Stage 3 - Execution Plan")
print("=" * 60)
PLAN_SYSTEM = """
You are a senior data engineer creating a safe execution plan.
Inputs:
- change_json
- impact_json
Return JSON:
{
"plan_steps": [{"id": "str", "description": "str"}],
"prechecks": [str],
"postchecks": [str],
"rollback": [str]
}
Guidance:
- Include practical pre/post checks.
- Keep steps executable and concise.
""".strip()
plan_user = "\n\n".join(["CHANGE_JSON:\n" + json.dumps(change_json, ensure_ascii=False), "IMPACT_JSON:\n" + json.dumps(impact_json, ensure_ascii=False)])
plan_json, plan_agent_result = run_schemaflow_json_agent(name="SchemaFlow Plan Agent", instructions=PLAN_SYSTEM, prompt=plan_user, output_schema=PLAN_OUTPUT_SCHEMA, workflow_name="SchemaFlow Stage 3 Execution Plan", metadata={"stage": "execution_plan"})
if isinstance(plan_json, dict):
plan_json.setdefault("plan_steps", [])
plan_json.setdefault("prechecks", [])
plan_json.setdefault("postchecks", [])
plan_json.setdefault("rollback", [])
pretty(plan_json)
# =============================================================
# Stage 4 - SQL Generation
# =============================================================
print("\n" + "=" * 60)
print("Stage 4 - SQL Generation")
print("=" * 60)
SQL_SYSTEM = """
You are a senior data engineer producing SQL for multi-layer data stacks.
Output a SINGLE plaintext script with FOUR sections in order:
1) -- === LANDING (ODS) ===
2) -- === STAGING (STG) ===
3) -- === CORE (DIM/FACT/VIEW) ===
4) -- === MARTS (SERVING) ===
Rules:
- PostgreSQL dialect.
- Prefer idempotent DDL where possible.
- Propagate requested changes through downstream layers.
- Include concise assumptions as comments.
""".strip()
sql_user = "\n\n".join(["CHANGE_JSON:\n" + json.dumps(change_json, ensure_ascii=False), "PLAN_JSON:\n" + json.dumps(plan_json, ensure_ascii=False)])
sql_text, sql_agent_result = run_schemaflow_text_agent(name="SchemaFlow SQL Agent", instructions=SQL_SYSTEM, prompt=sql_user, workflow_name="SchemaFlow Stage 4 SQL Generation", metadata={"stage": "sql_generation"})
print(sql_text[:5000])
flush_traces()
Barreras de protección de salida de las etapas 3-4
Esta celda de barrera de protección valida el plan y el borrador de SQL antes de que el notebook pase a las verificaciones finales de cordura de SQL.
Las verificaciones comprueban que:
- las cuatro secciones del plan están pobladas:
plan_stepsprecheckspostchecksrollback
- el tipo de datos solicitado en
CHANGE_TEXTaparece en el SQL generado - las solicitudes anulables no crean accidentalmente restricciones
NOT NULL - las solicitudes explícitas
NOT NULLse reflejan cuando están presentes
Estas verificaciones complementan la Etapa 5. Las barreras de protección de las etapas 3-4 se centran en la completitud del plan y la consistencia semántica, mientras que la Etapa 5 se centra en los términos y acciones SQL esperados.
# Stages 3-4 Output Guardrails - inspects plan_json (Plan) and sql_text (SQL).
import re as _re
stages_3_4_guardrails = []
with trace("SchemaFlow Stages 3-4 Guardrails", group_id=SCHEMAFLOW_TRACE_GROUP_ID, metadata={"stage": "stages_3_4_guardrails"}):
def _check(name, ok, detail=""):
ok = bool(ok)
stages_3_4_guardrails.append({"name": name, "ok": ok, "detail": detail})
with guardrail_span(name, triggered=not ok):
trace_function_result(name + " detail", output_obj={"ok": ok, "detail": detail})
_plan = plan_json if isinstance(plan_json, dict) else {}
_plan_missing = [k for k in ["plan_steps", "prechecks", "postchecks", "rollback"] if not _plan.get(k)]
_check("plan_sections_populated", not _plan_missing, "all four populated" if not _plan_missing else f"empty: {_plan_missing}")
_dtype_match = _re.search(r"\b(?:add\s+\w+\s+|column\s+\w+\s+)((?:VAR)?CHAR\s*\([^)]*\)|TEXT|INTEGER|INT|BIGINT|BOOLEAN|DATE|TIMESTAMP|NUMERIC\s*\([^)]*\)|DECIMAL\s*\([^)]*\)|FLOAT|DOUBLE)", CHANGE_TEXT, flags=_re.IGNORECASE)
if _dtype_match:
_dtype = " ".join(_dtype_match.group(1).upper().split())
_check("data_type_propagated_to_sql", _dtype.lower() in sql_text.lower(), f"expected '{_dtype}' in SQL")
else:
_check("data_type_propagated_to_sql", True, "no data type referenced in CHANGE_TEXT (skipped)")
_change_lower = CHANGE_TEXT.lower()
_sql_lower = sql_text.lower()
_expected_cols = []
for _op in (change_json.get("operations") if isinstance(change_json, dict) else []) or []:
_details = _op.get("details") if isinstance(_op, dict) else None
if isinstance(_details, dict):
for _key in ("column", "column_name", "name"):
_val = _details.get(_key)
if isinstance(_val, str) and _val.strip():
_expected_cols.append(_val.strip().lower())
if "not null" in _change_lower:
_check("nullability_matches_request", "not null" in _sql_lower, "request: NOT NULL")
elif "nullable" in _change_lower:
_ddl_lines = []
for line in sql_text.split("\n"):
_line = line.strip().lower()
if not any(c in _line for c in _expected_cols):
continue
if "add column" in _line or any(_line.startswith(c + " ") or _line.startswith(c + "\t") for c in _expected_cols):
_ddl_lines.append(line.strip())
_bad_lines = [line for line in _ddl_lines if "not null" in line.lower()]
_check("nullability_matches_request", not _bad_lines, "no NOT NULL on nullable column DDL" if not _bad_lines else f"NOT NULL conflict in {len(_bad_lines)} DDL line(s)")
else:
_check("nullability_matches_request", True, "no explicit nullability requested (skipped)")
stages_3_4_guardrails_passed = all(c["ok"] for c in stages_3_4_guardrails)
trace_function_result("Stages 3-4 guardrails summary", output_obj={"passed": stages_3_4_guardrails_passed, "checks": stages_3_4_guardrails})
flush_traces()
print(f"Stages 3-4 Output Guardrails: {'PASS' if stages_3_4_guardrails_passed else 'FAIL'}")
for _c in stages_3_4_guardrails:
_flag = "OK " if _c["ok"] else "FAIL"
print(f" [{_flag}] {_c['name']:35s} {_c['detail']}")