dilutionrisk-mcp / document_query.py
mzx's picture
Exchange-aware API: accept NZX alongside ASX
952655a verified
Raw
History Blame Contribute Delete
3.92 kB
"""Stable, filtered document pagination over revision-pinned Parquet."""
from __future__ import annotations
import hashlib,json
from datetime import timezone
from pathlib import Path
from duckdb_runtime import query_rows
def check_exchange(value):
if value not in ("ASX","NZX"):raise ValueError("invalid exchange")
return value
def fingerprint(filters):return hashlib.sha256(json.dumps(filters,sort_keys=True,separators=(",", ":"),default=str).encode()).hexdigest()
def artifacts(cache_dir,table):
files=sorted(str(path) for path in (Path(cache_dir)/"data/v1"/table).rglob("*.parquet"))
if not files:raise RuntimeError("No %s Parquet artifacts"%table)
return files
def iso(value):
if value is None:return None
if hasattr(value,"astimezone"):return value.astimezone(timezone.utc).isoformat().replace("+00:00","Z")
return value.isoformat()
def search_documents(cache_dir,*,exchange="ASX",ticker=None,instrument_id=None,doc_type=None,date_from=None,date_to=None,direction="desc",limit=25,after=None):
if direction not in {"asc","desc"}:raise ValueError("invalid direction")
conditions=[f"d.exchange='{check_exchange(exchange)}'"];params=[artifacts(cache_dir,"documents")]
if ticker:conditions.append("d.ticker=?");params.append(ticker)
if doc_type:conditions.append("d.doc_type=?");params.append(doc_type)
if date_from:conditions.append("d.announcement_date>=?::DATE");params.append(str(date_from))
if date_to:conditions.append("d.announcement_date<=?::DATE");params.append(str(date_to))
if instrument_id:
conditions.append("EXISTS (SELECT 1 FROM read_parquet(?) e WHERE e.document_id=d.document_id AND e.instrument_id=?)");params.extend([artifacts(cache_dir,"events"),instrument_id])
sort_time="coalesce(d.announced_at,cast(d.announcement_date as TIMESTAMPTZ))";op=">" if direction=="asc" else "<"
if after:
conditions.append(f"({sort_time} {op} ?::TIMESTAMPTZ OR ({sort_time}=?::TIMESTAMPTZ AND d.document_id {op} ?))");params.extend([after[0],after[0],after[1]])
sql=f"""SELECT d.document_id,d.exchange,d.ticker,d.source_document_id,d.title,d.doc_type,d.doc_type_source,d.doc_type_confidence,d.announced_at,d.announcement_date,d.raw_artifact_key,d.raw_size_bytes,d.markdown_artifact_key,d.markdown_size_bytes,d.publication_state FROM read_parquet(?) d WHERE {' AND '.join(conditions)} ORDER BY {sort_time} {direction.upper()},d.document_id {direction.upper()} LIMIT ?"""
params.append(limit+1);rows=query_rows(sql,params);names=["document_id","exchange","ticker","source_document_id","title","doc_type","doc_type_source","doc_type_confidence","announced_at","announcement_date","raw_artifact_key","raw_size_bytes","markdown_artifact_key","markdown_size_bytes","publication_state"]
items=[]
for row in rows[:limit]:
item=dict(zip(names,row));item["announced_at"]=iso(item["announced_at"]);item["announcement_date"]=iso(item["announcement_date"]);items.append(item)
next_values=[items[-1]["announced_at"] or items[-1]["announcement_date"]+"T00:00:00Z",items[-1]["document_id"]] if len(rows)>limit else None
return items,next_values
def get_document(cache_dir,document_id):
paths=artifacts(cache_dir,"documents");rows=query_rows("SELECT document_id,exchange,ticker,source_document_id,title,doc_type,doc_type_source,doc_type_confidence,announced_at,announcement_date,raw_artifact_key,raw_size_bytes,markdown_artifact_key,markdown_size_bytes,publication_state FROM read_parquet(?) WHERE document_id=? LIMIT 1",[paths,document_id]);row=rows[0] if rows else None
if not row:return None
names=["document_id","exchange","ticker","source_document_id","title","doc_type","doc_type_source","doc_type_confidence","announced_at","announcement_date","raw_artifact_key","raw_size_bytes","markdown_artifact_key","markdown_size_bytes","publication_state"];item=dict(zip(names,row));item["announced_at"]=iso(item["announced_at"]);item["announcement_date"]=iso(item["announcement_date"]);return item