# Trouve-KA — API FastAPI # Author: Simon-Pierre Boucher # Contact: contact@spboucher.ai """API de Trouve-KA. Publique : /api/search, /api/status, /api/submit, /api/health Protégée (X-Admin-Token) : /api/admin/* — contrôle du crawler, jamais public sans auth. Dégradation gracieuse : si OpenSearch tombe, /api/status répond quand même; si Postgres tombe, la recherche répond quand même (analytics sautées). """ import html import re import time from contextlib import asynccontextmanager from typing import Annotated from fastapi import Depends, FastAPI, Header, HTTPException, Query from fastapi.middleware.cors import CORSMiddleware from pydantic import BaseModel from trouveka.config import get_settings from trouveka.database import Database from trouveka.logging import get_logger from trouveka.queue import Coordination from trouveka.ranking import build_search_body from trouveka.search_core import SearchCore from trouveka.shared import canonicalize_url, display_url, extract_domain, is_http_url log = get_logger("api") settings = get_settings() db = Database(settings.database_url, pool_min=settings.pg_pool_min, pool_max=settings.pg_pool_max) coord = Coordination(settings.redis_url) search = SearchCore(settings.search_url, settings.search_index) @asynccontextmanager async def lifespan(_app: FastAPI): await db.connect() await search.ensure_index() yield await search.close() await coord.close() await db.close() app = FastAPI(title="Trouve-KA API", version="0.1.0", lifespan=lifespan) # Derrière ngrok (www.trouve-ka.com), le web app proxifie /api : CORS permissif inutile # en prod, mais pratique en dev local (web sur :3000, API sur :8080). app.add_middleware( CORSMiddleware, allow_origins=["http://localhost:3000", settings.public_url], allow_methods=["GET", "POST"], allow_headers=["*", "X-Admin-Token"], ) def require_admin(x_admin_token: Annotated[str | None, Header()] = None) -> None: if not x_admin_token or x_admin_token != settings.admin_token: raise HTTPException(status_code=401, detail="Jeton admin invalide") _TAG_RE = re.compile(r"<(?!/?em>)[^>]*>") def _safe_snippet(fragments: list[str]) -> str: """Ne laisse passer que / (highlight); tout le reste est échappé par OpenSearch.""" return _TAG_RE.sub("", " … ".join(fragments))[:400] BADGE_LABELS = {"government": "Gouvernement", "news": "Actualités", "education": "Éducation"} # ---------------------------------------------------------------------- publique @app.get("/api/health") async def health(): return {"ok": True, "search_ok": await search.ping()} @app.get("/api/search") async def api_search( q: str = Query(..., min_length=1, max_length=200), page: int = Query(1, ge=1, le=100), limit: int = Query(10, ge=1, le=50), language: str | None = Query(None, pattern="^(fr|en)$"), location: str | None = None, category: str | None = Query(None, max_length=40), quebec_only: bool = False, freshness: str | None = Query(None, pattern="^(day|week|month|year)$"), ): started = time.monotonic() query_text = f"{q} {location}" if location else q body, analysis = build_search_body( query_text, page=page, limit=limit, language=language, category=category, quebec_only=quebec_only, freshness=freshness, ) try: res = await search.search(body) except Exception: log.exception("recherche échouée", extra={"ctx": {"q": q}}) raise HTTPException(status_code=503, detail="Le moteur de recherche est temporairement indisponible") took_ms = int((time.monotonic() - started) * 1000) results = [] for hit in res["hits"]["hits"]: src = hit["_source"] highlight = hit.get("highlight", {}) fragments = highlight.get("body") or highlight.get("description") or [] snippet = _safe_snippet(fragments) if fragments else html.escape(src.get("description") or "")[:400] badges = [BADGE_LABELS[c] for c in src.get("categories", []) if c in BADGE_LABELS] if src.get("page_quebec_score", 0) >= 0.45 or src.get("domain_quebec_score", 0) >= 0.6: badges.insert(0, "Québec") results.append({ "title": src.get("title") or src["url"], "url": src["url"], "display_url": display_url(src["url"]), "snippet": snippet, "domain": src["domain"], "language": src.get("language"), "quebec_score": src.get("page_quebec_score", 0), "badges": badges, "published_at": src.get("published_at"), }) total = res["hits"]["total"]["value"] # Analytics agrégées, respectueuses de la vie privée — jamais bloquantes try: await db.record_search_query(q, analysis["language"], total, took_ms) except Exception: log.exception("analytics de recherche sautées") return {"query": q, "total": total, "took_ms": took_ms, "page": page, "limit": limit, "results": results} @app.get("/api/status") async def api_status(): snapshot: dict = {} try: snapshot = await db.status_snapshot() except Exception: log.exception("statut PG indisponible") search_ok = await search.ping() try: paused = await coord.is_paused() except Exception: paused = False return { "pages_indexed": snapshot.get("pages_indexed", 0), "domains_count": snapshot.get("domains_count", 0), "indexed_last_hour": snapshot.get("indexed_last_hour", 0), "fetched_last_hour": snapshot.get("fetched_last_hour", 0), "errors_last_hour": snapshot.get("errors_last_hour", 0), "frontier_pending": snapshot.get("frontier_pending", 0), "frontier_in_progress": snapshot.get("frontier_in_progress", 0), "crawler_state": "paused" if paused else "running", "search_ok": search_ok, } @app.get("/api/live") async def api_live(): """Dernière page visitée par TrouveKABot — alimente le flux temps réel du footer. Public mais volontairement minimal : domaine + URL + horodatage, rien d'interne. """ try: events = await db.recent_events(1) except Exception: return {"event": None} if not events: return {"event": None} e = events[0] return { "event": { "at": e["at"], "url": e["url"], "domain": extract_domain(e["url"]), "outcome": e["outcome"], } } class SubmitBody(BaseModel): url: str @app.post("/api/submit") async def api_submit(payload: SubmitBody): """Soumettre un site québécois. Soumission ≠ inclusion : le crawler valide.""" if not is_http_url(payload.url): raise HTTPException(status_code=422, detail="URL invalide (http/https seulement)") url = canonicalize_url(payload.url) domain = extract_domain(url) if url else None if not url or not domain: raise HTTPException(status_code=422, detail="URL invalide") await db.add_submission(url) await db.enqueue_url(url, domain, priority=0.7, depth=0) return { "accepted": True, "message": "Merci! Le site sera visité par TrouveKABot. La soumission ne garantit pas l'inclusion.", } # ---------------------------------------------------------------------- admin class SeedsBody(BaseModel): urls: list[str] class RecrawlBody(BaseModel): url: str | None = None domain: str | None = None class DomainBody(BaseModel): domain: str @app.get("/api/admin/overview", dependencies=[Depends(require_admin)]) async def admin_overview(): overview = await db.admin_overview() overview["paused"] = await coord.is_paused() try: overview["enrich_backlog"] = await coord.enrich_backlog() except Exception: overview["enrich_backlog"] = None return overview @app.get("/api/admin/recent", dependencies=[Depends(require_admin)]) async def admin_recent(limit: int = Query(50, ge=1, le=200)): return {"events": await db.recent_events(limit)} @app.post("/api/admin/pause", dependencies=[Depends(require_admin)]) async def admin_pause(): await coord.pause_crawler() return {"paused": True} @app.post("/api/admin/resume", dependencies=[Depends(require_admin)]) async def admin_resume(): await coord.resume_crawler() return {"paused": False} @app.post("/api/admin/seeds", dependencies=[Depends(require_admin)]) async def admin_seeds(payload: SeedsBody): added = 0 for raw in payload.urls[:500]: url = canonicalize_url(raw.strip()) domain = extract_domain(url) if url else None if url and domain: if await db.enqueue_url(url, domain, priority=1.0, depth=0, is_seed=True): added += 1 return {"added": added} @app.post("/api/admin/recrawl", dependencies=[Depends(require_admin)]) async def admin_recrawl(payload: RecrawlBody): if payload.url: url = canonicalize_url(payload.url) ok = await db.requeue_url(url) if url else False return {"requeued": 1 if ok else 0} if payload.domain: return {"requeued": await db.requeue_domain(payload.domain.lower())} raise HTTPException(status_code=422, detail="url ou domain requis") @app.post("/api/admin/domains/block", dependencies=[Depends(require_admin)]) async def admin_block_domain(payload: DomainBody): await db.block_domain(payload.domain.lower()) return {"blocked": payload.domain.lower()} @app.get("/api/admin/frontier", dependencies=[Depends(require_admin)]) async def admin_frontier( domain: str | None = None, status: str | None = Query(None, pattern="^(pending|in_progress|done|failed|blocked)$"), limit: int = Query(100, ge=1, le=500), ): return {"items": await db.frontier_inspect(domain, status, limit)} def main() -> None: import uvicorn uvicorn.run("trouveka.api.main:app", host=settings.api_host, port=settings.api_port, workers=1) if __name__ == "__main__": main()