You cannot select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

212 lines
7.8 KiB
Python

"""Event Parser: parse exhibit text → events + event_parses."""
from __future__ import annotations
import argparse
import asyncio
import datetime as dt
import uuid
from sqlalchemy import select
from libs.common.config import get_settings
from libs.common.file_store import read_exhibit
from libs.common.ids import event_id as make_event_id
from libs.common.ids import new_job_run_id
from libs.common.logging import bind_job_run_id, configure_logging, get_logger
from libs.db.models import Document, Event, EventParse, JobRun
from libs.db.session import get_session
from libs.parser.rule_parser import PARSER_VERSION, SCHEMA_VERSION, RuleBasedParser
from libs.parser.schema_validator import validate_parser_output
from libs.parser.text_normalizer import normalize_text
logger = get_logger(__name__)
_parser = RuleBasedParser()
async def run_event_parser(run_id: str) -> dict[str, int]:
settings = get_settings()
app_config = settings.get_app_config()
exhibit_types = app_config.get("pipeline", {}).get("exhibit_types", ["EX-99.1"])
stats = {"seen": 0, "valid": 0, "invalid": 0, "errors": 0}
async with get_session() as session:
job = JobRun(
job_run_id=uuid.UUID(run_id),
job_name="event_parser",
source_name="oracle",
run_date=dt.date.today(),
status="running",
)
session.add(job)
await session.flush()
result = await session.execute(
select(Document).where(Document.parsed_status == "ready_for_parse")
)
docs = result.scalars().all()
stats["seen"] = len(docs)
for doc in docs:
if not doc.accession_no:
stats["invalid"] += 1
doc.parsed_status = "failed"
continue
# Try each exhibit type until we find one
text: str | None = None
for exhibit_type in exhibit_types:
try:
text = read_exhibit(doc.accession_no, exhibit_type)
break
except FileNotFoundError:
continue
if text is None:
logger.warning("no_exhibit_text", accession_no=doc.accession_no)
doc.parsed_status = "failed"
stats["errors"] += 1
continue
normalized = normalize_text(text)
metadata = {
"filing_date": doc.filing_date.isoformat(),
"accepted_at_utc": (
doc.accepted_at_utc.isoformat() if doc.accepted_at_utc else None
),
"form_type": doc.form_type,
}
try:
output = _parser.parse(
document_id=doc.document_id,
form_type=doc.form_type,
text=normalized,
metadata=metadata,
)
output_dict = output.model_dump()
errors = validate_parser_output(output_dict)
if errors:
logger.warning(
"parse_validation_failed",
document_id=doc.document_id,
errors=errors[:3],
)
# Store invalid parse, no event row
event_id_str = make_event_id(doc.document_id, output.event_type)
# Create a minimal event row first (needed for FK)
event = Event(
event_id=event_id_str,
primary_document_id=doc.document_id,
issuer_id=doc.issuer_id,
symbol_id=doc.symbol_id,
event_type=output.event_type,
event_direction=output.event_direction,
event_date=doc.filing_date,
filed_at_utc=doc.accepted_at_utc,
parser_version=PARSER_VERSION,
parse_confidence=output.confidence.overall,
status="rejected",
)
session.add(event)
await session.flush()
parse_row = EventParse(
event_id=event_id_str,
parser_kind="rule",
parser_version=PARSER_VERSION,
schema_version=SCHEMA_VERSION,
output_json=output_dict,
validation_status="invalid",
validation_errors={"errors": errors},
)
session.add(parse_row)
doc.parsed_status = "failed"
stats["invalid"] += 1
else:
event_id_str = make_event_id(doc.document_id, output.event_type)
# Check duplicate event
existing_evt = await session.execute(
select(Event).where(Event.event_id == event_id_str)
)
if existing_evt.scalar_one_or_none() is not None:
logger.info("event_already_exists", event_id=event_id_str)
doc.parsed_status = "succeeded"
stats["valid"] += 1
continue
event = Event(
event_id=event_id_str,
primary_document_id=doc.document_id,
issuer_id=doc.issuer_id,
symbol_id=doc.symbol_id,
event_type=output.event_type,
event_direction=output.event_direction,
event_date=doc.filing_date,
filed_at_utc=doc.accepted_at_utc,
parser_version=PARSER_VERSION,
parse_confidence=output.confidence.overall,
status="pending",
)
session.add(event)
await session.flush()
parse_row = EventParse(
event_id=event_id_str,
parser_kind="rule",
parser_version=PARSER_VERSION,
schema_version=SCHEMA_VERSION,
output_json=output_dict,
validation_status="valid",
validation_errors=None,
)
session.add(parse_row)
doc.parsed_status = "succeeded"
stats["valid"] += 1
logger.info(
"event_created",
event_id=event_id_str,
event_type=output.event_type,
confidence=output.confidence.overall,
)
except Exception as exc:
logger.error(
"parse_error",
document_id=doc.document_id,
error=str(exc),
)
doc.parsed_status = "failed"
stats["errors"] += 1
doc.updated_at_utc = dt.datetime.now(tz=dt.UTC)
job.status = "succeeded" if stats["errors"] == 0 else "partial"
job.finished_at_utc = dt.datetime.now(tz=dt.UTC)
job.records_seen = stats["seen"]
job.records_written = stats["valid"]
job.error_count = stats["errors"] + stats["invalid"]
logger.info("event_parser_done", **stats)
return stats
def main() -> None:
parser = argparse.ArgumentParser(description="Event Parser")
parser.add_argument("--run-id", default=new_job_run_id())
args = parser.parse_args()
settings = get_settings()
configure_logging(settings.log_level)
bind_job_run_id(args.run_id)
asyncio.run(run_event_parser(args.run_id))
if __name__ == "__main__":
main()