llm-ready-data / app /api /v1 /keys_extract.py
validops-east-1's picture
common middlewre for /api/v1
b7dddbe
Raw
History Blame Contribute Delete
5.59 kB
from __future__ import annotations
import asyncio
import time
from typing import Any, Dict, List, Optional
from fastapi import APIRouter
from pydantic import BaseModel, Field, field_validator
from app.core.logger import get_logger
from app.core.thread_pool import thread_pool
from app.services.keys_extractor_service import KeysExtractor
logger = get_logger(__name__)
router = APIRouter()
class KeysExtractRequest(BaseModel):
data: Any = Field(..., description="JSON object or array to search")
key_names: Optional[List[str]] = Field(
default=None,
description=(
"Key names to look up at any depth. Mutually exclusive with `query`; "
"provide one of the two."
),
)
query: Optional[str] = Field(
default=None,
description=(
"Extended JSONPath expression (parsed with `jsonpath_ng.ext`) used to "
"select values. Supports filters, arithmetic and the union operator, e.g. "
"`$.store.book[?price > 10].title`. Each match is returned with its "
"full path inside the source document. Mutually exclusive with `key_names`."
),
)
result_limit: Optional[int] = Field(
None,
ge=1,
description=(
"Maximum results to return. For `key_names` mode this caps values per "
"key; for `query` mode it caps the number of matches returned. "
"Omit (or pass null) to return the full, uncapped result."
),
)
@field_validator("key_names")
@classmethod
def _validate_key_names(cls, v: Optional[List[str]]) -> Optional[List[str]]:
if v is None:
return v
if not v:
raise ValueError("key_names must be a non-empty list (omit it or supply a `query`)")
for kn in v:
if not isinstance(kn, str) or not kn:
raise ValueError("each key_name must be a non-empty string")
return v
@field_validator("query")
@classmethod
def _validate_query(cls, v: Optional[str]) -> Optional[str]:
if v is None:
return v
if not isinstance(v, str) or not v.strip():
raise ValueError("query must be a non-empty string")
return v
class KeysExtractResponse(BaseModel):
success: bool
time_ms: float
data: Dict[str, Any]
error_message: Optional[str] = None
class KeysExtractBatchRequest(BaseModel):
requests: List[KeysExtractRequest] = Field(
...,
min_length=1,
max_length=50,
description="List of extraction requests (1-50) to process concurrently",
)
class KeysExtractBatchResponse(BaseModel):
success: bool
time_ms: float
results: List[KeysExtractResponse]
def _error_response(start: float, message: str) -> KeysExtractResponse:
return KeysExtractResponse(
success=False,
time_ms=round((time.perf_counter() - start) * 1000, 3),
data={},
error_message=message,
)
async def _extract_single(body: KeysExtractRequest) -> KeysExtractResponse:
start = time.perf_counter()
has_key_names = bool(body.key_names)
has_query = body.query is not None
if has_key_names and has_query:
return _error_response(start, "provide either `key_names` or `query`, not both")
if not has_key_names and not has_query:
return _error_response(start, "either `key_names` or `query` must be provided")
if not isinstance(body.data, (dict, list)):
return _error_response(start, "`data` must be a JSON object or array")
try:
loop = asyncio.get_running_loop()
results = await loop.run_in_executor(
thread_pool,
KeysExtractor(body.data, body.key_names, body.result_limit, body.query).extract,
)
except (TypeError, ValueError) as exc:
return _error_response(start, str(exc))
except Exception as exc:
return _error_response(start, f"extraction failed: {exc}")
elapsed = round((time.perf_counter() - start) * 1000, 3)
return KeysExtractResponse(
success=True,
time_ms=elapsed,
data=results,
error_message=None,
)
@router.post(
"/keys/extract",
response_model=KeysExtractBatchResponse,
tags=["Keys Extractor"],
summary="Extract values from nested JSON objects (batch of up to 50 requests processed concurrently)",
)
async def extract_keys_batch(body: KeysExtractBatchRequest):
"""Process up to 50 extraction requests concurrently.
Each request supports two modes:
* **key_names** (default) -- recursively walk any JSON object/array and return
every value attached to the supplied key names, regardless of how deeply
nested they are. Returns ``{key_name: [values...]}``.
* **query** -- run an extended JSONPath expression and return each match with
its full path inside the source document. Returns
``{"matches": [{"path": "users.[0].role", "value": "admin"}, ...]}``.
"""
start = time.perf_counter()
logger.info("keys_extract_batch | start count=%d", len(body.requests))
tasks = [_extract_single(req) for req in body.requests]
results = await asyncio.gather(*tasks)
elapsed = round((time.perf_counter() - start) * 1000, 3)
all_ok = all(r.success for r in results)
logger.info(
"keys_extract_batch | done total=%d succeeded=%d time_ms=%s",
len(results),
sum(1 for r in results if r.success),
elapsed,
)
return KeysExtractBatchResponse(success=all_ok, time_ms=elapsed, results=results)