Introduce reusable ParserProfile and ParseChannel, pair jobs with enqueue flatten, Events admin CRUD, and drop inline/legacy parser-builder UI and aliases. Co-authored-by: Cursor <cursoragent@cursor.com>
436 lines
15 KiB
Python
436 lines
15 KiB
Python
from datetime import datetime
|
|
|
|
from fastapi import APIRouter, Depends, HTTPException, Query
|
|
from sqlalchemy.orm import Session
|
|
|
|
from ..database import get_db
|
|
from ..models import Consumer, Event, MapObject, ParseChannel, ParseJob, ParserProfile
|
|
from ..schemas import (
|
|
AnalyticsSummary,
|
|
ConsumerCreate,
|
|
ConsumerRead,
|
|
ConsumerUpdate,
|
|
EventCreate,
|
|
EventListResponse,
|
|
EventRead,
|
|
EventUpdate,
|
|
ParseJobCreate,
|
|
ParseJobRead,
|
|
ParseJobUpdate,
|
|
TimelinePoint,
|
|
TopItem,
|
|
)
|
|
from ..services.analytics import (
|
|
get_analytics_summary,
|
|
get_timeline,
|
|
get_top_localities,
|
|
get_top_regions,
|
|
)
|
|
from ..services.events_query import build_events_query
|
|
from ..services.filtering import (
|
|
consumer_to_read,
|
|
create_consumer,
|
|
list_consumers,
|
|
rotate_consumer_key,
|
|
sync_event_to_map_object,
|
|
update_consumer,
|
|
)
|
|
from ..services.jobs import enqueue_parse_job, flatten_pair_config
|
|
|
|
router = APIRouter(prefix="/admin", tags=["admin"])
|
|
|
|
|
|
def _validated_source_config(source_type: str, source_config: dict) -> dict:
|
|
try:
|
|
from contracts.queues import family_for_source
|
|
from contracts.sources import parse_source_config
|
|
|
|
family_for_source(source_type)
|
|
return parse_source_config(source_type, source_config or {}).model_dump()
|
|
except Exception as exc:
|
|
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
|
|
|
|
|
def _job_to_read(job: ParseJob) -> ParseJobRead:
|
|
return ParseJobRead(
|
|
id=job.id,
|
|
source_type=job.source_type,
|
|
source_config=job.source_config or {},
|
|
profile_id=job.profile_id,
|
|
channel_id=job.channel_id,
|
|
schedule=job.schedule,
|
|
interval_seconds=job.interval_seconds,
|
|
is_active=job.is_active,
|
|
status=job.status,
|
|
last_run_at=job.last_run_at,
|
|
last_error=job.last_error,
|
|
created_at=job.created_at,
|
|
profile_name=job.profile.name if job.profile else None,
|
|
channel_name=job.channel.name if job.channel else None,
|
|
channel_handle=job.channel.channel if job.channel else None,
|
|
)
|
|
|
|
|
|
@router.post("/jobs", response_model=ParseJobRead, status_code=201)
|
|
def create_parse_job(payload: ParseJobCreate, db: Session = Depends(get_db)):
|
|
pair_mode = payload.profile_id is not None or payload.channel_id is not None
|
|
if pair_mode:
|
|
if payload.profile_id is None or payload.channel_id is None:
|
|
raise HTTPException(
|
|
status_code=400,
|
|
detail="Both profile_id and channel_id are required to create a pair",
|
|
)
|
|
profile = db.query(ParserProfile).filter(ParserProfile.id == payload.profile_id).first()
|
|
channel = db.query(ParseChannel).filter(ParseChannel.id == payload.channel_id).first()
|
|
if not profile:
|
|
raise HTTPException(status_code=404, detail="Profile not found")
|
|
if not channel:
|
|
raise HTTPException(status_code=404, detail="Channel not found")
|
|
if not profile.heuristic_profile:
|
|
raise HTTPException(status_code=400, detail="Profile has no heuristic_profile")
|
|
if not channel.is_active:
|
|
raise HTTPException(status_code=400, detail="Channel is inactive")
|
|
|
|
try:
|
|
source_config = flatten_pair_config(channel, profile, limit=payload.limit)
|
|
except Exception as exc:
|
|
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
|
|
|
job = ParseJob(
|
|
source_type=channel.source_type or "telegram",
|
|
source_config=source_config,
|
|
profile_id=profile.id,
|
|
channel_id=channel.id,
|
|
schedule=payload.schedule,
|
|
interval_seconds=payload.interval_seconds,
|
|
is_active=payload.is_active,
|
|
status="queued",
|
|
)
|
|
else:
|
|
source_config = _validated_source_config(payload.source_type, payload.source_config)
|
|
job = ParseJob(
|
|
source_type=payload.source_type,
|
|
source_config=source_config,
|
|
schedule=payload.schedule,
|
|
interval_seconds=payload.interval_seconds,
|
|
is_active=payload.is_active,
|
|
status="queued",
|
|
)
|
|
|
|
db.add(job)
|
|
db.commit()
|
|
db.refresh(job)
|
|
|
|
enqueue_parse_job(db, job)
|
|
db.refresh(job)
|
|
return _job_to_read(job)
|
|
|
|
|
|
@router.get("/jobs", response_model=list[ParseJobRead])
|
|
def list_parse_jobs(db: Session = Depends(get_db)):
|
|
jobs = db.query(ParseJob).order_by(ParseJob.id.desc()).all()
|
|
return [_job_to_read(job) for job in jobs]
|
|
|
|
|
|
@router.get("/jobs/{job_id}", response_model=ParseJobRead)
|
|
def get_parse_job(job_id: int, 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")
|
|
return _job_to_read(job)
|
|
|
|
|
|
@router.post("/jobs/{job_id}/retry", response_model=ParseJobRead)
|
|
def retry_parse_job(job_id: int, 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")
|
|
if job.status in ("queued", "running"):
|
|
raise HTTPException(status_code=409, detail="Job is already running or queued")
|
|
|
|
job.status = "queued"
|
|
job.last_error = None
|
|
db.commit()
|
|
db.refresh(job)
|
|
|
|
try:
|
|
enqueue_parse_job(db, job)
|
|
except ValueError as exc:
|
|
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
|
db.refresh(job)
|
|
return _job_to_read(job)
|
|
|
|
|
|
@router.patch("/jobs/{job_id}", response_model=ParseJobRead)
|
|
def update_parse_job(
|
|
job_id: int,
|
|
payload: ParseJobUpdate,
|
|
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")
|
|
|
|
if payload.profile_id is not None:
|
|
profile = db.query(ParserProfile).filter(ParserProfile.id == payload.profile_id).first()
|
|
if not profile:
|
|
raise HTTPException(status_code=404, detail="Profile not found")
|
|
job.profile_id = profile.id
|
|
if payload.channel_id is not None:
|
|
channel = db.query(ParseChannel).filter(ParseChannel.id == payload.channel_id).first()
|
|
if not channel:
|
|
raise HTTPException(status_code=404, detail="Channel not found")
|
|
job.channel_id = channel.id
|
|
|
|
if job.profile_id and job.channel_id:
|
|
profile = db.query(ParserProfile).filter(ParserProfile.id == job.profile_id).first()
|
|
channel = db.query(ParseChannel).filter(ParseChannel.id == job.channel_id).first()
|
|
if profile and channel:
|
|
existing_limit = (job.source_config or {}).get("limit", 100)
|
|
limit = payload.limit if payload.limit is not None else existing_limit
|
|
if not isinstance(limit, int):
|
|
limit = 100
|
|
try:
|
|
job.source_config = flatten_pair_config(channel, profile, limit=limit)
|
|
except Exception as exc:
|
|
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
|
elif payload.source_config is not None:
|
|
job.source_config = _validated_source_config(job.source_type, payload.source_config)
|
|
elif payload.limit is not None and isinstance(job.source_config, dict):
|
|
cfg = dict(job.source_config)
|
|
cfg["limit"] = payload.limit
|
|
job.source_config = _validated_source_config(job.source_type, cfg)
|
|
|
|
if payload.interval_seconds is not None:
|
|
job.interval_seconds = payload.interval_seconds
|
|
if payload.is_active is not None:
|
|
job.is_active = payload.is_active
|
|
|
|
db.commit()
|
|
db.refresh(job)
|
|
return _job_to_read(job)
|
|
|
|
|
|
@router.delete("/jobs/{job_id}", status_code=204)
|
|
def delete_parse_job(job_id: int, 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")
|
|
if job.status == "running":
|
|
raise HTTPException(status_code=409, detail="Cannot delete a running job")
|
|
|
|
db.delete(job)
|
|
db.commit()
|
|
|
|
|
|
@router.get("/events", response_model=EventListResponse)
|
|
def list_events(
|
|
limit: int = Query(default=50, ge=1, le=1000),
|
|
offset: int = Query(default=0, ge=0),
|
|
source_type: str | None = None,
|
|
region: str | None = None,
|
|
topic: str | None = None,
|
|
locality: str | None = None,
|
|
date_from: datetime | None = None,
|
|
date_to: datetime | None = None,
|
|
search: str | None = None,
|
|
db: Session = Depends(get_db),
|
|
):
|
|
query = build_events_query(
|
|
db,
|
|
source_type=source_type,
|
|
region=region,
|
|
topic=topic,
|
|
locality=locality,
|
|
date_from=date_from,
|
|
date_to=date_to,
|
|
search=search,
|
|
)
|
|
total = query.count()
|
|
items = query.offset(offset).limit(limit).all()
|
|
return EventListResponse(items=items, total=total)
|
|
|
|
|
|
def _ensure_unique_source_url(db: Session, source_url: str, *, exclude_id: int | None = None) -> None:
|
|
query = db.query(Event).filter(Event.source_url == source_url)
|
|
if exclude_id is not None:
|
|
query = query.filter(Event.id != exclude_id)
|
|
if query.first():
|
|
raise HTTPException(status_code=409, detail="Event with this source_url already exists")
|
|
|
|
|
|
def _sync_or_clear_map_object(db: Session, event: Event) -> None:
|
|
if event.latitude is not None and event.longitude is not None:
|
|
sync_event_to_map_object(db, event)
|
|
return
|
|
linked = db.query(MapObject).filter(MapObject.event_id == event.id).first()
|
|
if linked:
|
|
db.delete(linked)
|
|
|
|
|
|
@router.post("/events", response_model=EventRead, status_code=201)
|
|
def create_event(payload: EventCreate, db: Session = Depends(get_db)):
|
|
source_url = payload.source_url.strip()
|
|
if not source_url:
|
|
raise HTTPException(status_code=400, detail="source_url is required")
|
|
_ensure_unique_source_url(db, source_url)
|
|
|
|
if (payload.latitude is None) != (payload.longitude is None):
|
|
raise HTTPException(status_code=400, detail="latitude and longitude must be set together")
|
|
|
|
event = Event(
|
|
source_type=payload.source_type.strip() or "manual",
|
|
source_url=source_url,
|
|
raw_text=payload.raw_text or "",
|
|
title=payload.title or "",
|
|
description=payload.description or "",
|
|
locality=payload.locality or "",
|
|
latitude=payload.latitude,
|
|
longitude=payload.longitude,
|
|
event_date=payload.event_date,
|
|
region=payload.region,
|
|
topic=payload.topic,
|
|
tags=payload.tags,
|
|
metadata_=payload.metadata,
|
|
)
|
|
db.add(event)
|
|
db.flush()
|
|
_sync_or_clear_map_object(db, event)
|
|
db.commit()
|
|
db.refresh(event)
|
|
return event
|
|
|
|
|
|
@router.get("/events/{event_id}", response_model=EventRead)
|
|
def get_event(event_id: int, db: Session = Depends(get_db)):
|
|
event = db.query(Event).filter(Event.id == event_id).first()
|
|
if not event:
|
|
raise HTTPException(status_code=404, detail="Event not found")
|
|
return event
|
|
|
|
|
|
@router.patch("/events/{event_id}", response_model=EventRead)
|
|
def update_event(event_id: int, payload: EventUpdate, db: Session = Depends(get_db)):
|
|
event = db.query(Event).filter(Event.id == event_id).first()
|
|
if not event:
|
|
raise HTTPException(status_code=404, detail="Event not found")
|
|
|
|
data = payload.model_dump(exclude_unset=True)
|
|
if "source_url" in data and data["source_url"] is not None:
|
|
source_url = data["source_url"].strip()
|
|
if not source_url:
|
|
raise HTTPException(status_code=400, detail="source_url cannot be empty")
|
|
_ensure_unique_source_url(db, source_url, exclude_id=event.id)
|
|
event.source_url = source_url
|
|
data.pop("source_url")
|
|
|
|
if "metadata" in data:
|
|
event.metadata_ = data.pop("metadata")
|
|
|
|
if "source_type" in data and data["source_type"] is not None:
|
|
event.source_type = data.pop("source_type").strip() or event.source_type
|
|
|
|
for field, value in data.items():
|
|
setattr(event, field, value)
|
|
|
|
lat = event.latitude
|
|
lon = event.longitude
|
|
if (lat is None) != (lon is None):
|
|
raise HTTPException(status_code=400, detail="latitude and longitude must be set together")
|
|
|
|
_sync_or_clear_map_object(db, event)
|
|
db.commit()
|
|
db.refresh(event)
|
|
return event
|
|
|
|
|
|
@router.delete("/events/{event_id}", status_code=204)
|
|
def delete_event(event_id: int, db: Session = Depends(get_db)):
|
|
event = db.query(Event).filter(Event.id == event_id).first()
|
|
if not event:
|
|
raise HTTPException(status_code=404, detail="Event not found")
|
|
|
|
linked = db.query(MapObject).filter(MapObject.event_id == event.id).all()
|
|
for obj in linked:
|
|
db.delete(obj)
|
|
db.delete(event)
|
|
db.commit()
|
|
|
|
|
|
@router.get("/analytics/summary", response_model=AnalyticsSummary)
|
|
def analytics_summary(db: Session = Depends(get_db)):
|
|
return get_analytics_summary(db)
|
|
|
|
|
|
@router.get("/analytics/timeline", response_model=list[TimelinePoint])
|
|
def analytics_timeline(
|
|
days: int = Query(default=30, ge=1, le=365),
|
|
db: Session = Depends(get_db),
|
|
):
|
|
return get_timeline(db, days=days)
|
|
|
|
|
|
@router.get("/analytics/top-localities", response_model=list[TopItem])
|
|
def analytics_top_localities(
|
|
limit: int = Query(default=10, ge=1, le=100),
|
|
db: Session = Depends(get_db),
|
|
):
|
|
return get_top_localities(db, limit=limit)
|
|
|
|
|
|
@router.get("/analytics/top-regions", response_model=list[TopItem])
|
|
def analytics_top_regions(
|
|
limit: int = Query(default=10, ge=1, le=100),
|
|
db: Session = Depends(get_db),
|
|
):
|
|
return get_top_regions(db, limit=limit)
|
|
|
|
|
|
@router.get("/consumers", response_model=list[ConsumerRead])
|
|
def list_pi_consumers(db: Session = Depends(get_db)):
|
|
return [consumer_to_read(c) for c in list_consumers(db)]
|
|
|
|
|
|
@router.post("/consumers", response_model=ConsumerRead, status_code=201)
|
|
def create_pi_consumer(payload: ConsumerCreate, db: Session = Depends(get_db)):
|
|
consumer, api_key = create_consumer(
|
|
db,
|
|
name=payload.name,
|
|
regions=payload.regions,
|
|
topics=payload.topics,
|
|
date_from=payload.date_from,
|
|
)
|
|
return consumer_to_read(consumer, api_key=api_key)
|
|
|
|
|
|
@router.patch("/consumers/{consumer_id}", response_model=ConsumerRead)
|
|
def patch_pi_consumer(
|
|
consumer_id: int,
|
|
payload: ConsumerUpdate,
|
|
db: Session = Depends(get_db),
|
|
):
|
|
consumer = db.query(Consumer).filter(Consumer.id == consumer_id).first()
|
|
if not consumer:
|
|
raise HTTPException(status_code=404, detail="Consumer not found")
|
|
|
|
consumer = update_consumer(
|
|
db,
|
|
consumer,
|
|
name=payload.name,
|
|
is_active=payload.is_active,
|
|
regions=payload.regions,
|
|
topics=payload.topics,
|
|
date_from=payload.date_from,
|
|
)
|
|
return consumer_to_read(consumer)
|
|
|
|
|
|
@router.post("/consumers/{consumer_id}/rotate-key", response_model=ConsumerRead)
|
|
def rotate_pi_consumer_key(consumer_id: int, db: Session = Depends(get_db)):
|
|
consumer = db.query(Consumer).filter(Consumer.id == consumer_id).first()
|
|
if not consumer:
|
|
raise HTTPException(status_code=404, detail="Consumer not found")
|
|
|
|
consumer, api_key = rotate_consumer_key(db, consumer)
|
|
return consumer_to_read(consumer, api_key=api_key)
|