from datetime import datetime from fastapi import APIRouter, Depends, HTTPException, Query from sqlalchemy.orm import Session from ..database import get_db from ..deps import verify_admin 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"], dependencies=[Depends(verify_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") kind = (profile.kind or "heuristic").strip().lower() if kind == "llm": if not profile.llm_profile: raise HTTPException(status_code=400, detail="LLM profile has no llm_profile") elif 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") from ..services.job_stale import is_stale_job if job.status in ("queued", "running") and not is_stale_job(job): 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": from ..services.job_stale import is_stale_job if not is_stale_job(job): raise HTTPException(status_code=409, detail="Cannot delete a running job") # Stale running — allow delete after marking failed for audit trail job.status = "failed" job.last_error = "Deleted while stale running" db.commit() 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)