Files
gitrusprusandCursor 3f9dc6643b Add reusable LLM parser profiles with multi-event extract.
Support kind=llm profiles (instruction/schema), optional multi-event posts via #eN URLs, and recover stale running/queued parse jobs after worker crashes.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-13 20:22:38 +03:00

104 lines
2.8 KiB
Python

from datetime import datetime, timezone
from fastapi import APIRouter, Depends, HTTPException
from sqlalchemy.orm import Session
from ..database import get_db
from ..deps import verify_internal_token
from ..models import ParseJob
from ..schemas import (
IngestRequest,
IngestResponse,
ListenerSubscription,
VpnSettingsInternalRead,
)
from ..services.ingest import ingest_events
from ..services.jobs import resolve_job_source_config
from ..services.vpn import ensure_vpn_settings, to_internal_read
router = APIRouter(prefix="/internal", tags=["internal"])
@router.post("/ingest", response_model=IngestResponse)
def internal_ingest(
payload: IngestRequest,
_: None = Depends(verify_internal_token),
db: Session = Depends(get_db),
):
ingested, updated, skipped, map_synced = ingest_events(
db,
payload.events,
payload.job_id,
touch_job=not payload.listener,
)
return IngestResponse(
ingested=ingested,
updated=updated,
skipped=skipped,
map_objects_synced=map_synced,
)
@router.patch("/jobs/{job_id}")
def update_job_status(
job_id: int,
status: str,
error: str | None = None,
_: None = Depends(verify_internal_token),
db: Session = Depends(get_db),
):
job = db.query(ParseJob).filter(ParseJob.id == job_id).first()
if not job:
raise HTTPException(status_code=404, detail="Job not found")
job.status = status
# Anchor staleness detection: running/queued start, and terminal finish
if status in ("running", "queued", "completed", "failed"):
job.last_run_at = datetime.now(timezone.utc)
job.last_error = error
db.commit()
return {"id": job.id, "status": job.status}
@router.get("/listener/subscriptions", response_model=list[ListenerSubscription])
def listener_subscriptions(
_: None = Depends(verify_internal_token),
db: Session = Depends(get_db),
):
jobs = (
db.query(ParseJob)
.filter(
ParseJob.is_active.is_(True),
ParseJob.source_type == "telegram",
)
.order_by(ParseJob.id.asc())
.all()
)
result: list[ListenerSubscription] = []
for job in jobs:
try:
source_config = resolve_job_source_config(db, job)
except ValueError:
continue
channel = source_config.get("channel")
if not channel:
continue
result.append(
ListenerSubscription(
job_id=job.id,
channel=str(channel),
source_config=dict(source_config),
)
)
db.commit()
return result
@router.get("/vpn", response_model=VpnSettingsInternalRead)
def internal_vpn_settings(
_: None = Depends(verify_internal_token),
db: Session = Depends(get_db),
):
row = ensure_vpn_settings(db)
return to_internal_read(row)