Files

788 lines
29 KiB
Python

"""Data workflow builder agent.
Creates reusable extraction/transformation workflow specs from uploaded CSV,
Excel, JSON, JSON-LD, or other structured data and a caller-provided JSON
Schema. The core workflow builder is deterministic so smoke/MCP calls do not
require an LLM credential. The generated workflow is saved as an artifact and,
when the managed database is available, recorded with an execution receipt for
platform acceptance evidence.
"""
from __future__ import annotations
import base64
import csv
import hashlib
import io
import json
import os
import re
import uuid
from datetime import datetime, timezone
from pathlib import Path
from typing import Annotated, Any, Literal
from pydantic import BaseModel, Field, ValidationError, field_validator
import a2a_pack as a2a
from a2a_pack import (
A2AAgent,
AgentDatabase,
AgentDatabaseEnv,
AgentDatabaseMigrations,
AgentPlatformResources,
FileUpload,
PlatformUserAuth,
Pricing,
Resources,
RunContext,
UploadedFile,
WorkspaceAccess,
WorkspaceMode,
)
try: # Optional at card-load time; required only when DATABASE_URL exists.
import psycopg
except Exception: # pragma: no cover - sandbox card load before optional deps
psycopg = None # type: ignore[assignment]
MAX_UPLOAD_BYTES = 2 * 1024 * 1024
MAX_SCHEMA_BYTES = 64 * 1024
MAX_SAMPLE_ROWS = 25
WORKFLOW_OUTPUT_PREFIX = "outputs/data-workflows/"
DB_NAME = "type-data-csv-excel-7rh4pz-057e66-data"
class TypeDataCsvExcel7rh4pz057e66Config(BaseModel):
max_upload_bytes: int = MAX_UPLOAD_BYTES
class BrowserDataUpload(BaseModel):
filename: str = Field(..., min_length=1, max_length=160)
media_type: str = Field(..., min_length=1, max_length=120)
data_base64: str = Field(..., min_length=1, max_length=MAX_UPLOAD_BYTES * 2)
@field_validator("filename")
@classmethod
def safe_filename(cls, value: str) -> str:
clean = value.strip().replace("\\", "/").split("/")[-1]
if not clean or clean in {".", ".."}:
raise ValueError("filename must be a safe basename")
if not re.fullmatch(r"[A-Za-z0-9._ ()@+-]{1,160}", clean):
raise ValueError("filename contains unsupported characters")
return clean
class WorkflowResult(BaseModel):
status: Literal["ok", "validation_error", "storage_warning"]
workflow_id: str
workflow_name: str
source_format: str
fields_detected: list[str]
mapped_fields: list[dict[str, Any]]
unmapped_required_fields: list[str]
workflow_path: str
artifact_uri: str | None = None
receipt_id: str
persisted: bool
warnings: list[str] = Field(default_factory=list)
class TypeDataCsvExcel7rh4pz057e66(A2AAgent[TypeDataCsvExcel7rh4pz057e66Config, PlatformUserAuth]):
name = "type-data-csv-excel-7rh4pz-057e66"
description = (
"Builds reusable extraction and transformation workflows from CSV, "
"Excel, JSON, JSON-LD, and other structured data into a requested JSON Schema."
)
version = "0.1.0"
config_model = TypeDataCsvExcel7rh4pz057e66Config
auth_model = PlatformUserAuth
pricing = Pricing(
price_per_call_usd=0.0,
caller_pays_llm=False,
notes="Deterministic workflow generation; no LLM credential is required.",
)
resources = Resources(cpu="500m", memory="512Mi", max_runtime_seconds=600)
workspace_access = WorkspaceAccess.dynamic(
max_files=32,
allowed_modes=(WorkspaceMode.READ_ONLY, WorkspaceMode.READ_WRITE_OVERLAY),
require_reason=False,
max_total_size_bytes=10 * 1024 * 1024,
)
platform_resources = AgentPlatformResources(
databases=(
AgentDatabase(
name=DB_NAME,
scope="user",
access_mode="read_write",
env=AgentDatabaseEnv(url="DATABASE_URL"),
migrations=AgentDatabaseMigrations(path="db/migrations"),
),
)
)
tools_used = ("jsonschema", "csv", "openpyxl")
@a2a.tool(
description=(
"Create and save a reusable extraction/transformation workflow from "
"a browser-uploaded base64 data file and a target JSON Schema."
),
timeout_seconds=120,
cost_class="standard",
grant_mode="read_write_overlay",
grant_allow_patterns=("outputs/data-workflows/**",),
grant_outputs_prefix=WORKFLOW_OUTPUT_PREFIX,
grant_write_prefixes=(WORKFLOW_OUTPUT_PREFIX,),
)
async def create_workflow_from_browser_upload(
self,
ctx: RunContext[PlatformUserAuth],
source_file: BrowserDataUpload,
output_schema: dict[str, Any],
workflow_name: str = "data-transform-workflow",
) -> WorkflowResult:
"""Browser-safe bounded base64 upload bridge for the packed frontend."""
try:
data = _decode_browser_upload(source_file, max_bytes=self.config.max_upload_bytes)
except ValueError as exc:
return WorkflowResult(
status="validation_error",
workflow_id="",
workflow_name=workflow_name,
source_format="unknown",
fields_detected=[],
mapped_fields=[],
unmapped_required_fields=[],
workflow_path="",
receipt_id="",
persisted=False,
warnings=[str(exc)],
)
return await self._create_workflow(
ctx=ctx,
filename=source_file.filename,
media_type=source_file.media_type,
data=data,
output_schema=output_schema,
workflow_name=workflow_name,
upload_mode="browser_base64",
)
@a2a.tool(
description=(
"Create and save a reusable extraction/transformation workflow from "
"an Agent API FileUpload and a target JSON Schema."
),
timeout_seconds=120,
cost_class="standard",
grant_mode="read_write_overlay",
grant_allow_patterns=("outputs/data-workflows/**",),
grant_outputs_prefix=WORKFLOW_OUTPUT_PREFIX,
grant_write_prefixes=(WORKFLOW_OUTPUT_PREFIX,),
)
async def create_workflow_from_file_upload(
self,
ctx: RunContext[PlatformUserAuth],
source_file: Annotated[
UploadedFile,
FileUpload(
accept=(
"text/csv",
"application/csv",
"application/json",
"application/ld+json",
"application/vnd.ms-excel",
"application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
"text/plain",
),
max_bytes=MAX_UPLOAD_BYTES,
description="CSV, Excel, JSON, JSON-LD, or text data file",
),
],
output_schema: dict[str, Any],
workflow_name: str = "data-transform-workflow",
) -> WorkflowResult:
"""Typed FileUpload path for external Agent API/MCP-compatible clients."""
if source_file.size_bytes and source_file.size_bytes > self.config.max_upload_bytes:
return WorkflowResult(
status="validation_error",
workflow_id="",
workflow_name=workflow_name,
source_format="unknown",
fields_detected=[],
mapped_fields=[],
unmapped_required_fields=[],
workflow_path="",
receipt_id="",
persisted=False,
warnings=[f"file exceeds {self.config.max_upload_bytes} byte limit"],
)
data = _read_uploaded_workspace_file(ctx, source_file.path)
if len(data) > self.config.max_upload_bytes:
return WorkflowResult(
status="validation_error",
workflow_id="",
workflow_name=workflow_name,
source_format="unknown",
fields_detected=[],
mapped_fields=[],
unmapped_required_fields=[],
workflow_path="",
receipt_id="",
persisted=False,
warnings=[f"file exceeds {self.config.max_upload_bytes} byte limit"],
)
return await self._create_workflow(
ctx=ctx,
filename=source_file.filename,
media_type=source_file.media_type,
data=data,
output_schema=output_schema,
workflow_name=workflow_name,
upload_mode="agent_file_upload",
)
@a2a.tool(
description="List recent saved workflow receipts for the signed-in user.",
timeout_seconds=30,
cost_class="cheap",
idempotent=True,
)
async def list_workflow_receipts(
self,
ctx: RunContext[PlatformUserAuth],
limit: int = 10,
) -> dict[str, Any]:
safe_limit = max(1, min(int(limit), 50))
rows = _list_receipts(_tenant_key(ctx), safe_limit)
return {"status": "ok", "receipts": rows, "persisted": bool(os.environ.get("DATABASE_URL"))}
async def _create_workflow(
self,
*,
ctx: RunContext[PlatformUserAuth],
filename: str,
media_type: str,
data: bytes,
output_schema: dict[str, Any],
workflow_name: str,
upload_mode: str,
) -> WorkflowResult:
tenant = _tenant_key(ctx)
warnings: list[str] = []
schema_error = _validate_json_schema(output_schema)
if schema_error:
return WorkflowResult(
status="validation_error",
workflow_id="",
workflow_name=workflow_name,
source_format="unknown",
fields_detected=[],
mapped_fields=[],
unmapped_required_fields=[],
workflow_path="",
receipt_id="",
persisted=False,
warnings=[schema_error],
)
if len(json.dumps(output_schema, default=str).encode("utf-8")) > MAX_SCHEMA_BYTES:
return WorkflowResult(
status="validation_error",
workflow_id="",
workflow_name=workflow_name,
source_format="unknown",
fields_detected=[],
mapped_fields=[],
unmapped_required_fields=[],
workflow_path="",
receipt_id="",
persisted=False,
warnings=[f"output_schema exceeds {MAX_SCHEMA_BYTES} bytes"],
)
profile = _profile_source_data(filename, media_type, data)
warnings.extend(profile.get("warnings", []))
workflow_id = _stable_workflow_id(filename, data, output_schema, workflow_name)
receipt_id = f"rcpt_{uuid.uuid4().hex}"
workflow_path = f"{WORKFLOW_OUTPUT_PREFIX}{workflow_id}.workflow.json"
workflow_spec = _build_workflow_spec(
workflow_id=workflow_id,
workflow_name=workflow_name,
filename=filename,
media_type=media_type,
upload_mode=upload_mode,
profile=profile,
output_schema=output_schema,
receipt_id=receipt_id,
)
artifact_uri: str | None = None
ref = await ctx.write_artifact(
f"{workflow_id}.workflow.json",
json.dumps(workflow_spec, indent=2, sort_keys=True).encode("utf-8"),
"application/json",
)
await ctx.emit_artifact(ref)
artifact_uri = ref.uri
await _write_workspace_output(ctx, workflow_path, workflow_spec, warnings)
persisted = _persist_workflow_and_receipt(
tenant_key=tenant,
workflow_id=workflow_id,
receipt_id=receipt_id,
workflow_name=workflow_name,
workflow_path=workflow_path,
source_format=str(profile["format"]),
source_filename=filename,
spec=workflow_spec,
result_summary={
"mapped_count": len(workflow_spec["mapping"]["fields"]),
"unmapped_required_fields": workflow_spec["mapping"]["unmapped_required_fields"],
"artifact_uri": artifact_uri,
},
warnings=warnings,
)
if not persisted:
warnings.append("DATABASE_URL is not configured; workflow artifact was created but database receipt was not persisted")
return WorkflowResult(
status="ok" if persisted else "storage_warning",
workflow_id=workflow_id,
workflow_name=workflow_name,
source_format=str(profile["format"]),
fields_detected=list(profile.get("fields", [])),
mapped_fields=list(workflow_spec["mapping"]["fields"]),
unmapped_required_fields=list(workflow_spec["mapping"]["unmapped_required_fields"]),
workflow_path=workflow_path,
artifact_uri=artifact_uri,
receipt_id=receipt_id,
persisted=persisted,
warnings=warnings,
)
def _tenant_key(ctx: RunContext[PlatformUserAuth]) -> str:
auth = ctx.auth
stable_id = auth.user_id if auth.user_id is not None else auth.sub
if stable_id is None or str(stable_id).strip() == "":
raise PermissionError("stable platform identity required")
return f"user:{stable_id}"
def _decode_browser_upload(source_file: BrowserDataUpload, *, max_bytes: int) -> bytes:
if source_file.media_type not in _accepted_media_types():
raise ValueError(f"unsupported media_type {source_file.media_type!r}")
try:
data = base64.b64decode(source_file.data_base64, validate=True)
except Exception as exc: # noqa: BLE001
raise ValueError("data_base64 must be valid base64") from exc
if not data:
raise ValueError("uploaded file is empty")
if len(data) > max_bytes:
raise ValueError(f"uploaded file exceeds {max_bytes} byte limit")
return data
def _accepted_media_types() -> set[str]:
return {
"text/csv",
"application/csv",
"application/json",
"application/ld+json",
"application/vnd.ms-excel",
"application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
"text/plain",
"application/octet-stream",
}
def _read_uploaded_workspace_file(ctx: RunContext[PlatformUserAuth], path: str) -> bytes:
workspace = ctx.workspace
reader = getattr(workspace, "read_bytes", None)
if callable(reader):
return bytes(reader(path))
raise RuntimeError("workspace file reader is unavailable for uploaded file")
async def _write_workspace_output(ctx: RunContext[PlatformUserAuth], path: str, payload: dict[str, Any], warnings: list[str]) -> None:
data = json.dumps(payload, indent=2, sort_keys=True).encode("utf-8")
try:
workspace = ctx.workspace
except PermissionError as exc:
warnings.append(f"workspace write skipped: {str(exc)[:160]}")
return
writer = getattr(workspace, "write_bytes", None)
if callable(writer):
writer(path, data)
return
try:
grant = await workspace.request_access(
files=(path,),
mode=WorkspaceMode.READ_WRITE_OVERLAY,
reason="Save generated data workflow spec",
purpose="workflow output",
)
view = await workspace.open_view(
purpose="workflow output",
hints=(path,),
max_files=1,
mode=WorkspaceMode.READ_WRITE_OVERLAY,
reason="Save generated data workflow spec",
)
await view.write(path, data)
_ = grant
except Exception as exc: # noqa: BLE001
warnings.append(f"workspace write warning: {type(exc).__name__}: {str(exc)[:160]}")
def _validate_json_schema(schema: dict[str, Any]) -> str | None:
if not isinstance(schema, dict):
return "output_schema must be a JSON object"
if schema.get("type") != "object":
return "output_schema must describe a JSON object with type='object'"
props = schema.get("properties")
if not isinstance(props, dict) or not props:
return "output_schema.properties must be a non-empty object"
required = schema.get("required", [])
if required is not None and not isinstance(required, list):
return "output_schema.required must be a list when provided"
try:
json.dumps(schema)
except Exception as exc: # noqa: BLE001
return f"output_schema must be JSON serializable: {exc}"
return None
def _profile_source_data(filename: str, media_type: str, data: bytes) -> dict[str, Any]:
ext = Path(filename.lower()).suffix
text = data.decode("utf-8-sig", errors="replace")
warnings: list[str] = []
if media_type in {"application/json", "application/ld+json"} or ext in {".json", ".jsonld"}:
try:
parsed = json.loads(text)
rows, fields = _json_rows_and_fields(parsed)
return {
"format": "jsonld" if media_type == "application/ld+json" or ext == ".jsonld" else "json",
"fields": fields,
"sample_rows": rows[:MAX_SAMPLE_ROWS],
"row_count_sampled": min(len(rows), MAX_SAMPLE_ROWS),
"warnings": warnings,
}
except json.JSONDecodeError as exc:
warnings.append(f"JSON parse warning: {exc.msg}")
if ext in {".xlsx", ".xlsm", ".xltx", ".xltm"} or "spreadsheetml" in media_type:
try:
return _profile_excel(data, warnings)
except Exception as exc: # noqa: BLE001
warnings.append(f"Excel parse warning: {type(exc).__name__}: {str(exc)[:120]}; falling back to delimited text")
delimiter = _sniff_delimiter(text)
try:
reader = csv.DictReader(io.StringIO(text), delimiter=delimiter)
fields = [str(f).strip() for f in (reader.fieldnames or []) if str(f).strip()]
rows = []
for idx, row in enumerate(reader):
if idx >= MAX_SAMPLE_ROWS:
break
rows.append({str(k): v for k, v in row.items() if k is not None})
if fields:
return {
"format": "csv" if delimiter == "," else "delimited_text",
"fields": fields,
"sample_rows": rows,
"row_count_sampled": len(rows),
"delimiter": delimiter,
"warnings": warnings,
}
except Exception as exc: # noqa: BLE001
warnings.append(f"CSV parse warning: {type(exc).__name__}: {str(exc)[:120]}")
lines = [line.strip() for line in text.splitlines() if line.strip()][:MAX_SAMPLE_ROWS]
return {
"format": "text",
"fields": [],
"sample_rows": [{"line": line} for line in lines],
"row_count_sampled": len(lines),
"warnings": warnings,
}
def _profile_excel(data: bytes, warnings: list[str]) -> dict[str, Any]:
from openpyxl import load_workbook
wb = load_workbook(io.BytesIO(data), read_only=True, data_only=True)
ws = wb[wb.sheetnames[0]]
rows_iter = ws.iter_rows(values_only=True)
header = next(rows_iter, None)
fields = [str(cell).strip() for cell in (header or []) if cell is not None and str(cell).strip()]
rows: list[dict[str, Any]] = []
for idx, row in enumerate(rows_iter):
if idx >= MAX_SAMPLE_ROWS:
break
rows.append({fields[i]: row[i] for i in range(min(len(fields), len(row)))})
if not fields:
warnings.append("first Excel sheet has no header row")
return {
"format": "excel",
"fields": fields,
"sample_rows": rows,
"row_count_sampled": len(rows),
"sheet": ws.title,
"warnings": warnings,
}
def _json_rows_and_fields(parsed: Any) -> tuple[list[dict[str, Any]], list[str]]:
if isinstance(parsed, list):
rows = [item for item in parsed if isinstance(item, dict)]
elif isinstance(parsed, dict):
if "@graph" in parsed and isinstance(parsed["@graph"], list):
rows = [item for item in parsed["@graph"] if isinstance(item, dict)]
else:
list_values = [value for value in parsed.values() if isinstance(value, list)]
rows = [item for item in (list_values[0] if list_values else [parsed]) if isinstance(item, dict)]
else:
rows = []
fields: list[str] = []
for row in rows[:MAX_SAMPLE_ROWS]:
for key in row:
if key not in fields:
fields.append(str(key))
return rows, fields
def _sniff_delimiter(text: str) -> str:
sample = text[:4096]
try:
return csv.Sniffer().sniff(sample, delimiters=",\t;|").delimiter
except Exception: # noqa: BLE001
return ","
def _build_workflow_spec(
*,
workflow_id: str,
workflow_name: str,
filename: str,
media_type: str,
upload_mode: str,
profile: dict[str, Any],
output_schema: dict[str, Any],
receipt_id: str,
) -> dict[str, Any]:
source_fields = list(profile.get("fields", []))
properties = output_schema.get("properties", {}) if isinstance(output_schema.get("properties"), dict) else {}
required = [str(item) for item in output_schema.get("required", [])]
mappings: list[dict[str, Any]] = []
unmapped_required: list[str] = []
for target, prop_schema in properties.items():
source = _best_field_match(str(target), source_fields)
target_type = prop_schema.get("type") if isinstance(prop_schema, dict) else None
rule = {
"target_field": str(target),
"source_field": source,
"required": str(target) in required,
"transformations": _transformations_for_type(target_type, prop_schema if isinstance(prop_schema, dict) else {}),
"confidence": 1.0 if source == target else 0.82 if source else 0.0,
}
mappings.append(rule)
if not source and str(target) in required:
unmapped_required.append(str(target))
return {
"workflow_version": "1.0",
"workflow_id": workflow_id,
"workflow_name": workflow_name.strip()[:120] or "data-transform-workflow",
"created_at": datetime.now(timezone.utc).isoformat(),
"source": {
"filename": filename,
"media_type": media_type,
"upload_mode": upload_mode,
"format": profile.get("format"),
"fields": source_fields,
"sample_rows": _json_safe(profile.get("sample_rows", [])),
"sample_limit": MAX_SAMPLE_ROWS,
},
"target_schema": output_schema,
"mapping": {
"strategy": "normalized_field_name_match_with_type_coercion",
"fields": mappings,
"unmapped_required_fields": unmapped_required,
},
"validation": {
"required_fields": required,
"additional_properties": output_schema.get("additionalProperties", True),
"gates": [
"parse source with declared source.format parser",
"apply per-field mapping and transformations",
"validate each output object against target_schema",
"send rows with unmapped required fields to review queue",
],
},
"receipt": {"receipt_id": receipt_id, "kind": "workflow_created"},
}
def _best_field_match(target: str, source_fields: list[str]) -> str | None:
if target in source_fields:
return target
normalized_target = _normalize_name(target)
normalized = {_normalize_name(field): field for field in source_fields}
if normalized_target in normalized:
return normalized[normalized_target]
for norm, original in normalized.items():
if normalized_target and (normalized_target in norm or norm in normalized_target):
return original
return None
def _normalize_name(value: str) -> str:
return re.sub(r"[^a-z0-9]+", "", value.lower())
def _transformations_for_type(target_type: Any, schema: dict[str, Any]) -> list[str]:
t = target_type[0] if isinstance(target_type, list) and target_type else target_type
steps = ["trim strings", "treat blank values as null"]
if t in {"integer", "number"}:
steps.append("parse numeric value")
elif t == "boolean":
steps.append("parse boolean synonyms true/false/yes/no/1/0")
elif t == "array":
steps.append("wrap scalar values or split delimited lists when configured")
elif t == "object":
steps.append("parse nested JSON object when source value is text")
else:
steps.append("coerce to string when non-null")
if "enum" in schema:
steps.append("validate against enum values")
if schema.get("format") == "date" or schema.get("format") == "date-time":
steps.append("normalize date/date-time to ISO 8601; flag ambiguous dates")
return steps
def _stable_workflow_id(filename: str, data: bytes, schema: dict[str, Any], workflow_name: str) -> str:
digest = hashlib.sha256()
digest.update(filename.encode("utf-8"))
digest.update(data[:65536])
digest.update(json.dumps(schema, sort_keys=True, default=str).encode("utf-8"))
digest.update(workflow_name.encode("utf-8"))
return f"wf_{digest.hexdigest()[:16]}"
def _json_safe(value: Any) -> Any:
try:
json.dumps(value)
return value
except TypeError:
return json.loads(json.dumps(value, default=str))
def _db_connect() -> Any | None:
url = os.environ.get("DATABASE_URL")
if not url or psycopg is None:
return None
return psycopg.connect(
url,
options="-c statement_timeout=5000 -c lock_timeout=2000 -c idle_in_transaction_session_timeout=10000",
)
def _persist_workflow_and_receipt(
*,
tenant_key: str,
workflow_id: str,
receipt_id: str,
workflow_name: str,
workflow_path: str,
source_format: str,
source_filename: str,
spec: dict[str, Any],
result_summary: dict[str, Any],
warnings: list[str],
) -> bool:
conn = _db_connect()
if conn is None:
return False
try:
with conn:
with conn.cursor() as cur:
cur.execute(
"""
INSERT INTO data_workflows
(tenant_key, workflow_id, workflow_name, workflow_path, source_format, source_filename, spec)
VALUES (%s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (tenant_key, workflow_id) DO UPDATE SET
workflow_name = EXCLUDED.workflow_name,
workflow_path = EXCLUDED.workflow_path,
source_format = EXCLUDED.source_format,
source_filename = EXCLUDED.source_filename,
spec = EXCLUDED.spec,
updated_at = NOW()
""",
(
tenant_key,
workflow_id,
workflow_name,
workflow_path,
source_format,
source_filename,
json.dumps(spec),
),
)
cur.execute(
"""
INSERT INTO execution_receipts
(tenant_key, receipt_id, workflow_id, skill_name, status, input_hash, result_summary, warnings)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s)
""",
(
tenant_key,
receipt_id,
workflow_id,
"create_workflow",
"ok",
hashlib.sha256(json.dumps(spec.get("source", {}), sort_keys=True, default=str).encode("utf-8")).hexdigest(),
json.dumps(result_summary),
json.dumps(warnings),
),
)
return True
except Exception: # noqa: BLE001
return False
finally:
conn.close()
def _list_receipts(tenant_key: str, limit: int) -> list[dict[str, Any]]:
conn = _db_connect()
if conn is None:
return []
try:
with conn:
with conn.cursor() as cur:
cur.execute(
"""
SELECT receipt_id, workflow_id, skill_name, status, result_summary, created_at
FROM execution_receipts
WHERE tenant_key = %s
ORDER BY created_at DESC
LIMIT %s
""",
(tenant_key, limit),
)
rows = []
for receipt_id, workflow_id, skill_name, status, result_summary, created_at in cur.fetchall():
rows.append(
{
"receipt_id": receipt_id,
"workflow_id": workflow_id,
"skill_name": skill_name,
"status": status,
"result_summary": result_summary,
"created_at": created_at.isoformat() if hasattr(created_at, "isoformat") else str(created_at),
}
)
return rows
finally:
conn.close()