Spaces:
Runtime error
Runtime error
File size: 6,333 Bytes
04c4194 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 | """
DocuLens — Export endpoints.
Provides CSV and JSON export of extraction results
from pipeline runs stored in Supabase.
"""
import csv
import io
import json
import logging
from fastapi import APIRouter, HTTPException, Depends, Query
from fastapi.responses import StreamingResponse
from middleware import check_rate_limit
logger = logging.getLogger(__name__)
export_router = APIRouter()
# ---------------------------------------------------------------------------
# Helpers
# ---------------------------------------------------------------------------
def _get_run_data(run_id: str) -> dict:
"""Fetch extraction_data from a pipeline run."""
from db.supabase import get_client
client = get_client()
if not client:
raise HTTPException(status_code=503, detail="Database not configured")
resp = (
client.table("pipeline_runs")
.select("id, use_case, extraction_data, overall_result, processing_time_ms, started_at")
.eq("id", run_id)
.execute()
)
if not resp.data:
raise HTTPException(status_code=404, detail=f"Run '{run_id}' not found")
return resp.data[0]
def _flatten_fields(data: dict) -> dict:
"""Extract flat key-value fields from extraction_data."""
fields = {}
# extracted_fields (primary — works for all use cases)
if "extracted_fields" in data and isinstance(data["extracted_fields"], dict):
for k, v in data["extracted_fields"].items():
fields[k] = v
# Legacy named fields (invoice-specific)
skip_keys = {
"line_items", "bounding_boxes", "extracted_fields",
"extraction_method", "model_id", "page_image",
"page_images", "processing_time_ms", "confidence",
"document_type", "raw_text",
}
for k, v in data.items():
if k not in skip_keys and k not in fields and not isinstance(v, (dict, list)):
fields[k] = v
return fields
def _extraction_to_csv(data: dict) -> str:
"""Convert extraction data to CSV string."""
output = io.StringIO()
# Section 1: Header fields
fields = _flatten_fields(data)
if fields:
writer = csv.writer(output)
writer.writerow(["Field", "Value"])
for k, v in fields.items():
writer.writerow([k, v])
output.write("\n")
# Section 2: Line items
line_items = data.get("line_items", [])
if line_items:
# Collect all unique keys across line items
all_keys = []
seen = set()
for item in line_items:
if isinstance(item, dict):
for k in item.keys():
if k not in seen:
all_keys.append(k)
seen.add(k)
writer = csv.writer(output)
writer.writerow(all_keys)
for item in line_items:
if isinstance(item, dict):
writer.writerow([item.get(k, "") for k in all_keys])
return output.getvalue()
def _extraction_to_json(data: dict, run_meta: dict) -> dict:
"""Structure extraction data for JSON export."""
fields = _flatten_fields(data)
return {
"run_id": run_meta.get("id"),
"use_case": run_meta.get("use_case"),
"status": run_meta.get("overall_result"),
"processed_at": run_meta.get("started_at"),
"processing_time_ms": run_meta.get("processing_time_ms"),
"document_type": data.get("document_type", ""),
"confidence": data.get("confidence", 0),
"fields": fields,
"line_items": data.get("line_items", []),
}
# ---------------------------------------------------------------------------
# Endpoints
# ---------------------------------------------------------------------------
@export_router.get("/api/v1/runs/{run_id}/export")
async def export_run(
run_id: str,
format: str = Query("json", regex="^(json|csv)$"),
api_key: str = Depends(check_rate_limit),
):
"""
Export extraction results from a pipeline run.
Query params:
- format: "json" (default) or "csv"
"""
run = _get_run_data(run_id)
data = run.get("extraction_data") or {}
if not data:
raise HTTPException(status_code=404, detail="No extraction data for this run")
if format == "csv":
csv_content = _extraction_to_csv(data)
return StreamingResponse(
io.BytesIO(csv_content.encode("utf-8")),
media_type="text/csv",
headers={
"Content-Disposition": f'attachment; filename="run_{run_id[:8]}.csv"',
},
)
else:
export_data = _extraction_to_json(data, run)
json_content = json.dumps(export_data, indent=2, default=str)
return StreamingResponse(
io.BytesIO(json_content.encode("utf-8")),
media_type="application/json",
headers={
"Content-Disposition": f'attachment; filename="run_{run_id[:8]}.json"',
},
)
@export_router.post("/api/v1/export/inline")
async def export_inline(
data: dict,
format: str = Query("json", regex="^(json|csv)$"),
api_key: str = Depends(check_rate_limit),
):
"""
Export extraction data directly (without a saved run).
Accepts the extraction result JSON in the request body.
Useful for exporting results from in-progress extractions
that haven't been persisted yet.
"""
if format == "csv":
csv_content = _extraction_to_csv(data)
return StreamingResponse(
io.BytesIO(csv_content.encode("utf-8")),
media_type="text/csv",
headers={
"Content-Disposition": 'attachment; filename="extraction.csv"',
},
)
else:
fields = _flatten_fields(data)
export_data = {
"document_type": data.get("document_type", ""),
"confidence": data.get("confidence", 0),
"fields": fields,
"line_items": data.get("line_items", []),
}
json_content = json.dumps(export_data, indent=2, default=str)
return StreamingResponse(
io.BytesIO(json_content.encode("utf-8")),
media_type="application/json",
headers={
"Content-Disposition": 'attachment; filename="extraction.json"',
},
)
|