# 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."""
snippet = _TAG_RE.sub("", " … ".join(fragments))[:400]
# La troncature peut couper une balise en deux (« …Plombier]*$", "", snippet)
if snippet.count("") > snippet.count(""):
snippet += ""
return snippet
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)$"),
images: bool = False,
):
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,
images_only=images,
)
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"),
"image": src.get("image_url"),
})
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()