from __future__ import annotations from app.routes.reco_proxy import register_recommendation_routes import json import os from datetime import datetime from typing import Any, Optional import httpx from fastapi import APIRouter, File, HTTPException, UploadFile from pydantic import BaseModel, Field from app.config import settings from app.db import execute, fetch_all, fetch_one from app.services import herman as herman_service from app.services import ollama from app.services.auto_ingest import last_auto_sync, run_full_auto_sync from app.services.client_360 import get_client_360 from app.services import nas_folders from app.services.monitor import add_site, build_parse_intelligence, get_parse_page, list_parse_results, remove_site, trigger_crawl from app.services.agent_events_log import log_agent_event admin_router = APIRouter(prefix="/api/admin", tags=["admin-api"]) register_recommendation_routes(admin_router) ai_router = APIRouter(prefix="/api/ai", tags=["ai"]) herman_api = APIRouter(prefix="/api/herman", tags=["herman-api"]) voice_api = APIRouter(prefix="/api/voice", tags=["voice-api"]) def _serialize(row: dict | None) -> dict | None: if not row: return None out = dict(row) for k, v in list(out.items()): if hasattr(v, "isoformat"): out[k] = v.isoformat() elif isinstance(v, (datetime,)): out[k] = v.isoformat() return out def _serialize_rows(rows: list) -> list: return [_serialize(r) for r in rows] # --- Pydantic models --- class ClientBody(BaseModel): name: str contact: Optional[str] = None email: Optional[str] = None stage: str = "intake" sector: Optional[str] = None mrr_estimate: Optional[float] = None notes: Optional[str] = None class DealBody(BaseModel): client_id: Optional[int] = None title: str value: float = 0 stage: str = "lead" agent_owner: str = "herman" next_action: Optional[str] = None deadline: Optional[str] = None class ProductBody(BaseModel): client_id: Optional[int] = None name: str status: str = "concept" margin_pct: Optional[float] = None moq: Optional[int] = None shelf_target: Optional[str] = None launch_date: Optional[str] = None class SupplierBody(BaseModel): name: str country: Optional[str] = None category: Optional[str] = None moq: Optional[int] = None lead_time_days: Optional[int] = None rating: Optional[float] = None contact: Optional[str] = None class MentionBody(BaseModel): platform: str text: str class AccountBody(BaseModel): platform: str username: str is_active: bool = True class ScheduledPostBody(BaseModel): account_id: int content: str scheduled_time: str status: str = "pending" class PostStatusBody(BaseModel): status: str class AgentRuleBody(BaseModel): name: str condition_type: str threshold: float action: str = "notify" is_active: bool = True class SiteBody(BaseModel): url: str name: str class MarginBody(BaseModel): cost: float sell: float class AIGenerateBody(BaseModel): platform: str = "Instagram" topic: str tone: str = "vriendelijk en professioneel" class AIPostDraftBody(BaseModel): account_id: Optional[int] = None platform: str = "Instagram" topic: str tone: str = "vriendelijk en professioneel" class HermanChatBody(BaseModel): message: str channel: str = "cockpit" class ImageGenerateBody(BaseModel): prompt: str = Field(..., min_length=3, max_length=2000) width: int = Field(default=1024, ge=256, le=1024) height: int = Field(default=1024, ge=256, le=1024) steps: int = Field(default=28, ge=5, le=40) quality: str = Field(default="hd") class ImageStartBody(BaseModel): prompt: str = Field(..., min_length=3, max_length=2000) negative_prompt: str = Field(default="blurry, low quality, watermark, text, ugly, deformed", max_length=2000) quality: str = Field(default="hd") class CrawlBody(BaseModel): site_id: Optional[int] = None class DelegateBody(BaseModel): agent: str task: str # --- Clients CRUD --- @admin_router.get("/clients") def list_clients(): return {"items": _serialize_rows(fetch_all( """SELECT c.*, (SELECT COUNT(*) FROM client_supermarket_links l WHERE l.client_id = c.id) + (SELECT COUNT(*) FROM supermarkets s WHERE s.client_id = c.id AND NOT EXISTS ( SELECT 1 FROM client_supermarket_links l2 WHERE l2.supermarket_id = s.id AND l2.client_id = c.id )) AS store_count FROM clients c ORDER BY c.updated_at DESC NULLS LAST, c.created_at DESC LIMIT 500""" ))} @admin_router.post("/clients") def create_client(body: ClientBody): row = fetch_one( """INSERT INTO clients (name, contact, email, stage, sector, mrr_estimate, notes, updated_at) VALUES (%s,%s,%s,%s,%s,%s,%s,NOW()) RETURNING id""", (body.name, body.contact, body.email, body.stage, body.sector, body.mrr_estimate, body.notes), ) cid = row["id"] try: nas_folders.ensure_client_folder(cid, body.name) except Exception: pass return {"ok": True, "item": _serialize(fetch_one("SELECT * FROM clients WHERE id=%s", (cid,)))} @admin_router.put("/clients/{client_id}") def update_client(client_id: int, body: ClientBody): execute( """UPDATE clients SET name=%s, contact=%s, email=%s, stage=%s, sector=%s, mrr_estimate=%s, notes=%s, updated_at=NOW() WHERE id=%s""", (body.name, body.contact, body.email, body.stage, body.sector, body.mrr_estimate, body.notes, client_id), ) return {"ok": True, "item": _serialize(fetch_one("SELECT * FROM clients WHERE id=%s", (client_id,)))} @admin_router.delete("/clients/{client_id}") def delete_client(client_id: int): execute("DELETE FROM clients WHERE id=%s", (client_id,)) return {"ok": True} @admin_router.get("/clients/stats") def clients_stats(): stats: dict[str, Any] = {} try: rows = fetch_all("SELECT stage, COUNT(*) AS n FROM clients GROUP BY stage") stats["by_stage"] = {r["stage"]: int(r["n"]) for r in rows} stats["total"] = sum(stats["by_stage"].values()) stats["active"] = stats["by_stage"].get("active", 0) row = fetch_one("SELECT COALESCE(SUM(mrr_estimate),0) AS s FROM clients WHERE stage NOT IN ('churned')") stats["mrr_total"] = float(row["s"] or 0) if row else 0 row = fetch_one( """SELECT COALESCE(SUM(d.value),0) AS s FROM deals d JOIN clients c ON c.id = d.client_id WHERE d.stage NOT IN ('won','lost')""" ) stats["pipeline_value"] = float(row["s"] or 0) if row else 0 row = fetch_one("SELECT COUNT(*) AS n FROM client_supermarket_links") stats["store_links"] = int(row["n"] or 0) if row else 0 row = fetch_one("SELECT COUNT(*) AS n FROM deals") stats["deals"] = int(row["n"] or 0) if row else 0 except Exception: stats = {"total": 0, "active": 0, "mrr_total": 0, "pipeline_value": 0, "store_links": 0, "deals": 0, "by_stage": {}} return {"ok": True, "stats": stats} @admin_router.get("/clients/{client_id}/detail") def client_detail(client_id: int): client = fetch_one("SELECT * FROM clients WHERE id = %s", (client_id,)) if not client: raise HTTPException(404, "Client not found") deals = fetch_all( """SELECT id, title, value, stage, next_action, deadline, updated_at FROM deals WHERE client_id = %s ORDER BY updated_at DESC NULLS LAST LIMIT 20""", (client_id,), ) stores = fetch_all( """ SELECT s.id, s.name, s.chain, s.city, s.province, s.postcode, s.phone, s.email, s.partnership_status, s.halal_certified, s.has_halal_section, s.manager_name, l.relationship_type, l.notes AS link_notes, l.deal_id, l.created_at AS linked_at FROM client_supermarket_links l JOIN supermarkets s ON s.id = l.supermarket_id WHERE l.client_id = %s ORDER BY s.chain, s.city, s.name """, (client_id,), ) direct = fetch_all( """ SELECT s.id, s.name, s.chain, s.city, s.province, s.postcode, s.phone, s.email, s.partnership_status, s.halal_certified, s.has_halal_section, s.manager_name, 'direct' AS relationship_type, NULL AS link_notes, s.deal_id, s.last_updated AS linked_at FROM supermarkets s WHERE s.client_id = %s AND NOT EXISTS ( SELECT 1 FROM client_supermarket_links l WHERE l.supermarket_id = s.id AND l.client_id = s.client_id ) ORDER BY s.chain, s.city, s.name """, (client_id,), ) merged: dict[int, dict] = {} for row in list(stores) + list(direct): merged[int(row["id"])] = row return { "ok": True, "client": _serialize(client), "deals": _serialize_rows(deals), "stores": _serialize_rows(list(merged.values())), } @admin_router.get("/clients/{client_id}/360") def client_360_view(client_id: int) -> dict[str, Any]: data = get_client_360(client_id) if not data.get("ok"): raise HTTPException(404, data.get("error", "not found")) return data @admin_router.get("/documents/links") def document_links_list( client_id: int | None = None, project_id: int | None = None, ) -> dict[str, Any]: clauses, params = [], [] if client_id: clauses.append("dl.client_id = %s") params.append(client_id) if project_id: clauses.append("dl.project_id = %s") params.append(project_id) where = (" WHERE " + " AND ".join(clauses)) if clauses else "" rows = fetch_all( f""" SELECT dl.*, c.name AS client_name, p.name AS project_name FROM document_links dl LEFT JOIN clients c ON c.id = dl.client_id LEFT JOIN cockpit_projects p ON p.id = dl.project_id {where} ORDER BY dl.created_at DESC LIMIT 500 """, tuple(params), ) return {"ok": True, "items": _serialize_rows(rows)} @admin_router.post("/documents/links") def document_links_create(body: DocumentLinkBody) -> dict[str, Any]: if not body.client_id and not body.project_id: raise HTTPException(400, "client_id of project_id vereist") row = fetch_one( """ INSERT INTO document_links (storage_path, is_folder, client_id, project_id, notes) VALUES (%s, %s, %s, %s, %s) RETURNING * """, (body.storage_path.strip(), body.is_folder, body.client_id, body.project_id, body.notes or ""), ) if body.project_id and row: try: from app.services import projects as projects_svc projects_svc.add_asset( body.project_id, "nas_file", body.storage_path.split("/")[-1], file_path=body.storage_path.strip(), source_agent="documents", ) except Exception: pass return {"ok": True, "item": _serialize(row)} @admin_router.delete("/documents/links/{link_id}") def document_links_delete(link_id: int) -> dict[str, Any]: execute("DELETE FROM document_links WHERE id = %s", (link_id,)) return {"ok": True} @admin_router.post("/brain/auto-sync") async def brain_auto_sync(force: bool = False) -> dict[str, Any]: return await run_full_auto_sync(force_scan=force) @admin_router.get("/brain/auto-sync/status") def brain_auto_sync_status() -> dict[str, Any]: last = last_auto_sync() return {"ok": True, "last": last} # --- Deals CRUD --- @admin_router.get("/deals") def list_deals(): return {"items": _serialize_rows(fetch_all( """SELECT d.*, c.name AS client_name FROM deals d LEFT JOIN clients c ON c.id = d.client_id ORDER BY d.updated_at DESC NULLS LAST LIMIT 500""" ))} @admin_router.post("/deals") def create_deal(body: DealBody): row = fetch_one( """INSERT INTO deals (client_id, title, value, stage, agent_owner, next_action, deadline, updated_at) VALUES (%s,%s,%s,%s,%s,%s,%s,NOW()) RETURNING id""", (body.client_id, body.title, body.value, body.stage, body.agent_owner, body.next_action, body.deadline or None), ) did = row["id"] return {"ok": True, "item": _serialize(fetch_one("SELECT * FROM deals WHERE id=%s", (did,)))} @admin_router.put("/deals/{deal_id}") def update_deal(deal_id: int, body: DealBody): execute( """UPDATE deals SET client_id=%s, title=%s, value=%s, stage=%s, agent_owner=%s, next_action=%s, deadline=%s, updated_at=NOW() WHERE id=%s""", (body.client_id, body.title, body.value, body.stage, body.agent_owner, body.next_action, body.deadline or None, deal_id), ) return {"ok": True, "item": _serialize(fetch_one("SELECT * FROM deals WHERE id=%s", (deal_id,)))} @admin_router.delete("/deals/{deal_id}") def delete_deal(deal_id: int): execute("DELETE FROM deals WHERE id=%s", (deal_id,)) return {"ok": True} # --- Products CRUD --- @admin_router.get("/products/stats") def products_stats(): stats: dict[str, Any] = {} try: row = fetch_one("SELECT COUNT(*) AS n FROM products") stats["total"] = int(row["n"] or 0) if row else 0 rows = fetch_all("SELECT status, COUNT(*) AS n FROM products GROUP BY status") stats["by_status"] = {r["status"]: int(r["n"]) for r in rows} stats["active"] = stats["by_status"].get("active", 0) row = fetch_one("SELECT AVG(margin_pct) AS a FROM products WHERE margin_pct IS NOT NULL") stats["avg_margin"] = round(float(row["a"] or 0), 1) if row else 0 row = fetch_one("SELECT COUNT(*) AS n FROM products WHERE client_id IS NOT NULL") stats["with_client"] = int(row["n"] or 0) if row else 0 except Exception: stats = {"total": 0, "active": 0, "avg_margin": 0, "with_client": 0, "by_status": {}} return {"ok": True, "stats": stats} @admin_router.get("/products/{product_id}/detail") def product_detail(product_id: int): product = fetch_one( """SELECT p.*, c.name AS client_name, c.email AS client_email, c.stage AS client_stage FROM products p LEFT JOIN clients c ON c.id = p.client_id WHERE p.id = %s""", (product_id,), ) if not product: raise HTTPException(404, "Product not found") return {"ok": True, "product": _serialize(product)} @admin_router.get("/products") def list_products(): return {"items": _serialize_rows(fetch_all( """SELECT p.*, c.name AS client_name FROM products p LEFT JOIN clients c ON c.id = p.client_id ORDER BY p.created_at DESC LIMIT 500""" ))} @admin_router.post("/products") def create_product(body: ProductBody): row = fetch_one( """INSERT INTO products (client_id, name, status, margin_pct, moq, shelf_target, launch_date) VALUES (%s,%s,%s,%s,%s,%s,%s) RETURNING id""", (body.client_id, body.name, body.status, body.margin_pct, body.moq, body.shelf_target, body.launch_date or None), ) pid = row["id"] return {"ok": True, "item": _serialize(fetch_one("SELECT * FROM products WHERE id=%s", (pid,)))} @admin_router.put("/products/{product_id}") def update_product(product_id: int, body: ProductBody): execute( """UPDATE products SET client_id=%s, name=%s, status=%s, margin_pct=%s, moq=%s, shelf_target=%s, launch_date=%s WHERE id=%s""", (body.client_id, body.name, body.status, body.margin_pct, body.moq, body.shelf_target, body.launch_date or None, product_id), ) return {"ok": True, "item": _serialize(fetch_one("SELECT * FROM products WHERE id=%s", (product_id,)))} @admin_router.delete("/products/{product_id}") def delete_product(product_id: int): execute("DELETE FROM products WHERE id=%s", (product_id,)) return {"ok": True} @admin_router.post("/products/margin-calc") def margin_calc(body: MarginBody): if body.sell <= 0: raise HTTPException(400, "sell must be > 0") margin = ((body.sell - body.cost) / body.sell) * 100 return {"cost": body.cost, "sell": body.sell, "margin_pct": round(margin, 2), "profit": round(body.sell - body.cost, 2)} # --- Suppliers CRUD --- @admin_router.get("/suppliers/stats") def suppliers_stats(): stats: dict[str, Any] = {} try: row = fetch_one("SELECT COUNT(*) AS n FROM suppliers") stats["total"] = int(row["n"] or 0) if row else 0 row = fetch_one("SELECT AVG(rating) AS a FROM suppliers WHERE rating IS NOT NULL") stats["avg_rating"] = round(float(row["a"] or 0), 1) if row else 0 row = fetch_one("SELECT COUNT(DISTINCT country) AS n FROM suppliers WHERE country IS NOT NULL AND country <> ''") stats["countries"] = int(row["n"] or 0) if row else 0 row = fetch_one("SELECT AVG(lead_time_days) AS a FROM suppliers WHERE lead_time_days IS NOT NULL") stats["avg_lead"] = round(float(row["a"] or 0), 0) if row else 0 rows = fetch_all( "SELECT COALESCE(category, 'Overig') AS category, COUNT(*) AS n FROM suppliers GROUP BY category ORDER BY n DESC LIMIT 8" ) stats["by_category"] = {r["category"]: int(r["n"]) for r in rows} except Exception: stats = {"total": 0, "avg_rating": 0, "countries": 0, "avg_lead": 0, "by_category": {}} return {"ok": True, "stats": stats} @admin_router.get("/suppliers/{supplier_id}/detail") def supplier_detail(supplier_id: int): supplier = fetch_one("SELECT * FROM suppliers WHERE id = %s", (supplier_id,)) if not supplier: raise HTTPException(404, "Supplier not found") return {"ok": True, "supplier": _serialize(supplier)} @admin_router.get("/suppliers") def list_suppliers(): return {"items": _serialize_rows(fetch_all("SELECT * FROM suppliers ORDER BY name LIMIT 500"))} @admin_router.post("/suppliers") def create_supplier(body: SupplierBody): row = fetch_one( """INSERT INTO suppliers (name, country, category, moq, lead_time_days, rating, contact) VALUES (%s,%s,%s,%s,%s,%s,%s) RETURNING id""", (body.name, body.country, body.category, body.moq, body.lead_time_days, body.rating, body.contact), ) sid = row["id"] return {"ok": True, "item": _serialize(fetch_one("SELECT * FROM suppliers WHERE id=%s", (sid,)))} @admin_router.put("/suppliers/{supplier_id}") def update_supplier(supplier_id: int, body: SupplierBody): execute( """UPDATE suppliers SET name=%s, country=%s, category=%s, moq=%s, lead_time_days=%s, rating=%s, contact=%s WHERE id=%s""", (body.name, body.country, body.category, body.moq, body.lead_time_days, body.rating, body.contact, supplier_id), ) return {"ok": True, "item": _serialize(fetch_one("SELECT * FROM suppliers WHERE id=%s", (supplier_id,)))} @admin_router.delete("/suppliers/{supplier_id}") def delete_supplier(supplier_id: int): execute("DELETE FROM suppliers WHERE id=%s", (supplier_id,)) return {"ok": True} # --- Marketing --- @admin_router.get("/mentions") def list_mentions(platform: str = "all"): if platform == "all": rows = fetch_all("SELECT * FROM social_mentions ORDER BY created_at DESC LIMIT 200") else: rows = fetch_all( "SELECT * FROM social_mentions WHERE platform=%s ORDER BY created_at DESC LIMIT 200", (platform,), ) evaluate_agent_rules() return {"items": _serialize_rows(rows)} @admin_router.post("/mentions") def create_mention(body: MentionBody): text = body.text.strip() if not text: raise HTTPException(400, "text required") score = sentiment_score(text) row = fetch_one( "INSERT INTO social_mentions (platform, text, sentiment_score) VALUES (%s,%s,%s) RETURNING id", (body.platform, text, score), ) evaluate_agent_rules(row["id"]) return {"ok": True, "item": _serialize(fetch_one("SELECT * FROM social_mentions WHERE id=%s", (row["id"],)))} @admin_router.delete("/mentions/{mention_id}") def delete_mention(mention_id: int): execute("DELETE FROM social_mentions WHERE id=%s", (mention_id,)) return {"ok": True} @admin_router.get("/accounts") def list_accounts(): return {"items": _serialize_rows(fetch_all("SELECT * FROM social_accounts ORDER BY platform, username"))} @admin_router.post("/accounts") def create_account(body: AccountBody): row = fetch_one( "INSERT INTO social_accounts (platform, username, is_active) VALUES (%s,%s,%s) RETURNING id", (body.platform, body.username.strip(), body.is_active), ) return {"ok": True, "item": _serialize(fetch_one("SELECT * FROM social_accounts WHERE id=%s", (row["id"],)))} @admin_router.delete("/accounts/{account_id}") def delete_account(account_id: int): execute("DELETE FROM social_accounts WHERE id=%s", (account_id,)) return {"ok": True} @admin_router.get("/scheduled-posts") def list_scheduled_posts(): rows = fetch_all( """SELECT sp.*, sa.platform, sa.username FROM scheduled_posts sp JOIN social_accounts sa ON sp.account_id = sa.id ORDER BY sp.scheduled_time ASC""" ) return {"items": _serialize_rows(rows)} @admin_router.post("/scheduled-posts") def create_scheduled_post(body: ScheduledPostBody): row = fetch_one( "INSERT INTO scheduled_posts (account_id, content, scheduled_time, status) VALUES (%s,%s,%s,%s) RETURNING id", (body.account_id, body.content.strip(), body.scheduled_time, body.status), ) return {"ok": True, "item": _serialize(fetch_one("SELECT * FROM scheduled_posts WHERE id=%s", (row["id"],)))} @admin_router.patch("/scheduled-posts/{post_id}/status") def update_scheduled_status(post_id: int, body: PostStatusBody): allowed = {"pending", "approved", "posted", "rejected", "published", "cancelled"} if body.status not in allowed: raise HTTPException(400, "invalid status") if body.status in ("posted", "published"): execute( "UPDATE scheduled_posts SET status=%s, posted_at=NOW() WHERE id=%s", (body.status, post_id), ) else: execute("UPDATE scheduled_posts SET status=%s WHERE id=%s", (body.status, post_id)) return {"ok": True} @admin_router.put("/scheduled-posts/{post_id}") def update_scheduled_post(post_id: int, body: ScheduledPostBody): execute( "UPDATE scheduled_posts SET account_id=%s, content=%s, scheduled_time=%s, status=%s WHERE id=%s", (body.account_id, body.content, body.scheduled_time, body.status, post_id), ) return {"ok": True} @admin_router.delete("/scheduled-posts/{post_id}") def delete_scheduled_post(post_id: int): execute("DELETE FROM scheduled_posts WHERE id=%s", (post_id,)) return {"ok": True} @admin_router.post("/scheduled-posts/publish-due") def publish_due_scheduled_posts(): """Mark due scheduled posts as published (marketing cron).""" rows = fetch_all( """ SELECT id, content FROM scheduled_posts WHERE status = 'scheduled' AND scheduled_time <= NOW() ORDER BY scheduled_time ASC LIMIT 50 """ ) published = [] for row in rows: pid = int(row["id"]) execute( "UPDATE scheduled_posts SET status = 'published', posted_at = NOW() WHERE id = %s", (pid,), ) try: execute( """ INSERT INTO agent_events (agent_name, agent_type, event_type, title, body, status, channel, metadata) VALUES (%s, %s, %s, %s, %s, %s, %s, %s::jsonb) """, ( "marketing", "marketing", "post_published", f"Scheduled post #{pid} gepubliceerd", (row.get("content") or "")[:500], "completed", "cron", json.dumps({"post_id": pid}), ), ) except Exception: pass published.append(pid) return {"ok": True, "published": published, "count": len(published)} @admin_router.get("/agent-rules") def list_agent_rules(): evaluate_agent_rules() return {"items": _serialize_rows(fetch_all("SELECT * FROM agent_rules ORDER BY id"))} @admin_router.post("/agent-rules") def create_agent_rule(body: AgentRuleBody): row = fetch_one( """INSERT INTO agent_rules (name, condition_type, threshold, action, is_active) VALUES (%s,%s,%s,%s,%s) RETURNING id""", (body.name, body.condition_type, body.threshold, body.action, body.is_active), ) return {"ok": True, "item": _serialize(fetch_one("SELECT * FROM agent_rules WHERE id=%s", (row["id"],)))} @admin_router.patch("/agent-rules/{rule_id}/toggle") def toggle_agent_rule(rule_id: int): execute("UPDATE agent_rules SET is_active = NOT is_active WHERE id=%s", (rule_id,)) return {"ok": True} @admin_router.delete("/agent-rules/{rule_id}") def delete_agent_rule(rule_id: int): execute("DELETE FROM agent_rules WHERE id=%s", (rule_id,)) return {"ok": True} @admin_router.get("/agent-logs") def list_agent_logs(limit: int = 100): limit = max(1, min(limit, 500)) rows = fetch_all( """SELECT al.*, ar.name AS rule_name FROM agent_logs al LEFT JOIN agent_rules ar ON al.rule_id = ar.id ORDER BY al.created_at DESC LIMIT %s""", (limit,), ) return {"items": _serialize_rows(rows)} @admin_router.get("/analytics/social") def social_analytics(): account_stats = fetch_all( """SELECT sa.platform, sa.username, SUM(sa2.impressions) AS impressions, SUM(sa2.engagements) AS engagements, SUM(sa2.reach) AS reach FROM social_analytics sa2 JOIN social_accounts sa ON sa2.account_id = sa.id WHERE sa2.date > CURRENT_DATE - interval '7 days' GROUP BY sa.platform, sa.username ORDER BY impressions DESC""" ) daily = fetch_all( """SELECT date, SUM(impressions) AS impressions, SUM(engagements) AS engagements FROM social_analytics WHERE date > CURRENT_DATE - interval '7 days' GROUP BY date ORDER BY date""" ) mention_stats = fetch_one( "SELECT COUNT(*) AS cnt, COALESCE(AVG(sentiment_score),0) AS avg_sentiment FROM social_mentions" ) return { "account_stats": _serialize_rows(account_stats), "daily_stats": _serialize_rows(daily), "mentions": _serialize(mention_stats), } # --- Monitor --- @admin_router.get("/monitor/sites") def list_monitor_sites(): rows = fetch_all( """SELECT id, url, name, last_hash, last_crawled, is_active FROM monitored_sites ORDER BY last_crawled DESC NULLS LAST""" ) active = fetch_one("SELECT COUNT(*) AS c FROM monitored_sites WHERE is_active = TRUE") changes_24h = fetch_one( "SELECT COUNT(*) AS c FROM page_changes WHERE changed_at > NOW() - interval '24 hours'" ) return { "items": _serialize_rows(rows), "stats": { "active_sites": int(active["c"]) if active else 0, "changes_24h": int(changes_24h["c"]) if changes_24h else 0, }, } @admin_router.post("/monitor/sites") def create_monitor_site(body: SiteBody): try: item = add_site(body.url, body.name) except ValueError as exc: raise HTTPException(400, str(exc)) from exc for k, v in list(item.items()): if hasattr(v, "isoformat"): item[k] = v.isoformat() try: log_agent_event( "browser", "monitor_site_added", f"Site toegevoegd: {(body.name or body.url)[:80]}", body.url, channel="monitor", metadata={"url": body.url, "name": body.name, "site_id": item.get("id")}, ) except Exception: pass return {"ok": True, "item": item} @admin_router.delete("/monitor/sites/{site_id}") def delete_monitor_site(site_id: int, hard: bool = False): remove_site(site_id, soft=not hard) return {"ok": True} @admin_router.patch("/monitor/sites/{site_id}/toggle") def toggle_monitor_site(site_id: int): execute("UPDATE monitored_sites SET is_active = NOT is_active WHERE id=%s", (site_id,)) return {"ok": True} @admin_router.post("/monitor/trigger-crawl") def monitor_trigger_crawl(body: CrawlBody = CrawlBody()): result = trigger_crawl(body.site_id) summary_bits = [] for r in result.get("results") or []: if r.get("status") == "ERROR": summary_bits.append(f"{r.get('name') or r.get('url')}: fout") else: summary_bits.append( f"{r.get('title') or r.get('name')}: {r.get('word_count', 0)} woorden" ) try: log_agent_event( "browser", "monitor_crawl", f"Crawl klaar โ€” {result.get('sites', 0)} site(s), {result.get('changed', 0)} wijziging(en)", "; ".join(summary_bits)[:800], channel="monitor", metadata={"site_id": body.site_id, **{k: v for k, v in result.items() if k != "results"}}, ) except Exception: pass return result @admin_router.get("/monitor/parse-results") def monitor_parse_results(site_id: int | None = None, limit: int = 20): items = list_parse_results(site_id=site_id, limit=limit) return {"items": items, "count": len(items)} @admin_router.get("/monitor/intelligence") def monitor_intelligence(site_id: int | None = None, q: str | None = None, limit: int = 100): data = build_parse_intelligence(site_id=site_id, query=q, limit=min(limit, 100)) return data @admin_router.get("/monitor/parse-results/{page_id}") def monitor_parse_page_detail(page_id: int): page = get_parse_page(page_id) if not page: raise HTTPException(404, "Parse result not found") return {"page": page} @admin_router.get("/monitor/page-changes") def monitor_page_changes(limit: int = 50): rows = fetch_all( """SELECT pc.*, ms.url, ms.name FROM page_changes pc JOIN monitored_sites ms ON pc.site_id = ms.id ORDER BY pc.changed_at DESC LIMIT %s""", (max(1, min(limit, 200)),), ) return {"items": _serialize_rows(rows)} @admin_router.get("/monitor/crawl-logs") def monitor_crawl_logs(limit: int = 50): rows = fetch_all( """SELECT cl.*, ms.url, ms.name FROM crawl_logs cl LEFT JOIN monitored_sites ms ON cl.site_id = ms.id ORDER BY cl.logged_at DESC LIMIT %s""", (max(1, min(limit, 200)),), ) return {"items": _serialize_rows(rows)} # --- Agent events (approvals) --- @admin_router.patch("/agent-events/{event_id}/status") def patch_agent_event_status(event_id: int, body: PostStatusBody): if body.status not in ("approved", "rejected", "needs_approval", "completed"): raise HTTPException(400, "invalid status") execute("UPDATE agent_events SET status=%s WHERE id=%s", (body.status, event_id)) return {"ok": True} @admin_router.post("/delegate") async def delegate_task(body: DelegateBody): msg = f"[{body.agent}] {body.task}" result = await herman_service.chat(msg) return {"ok": True, "result": result} # --- AI --- async def _generate_social(prompt: str) -> str: try: return await ollama.generate(prompt) except Exception as exc: raise HTTPException(502, f"Ollama error: {exc}") from exc @ai_router.post("/generate-content") async def ai_generate_content(body: AIGenerateBody): topic = body.topic.strip() if not topic: raise HTTPException(400, "topic required") prompt = ( f"Schrijf een korte social media post voor {body.platform} over: {topic}. " f"Toon: {body.tone}. Maximaal 280 tekens. Alleen de posttekst, geen uitleg." ) content = await _generate_social(prompt) return {"ok": True, "content": content} @ai_router.post("/generate-image") async def ai_generate_image(body: ImageGenerateBody): import httpx try: async with httpx.AsyncClient(timeout=620.0) as client: r = await client.post( f"{settings.TOOLS_API_URL.rstrip('/')}/images/generate", json=body.model_dump(), ) r.raise_for_status() data = r.json() fn = data.get("filename", "") sub = data.get("subfolder", "") typ = data.get("type", "output") data["proxy_url"] = f"/api/ai/generated-image?filename={fn}&subfolder={sub}&type={typ}" return data except httpx.HTTPStatusError as exc: detail = exc.response.text[:300] if exc.response else str(exc) raise HTTPException(status_code=exc.response.status_code if exc.response else 502, detail=detail) from exc except Exception as exc: raise HTTPException(status_code=502, detail=str(exc)) from exc @ai_router.get("/generated-image") async def ai_generated_image( filename: str, subfolder: str = "", type: str = "output", ): import httpx from fastapi.responses import Response try: async with httpx.AsyncClient(timeout=60.0) as client: r = await client.get( f"{settings.TOOLS_API_URL.rstrip('/')}/images/view", params={"filename": filename, "subfolder": subfolder, "type": type}, ) r.raise_for_status() media = r.headers.get("content-type", "image/png") return Response(content=r.content, media_type=media) except Exception as exc: raise HTTPException(status_code=404, detail=str(exc)) from exc @ai_router.post("/generate-image/start") async def ai_generate_image_start(body: ImageStartBody): import httpx try: async with httpx.AsyncClient(timeout=30.0) as client: r = await client.post( f"{settings.TOOLS_API_URL.rstrip('/')}/images/generate/start", json=body.model_dump(), ) r.raise_for_status() return r.json() except httpx.HTTPStatusError as exc: detail = exc.response.text[:300] if exc.response else str(exc) raise HTTPException(status_code=exc.response.status_code if exc.response else 502, detail=detail) from exc except Exception as exc: raise HTTPException(status_code=502, detail=str(exc)) from exc @ai_router.get("/generate-image/progress/{prompt_id}") async def ai_generate_image_progress(prompt_id: str): import httpx try: async with httpx.AsyncClient(timeout=30.0) as client: r = await client.get(f"{settings.TOOLS_API_URL.rstrip('/')}/images/progress/{prompt_id}") r.raise_for_status() data = r.json() result = data.get("result") or {} if data.get("status") == "done" and result.get("filename"): fn = result["filename"] sub = result.get("subfolder", "") typ = result.get("type", "output") data["proxy_url"] = f"/api/ai/generated-image?filename={fn}&subfolder={sub}&type={typ}" return data except httpx.HTTPStatusError as exc: raise HTTPException(status_code=exc.response.status_code if exc.response else 502, detail=exc.response.text[:300]) from exc except Exception as exc: raise HTTPException(status_code=502, detail=str(exc)) from exc @ai_router.post("/generate-post-draft") async def ai_generate_post_draft(body: AIPostDraftBody): platform = body.platform if body.account_id: acc = fetch_one("SELECT platform, username FROM social_accounts WHERE id=%s", (body.account_id,)) if acc: platform = f"{acc['platform']} (@{acc['username']})" prompt = ( f"Schrijf een geplande social post voor {platform} over: {body.topic.strip()}. " f"Toon: {body.tone}. Geef titel + body." ) content = await _generate_social(prompt) return {"ok": True, "content": content, "platform": platform} @herman_api.post("/chat") async def herman_chat_api(body: HermanChatBody): if not body.message.strip(): raise HTTPException(400, "message required") result = await herman_service.chat(body.message.strip(), channel=(body.channel or "cockpit").strip()) return {"ok": True, **result} @voice_api.post("/transcribe") async def voice_transcribe(file: UploadFile = File(...)): url = "http://10.4.7.19:8877/v1/audio/transcriptions" data = await file.read() if not data: raise HTTPException(400, "empty file") try: async with httpx.AsyncClient(timeout=120.0) as client: resp = await client.post( url, files={"file": (file.filename or "audio.webm", data, file.content_type or "audio/webm")}, data={"model": "Systran/faster-whisper-base", "language": "nl"}, ) resp.raise_for_status() data = resp.json() return {"text": data.get("text", ""), "raw": data} except httpx.HTTPError as exc: raise HTTPException(502, f"transcribe proxy failed: {exc}") from exc # --- Document word analytics & sentiment --- @admin_router.get("/documents/summary") def documents_summary(): try: row = fetch_one( """ SELECT COUNT(*) AS documents, COALESCE(SUM(word_count), 0) AS total_words, COALESCE(AVG(sentiment_compound), 0) AS avg_sentiment FROM document_analytics """ ) unique = fetch_one( "SELECT COUNT(DISTINCT lemma) AS unique_words FROM document_word_counts WHERE NOT is_stopword" ) sentiment = fetch_all( """ SELECT sentiment_label, COUNT(*) AS cnt FROM document_analytics GROUP BY sentiment_label """ ) except Exception as exc: raise HTTPException(500, str(exc)) from exc return { "ok": True, "summary": { "documents": int(row["documents"] or 0) if row else 0, "total_words": int(row["total_words"] or 0) if row else 0, "avg_sentiment": round(float(row["avg_sentiment"] or 0), 4) if row else 0, "unique_words": int(unique["unique_words"] or 0) if unique else 0, "sentiment": {r["sentiment_label"]: int(r["cnt"]) for r in sentiment}, }, } @admin_router.get("/documents/words") def documents_words( q: str | None = None, limit: int = 40, stopwords: bool = False, ): limit = max(1, min(limit, 200)) clauses = [] params: list = [] if not stopwords: clauses.append("NOT is_stopword") if q and q.strip(): clauses.append("lemma ILIKE %s") params.append(f"%{q.strip()}%") where = f"WHERE {' AND '.join(clauses)}" if clauses else "" params.append(limit) try: rows = fetch_all( f""" SELECT lemma, MAX(token) AS token, SUM(count) AS total_count, COUNT(DISTINCT storage_path) AS document_count FROM document_word_counts {where} GROUP BY lemma ORDER BY total_count DESC LIMIT %s """, tuple(params), ) except Exception as exc: raise HTTPException(500, str(exc)) from exc return {"ok": True, "items": [_serialize(r) for r in rows]} @admin_router.get("/documents/list") def documents_list(limit: int = 50, label: str | None = None): limit = max(1, min(limit, 200)) clauses: list[str] = [] params: list[Any] = [] if label and label.strip(): clauses.append("%s = ANY(user_labels)") params.append(label.strip().lower()) where = f"WHERE {' AND '.join(clauses)}" if clauses else "" params.append(limit) try: rows = fetch_all( f""" SELECT filename, storage_path, doc_type, language, word_count, unique_lemmas, sentiment_label, sentiment_compound, sentiment_positive, sentiment_negative, sentiment_neutral, sentiment_subjectivity, extraction_method, analyzed_at, COALESCE(user_labels, '{{}}') AS user_labels, label_notes, labeled_at FROM document_analytics {where} ORDER BY analyzed_at DESC LIMIT %s """, tuple(params), ) except Exception as exc: raise HTTPException(500, str(exc)) from exc return {"ok": True, "items": [_serialize(r) for r in rows]} class DocumentLabelsBody(BaseModel): storage_path: str = Field(..., min_length=1) labels: list[str] = Field(default_factory=list) notes: str | None = None class PhotoLabelsBody(BaseModel): labels: list[str] = Field(default_factory=list) notes: str | None = None def _normalize_labels(raw: list[str]) -> list[str]: seen: set[str] = set() out: list[str] = [] for item in raw: tag = (item or "").strip().lower() if not tag or tag in seen: continue seen.add(tag) out.append(tag[:64]) return out @admin_router.patch("/documents/labels") def documents_patch_labels(body: DocumentLabelsBody) -> dict[str, Any]: labels = _normalize_labels(body.labels) notes = (body.notes or "").strip() or None try: row = fetch_one( """ UPDATE document_analytics SET user_labels = %s, label_notes = %s, labeled_at = NOW() WHERE storage_path = %s RETURNING filename, storage_path, user_labels, label_notes, labeled_at """, (labels, notes, body.storage_path), ) except Exception as exc: raise HTTPException(500, str(exc)) from exc if not row: execute( """ INSERT INTO document_analytics (storage_path, filename, user_labels, label_notes, labeled_at) VALUES (%s, %s, %s, %s, NOW()) ON CONFLICT (storage_path) DO UPDATE SET user_labels = EXCLUDED.user_labels, label_notes = EXCLUDED.label_notes, labeled_at = NOW() """, (body.storage_path, body.storage_path.split("/")[-1], labels, notes), ) row = fetch_one( """ SELECT filename, storage_path, user_labels, label_notes, labeled_at FROM document_analytics WHERE storage_path = %s """, (body.storage_path,), ) return {"ok": True, "document": _serialize(row)} @admin_router.patch("/photos/{photo_id}/labels") def photos_patch_labels(photo_id: int, body: PhotoLabelsBody) -> dict[str, Any]: labels = _normalize_labels(body.labels) notes = (body.notes or "").strip() or None try: row = fetch_one( """ UPDATE photo_imports SET user_labels = %s, label_notes = %s, labeled_at = NOW() WHERE id = %s RETURNING id, filename, storage_path, user_labels, label_notes, labeled_at """, (labels, notes, photo_id), ) except Exception as exc: raise HTTPException(500, str(exc)) from exc if not row: raise HTTPException(404, "Photo not found") return {"ok": True, "photo": _serialize(row)} @admin_router.get("/labels/vocabulary") def labels_vocabulary(limit: int = 80) -> dict[str, Any]: limit = max(1, min(limit, 200)) try: rows = fetch_all( """ SELECT label, COUNT(*) AS usage_count FROM ( SELECT unnest(COALESCE(user_labels, '{}')) AS label FROM document_analytics UNION ALL SELECT unnest(COALESCE(user_labels, '{}')) AS label FROM photo_imports ) t WHERE label IS NOT NULL AND label <> '' GROUP BY label ORDER BY usage_count DESC, label ASC LIMIT %s """, (limit,), ) except Exception as exc: raise HTTPException(500, str(exc)) from exc return {"ok": True, "labels": [_serialize(r) for r in rows]} PHOTO_LABEL_TAXONOMY = [ {"id": "productfoto", "label": "Productfoto", "emoji": "๐Ÿ“ฆ", "color": "#00e5ff"}, {"id": "prijslijst", "label": "Prijslijst", "emoji": "๐Ÿ’ฐ", "color": "#ffd700"}, {"id": "verpakking", "label": "Verpakking", "emoji": "๐Ÿ“‹", "color": "#b8ff3c"}, {"id": "halal", "label": "Halal certificaat", "emoji": "โœ…", "color": "#4ade80"}, {"id": "recept", "label": "Recept / technisch", "emoji": "๐Ÿงช", "color": "#a78bfa"}, {"id": "factuur", "label": "Factuur / bon", "emoji": "๐Ÿงพ", "color": "#fb7185"}, {"id": "etiket", "label": "Etiket", "emoji": "๐Ÿท๏ธ", "color": "#f97316"}, {"id": "marketing", "label": "Marketing", "emoji": "๐Ÿ“ฃ", "color": "#38bdf8"}, {"id": "schap", "label": "Schap / retail", "emoji": "๐Ÿ›’", "color": "#22d3ee"}, {"id": "menu", "label": "Menu / horeca", "emoji": "๐Ÿฝ๏ธ", "color": "#fbbf24"}, {"id": "overig", "label": "Overig", "emoji": "๐Ÿ“", "color": "#94a3b8"}, ] DOCUMENT_LABEL_TAXONOMY = [ {"id": "contract", "label": "Contract", "emoji": "๐Ÿ“„", "color": "#00e5ff"}, {"id": "presentatie", "label": "Presentatie", "emoji": "๐Ÿ“Š", "color": "#a78bfa"}, {"id": "rapport", "label": "Rapport", "emoji": "๐Ÿ“ˆ", "color": "#4ade80"}, {"id": "halal", "label": "Halal / certificaat", "emoji": "โœ…", "color": "#22c55e"}, {"id": "prijslijst", "label": "Prijslijst", "emoji": "๐Ÿ’ฐ", "color": "#ffd700"}, {"id": "marketing", "label": "Marketing", "emoji": "๐Ÿ“ฃ", "color": "#38bdf8"}, {"id": "juridisch", "label": "Juridisch", "emoji": "โš–๏ธ", "color": "#fb7185"}, {"id": "financieel", "label": "Financieel", "emoji": "๐Ÿงพ", "color": "#f97316"}, {"id": "product", "label": "Product info", "emoji": "๐Ÿ“ฆ", "color": "#b8ff3c"}, {"id": "overig", "label": "Overig", "emoji": "๐Ÿ“", "color": "#94a3b8"}, ] @admin_router.get("/labels/taxonomy") def labels_taxonomy() -> dict[str, Any]: return { "ok": True, "photo": PHOTO_LABEL_TAXONOMY, "document": DOCUMENT_LABEL_TAXONOMY, } class DoclingConvertBody(BaseModel): path: str = Field(..., min_length=1) ocr: bool = True tables: bool = True page_images: bool = False picture_images: bool = True force_full_page_ocr: bool = False formats: list[str] = Field(default_factory=lambda: ["markdown", "html", "json", "text", "doctags", "yaml"]) class DoclingBatchBody(BaseModel): paths: list[str] = Field(default_factory=list) limit: int = Field(default=10, ge=1, le=50) ext: str | None = None ocr: bool = True tables: bool = True page_images: bool = False picture_images: bool = True force_full_page_ocr: bool = False formats: list[str] = Field(default_factory=lambda: ["markdown", "json"]) @admin_router.get("/docling/capabilities") async def docling_capabilities_proxy() -> dict[str, Any]: doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") try: async with httpx.AsyncClient(timeout=30.0) as client: r = await client.get(f"{doc_url.rstrip('/')}/docling/capabilities") r.raise_for_status() return r.json() except Exception as exc: return {"ok": False, "installed": False, "error": str(exc)} @admin_router.post("/docling/convert") async def docling_convert_proxy(body: DoclingConvertBody) -> dict[str, Any]: doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") try: async with httpx.AsyncClient(timeout=600.0) as client: r = await client.post(f"{doc_url.rstrip('/')}/docling/convert", json=body.model_dump()) if r.status_code >= 400: raise HTTPException(r.status_code, r.text[:500]) return r.json() except HTTPException: raise except Exception as exc: raise HTTPException(502, str(exc)) from exc @admin_router.post("/docling/batch") async def docling_batch_proxy(body: DoclingBatchBody) -> dict[str, Any]: doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") try: async with httpx.AsyncClient(timeout=900.0) as client: r = await client.post(f"{doc_url.rstrip('/')}/docling/batch", json=body.model_dump()) if r.status_code >= 400: raise HTTPException(r.status_code, r.text[:500]) return r.json() except HTTPException: raise except Exception as exc: raise HTTPException(502, str(exc)) from exc @admin_router.get("/docling/history") async def docling_history_proxy(limit: int = 30) -> dict[str, Any]: doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") try: async with httpx.AsyncClient(timeout=30.0) as client: r = await client.get(f"{doc_url.rstrip('/')}/docling/history", params={"limit": limit}) r.raise_for_status() return r.json() except Exception as exc: return {"ok": False, "items": [], "error": str(exc)} @admin_router.get("/docling/result") async def docling_result_proxy(path: str, file_sig: str | None = None) -> dict[str, Any]: doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") params: dict[str, str] = {"path": path} if file_sig: params["file_sig"] = file_sig try: async with httpx.AsyncClient(timeout=30.0) as client: r = await client.get(f"{doc_url.rstrip('/')}/docling/result", params=params) if r.status_code == 404: return {"ok": False, "result": None} r.raise_for_status() return r.json() except Exception as exc: raise HTTPException(502, str(exc)) from exc class WorkspaceBody(BaseModel): storage_path: str = Field(..., min_length=1) content: str = "" source_format: str = "markdown" class NasWriteBody(BaseModel): path: str = Field(..., min_length=1) content: str = "" subdir: str = "Telegram/Exports" class PptxSlideBody(BaseModel): title: str = "" bullets: list[str] = Field(default_factory=list) class PptxCreateBody(BaseModel): title: str = Field(..., min_length=1, max_length=200) subtitle: str = Field(default="Foodlinkk", max_length=200) slides: list[PptxSlideBody] = Field(default_factory=list) filename: str = "" class DocumentsChatMessage(BaseModel): role: str = Field(..., pattern="^(user|assistant)$") content: str = "" class DocumentsChatBody(BaseModel): message: str = Field(..., min_length=1, max_length=8000) path_prefix: str = "" client_id: int | None = None project_id: int | None = None use_herman: bool = False history: list[DocumentsChatMessage] = Field(default_factory=list) class DocumentLinkBody(BaseModel): storage_path: str = Field(..., min_length=1) client_id: int | None = None project_id: int | None = None is_folder: bool = False notes: str = "" @admin_router.get("/docling/workspace") def docling_workspace_get(path: str) -> dict[str, Any]: try: row = fetch_one( "SELECT storage_path, source_format, content, updated_at FROM document_workspace WHERE storage_path = %s", (path,), ) except Exception as exc: raise HTTPException(500, str(exc)) from exc if not row: return {"ok": True, "workspace": None} return {"ok": True, "workspace": _serialize(row)} @admin_router.put("/docling/workspace") def docling_workspace_put(body: WorkspaceBody) -> dict[str, Any]: try: row = fetch_one( """ INSERT INTO document_workspace (storage_path, source_format, content, updated_at) VALUES (%s, %s, %s, NOW()) ON CONFLICT (storage_path) DO UPDATE SET source_format = EXCLUDED.source_format, content = EXCLUDED.content, updated_at = NOW() RETURNING storage_path, source_format, content, updated_at """, (body.storage_path, body.source_format, body.content), ) except Exception as exc: raise HTTPException(500, str(exc)) from exc return {"ok": True, "workspace": _serialize(row)} @admin_router.get("/docling/nas-read") async def docling_nas_read_proxy(path: str) -> dict[str, Any]: doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") try: async with httpx.AsyncClient(timeout=60.0) as client: r = await client.get( f"{doc_url.rstrip('/')}/docling/nas/read", params={"path": path}, ) if r.status_code >= 400: raise HTTPException(r.status_code, r.text[:500]) return r.json() except HTTPException: raise except Exception as exc: raise HTTPException(502, str(exc)) from exc @admin_router.get("/documents/nas-file") async def documents_nas_file_proxy(path: str): from fastapi.responses import Response doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") try: async with httpx.AsyncClient(timeout=120.0) as client: r = await client.get( f"{doc_url.rstrip('/')}/docling/nas/file", params={"path": path}, ) if r.status_code >= 400: raise HTTPException(r.status_code, r.text[:300]) media = r.headers.get("content-type", "application/octet-stream") return Response( content=r.content, media_type=media, headers={"Content-Disposition": f'inline; filename="{path.split("/")[-1]}"'}, ) except HTTPException: raise except Exception as exc: raise HTTPException(502, str(exc)) from exc @admin_router.post("/docling/save-nas") async def docling_save_nas(body: NasWriteBody) -> dict[str, Any]: doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") try: async with httpx.AsyncClient(timeout=60.0) as client: r = await client.post(f"{doc_url.rstrip('/')}/docling/nas/write", json=body.model_dump()) if r.status_code >= 400: raise HTTPException(r.status_code, r.text[:500]) return r.json() except HTTPException: raise except Exception as exc: raise HTTPException(502, str(exc)) from exc @admin_router.post("/pptx/create") async def pptx_create_proxy(body: PptxCreateBody) -> dict[str, Any]: doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") payload = { "title": body.title, "subtitle": body.subtitle, "filename": body.filename, "slides": [{"title": s.title, "bullets": s.bullets} for s in body.slides], } try: async with httpx.AsyncClient(timeout=120.0) as client: r = await client.post(f"{doc_url.rstrip('/')}/convert/create-pptx", json=payload) if r.status_code >= 400: raise HTTPException(r.status_code, r.text[:500]) return r.json() except HTTPException: raise except Exception as exc: raise HTTPException(502, str(exc)) from exc @admin_router.post("/pptx/pdf-to-pptx") async def pptx_from_pdf_proxy( path: str, mode: str = "smart", max_pages: int = 50, ) -> dict[str, Any]: doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") try: async with httpx.AsyncClient(timeout=300.0) as client: r = await client.post( f"{doc_url.rstrip('/')}/convert/pdf-to-pptx", params={"path": path, "mode": mode, "max_pages": max_pages}, ) if r.status_code >= 400: raise HTTPException(r.status_code, r.text[:500]) return r.json() except HTTPException: raise except Exception as exc: raise HTTPException(502, str(exc)) from exc @admin_router.get("/pptx/download") async def pptx_download_proxy(path: str): from fastapi.responses import Response doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") try: async with httpx.AsyncClient(timeout=60.0) as client: r = await client.get( f"{doc_url.rstrip('/')}/convert/pdf-to-pptx/download", params={"path": path}, ) if r.status_code != 200: raise HTTPException(r.status_code, r.text[:300]) return Response( content=r.content, media_type="application/vnd.openxmlformats-officedocument.presentationml.presentation", headers={"Content-Disposition": f'attachment; filename="{path.split("/")[-1]}"'}, ) except HTTPException: raise except Exception as exc: raise HTTPException(502, str(exc)) from exc @admin_router.get("/documents/{storage_path:path}/words") def document_words(storage_path: str, limit: int = 100, stopwords: bool = False): limit = max(1, min(limit, 500)) clauses = ["storage_path = %s"] params: list = [storage_path] if not stopwords: clauses.append("NOT is_stopword") try: rows = fetch_all( f""" SELECT lemma, token, count, pos_tag, is_stopword, language FROM document_word_counts WHERE {' AND '.join(clauses)} ORDER BY count DESC LIMIT %s """, tuple(params + [limit]), ) except Exception as exc: raise HTTPException(500, str(exc)) from exc return {"ok": True, "items": [_serialize(r) for r in rows]} BROWSER_AGENT_URL = os.getenv("BROWSER_AGENT_URL", "http://browser-agent:7790") class BrowserBrowseBody(BaseModel): url: str task: Optional[str] = None site_id: Optional[int] = None wait_seconds: float = 4.0 @admin_router.post("/browser/browse") async def browser_browse(body: BrowserBrowseBody) -> dict[str, Any]: import httpx try: async with httpx.AsyncClient(timeout=120.0) as client: r = await client.post(f"{BROWSER_AGENT_URL.rstrip('/')}/browse", json=body.model_dump()) r.raise_for_status() data = r.json() except httpx.HTTPStatusError as exc: detail = exc.response.text[:500] raise HTTPException(status_code=exc.response.status_code, detail=detail) from exc except Exception as exc: raise HTTPException(status_code=502, detail=str(exc)) from exc try: meta = data.get("session") or data steps = (meta.get("metadata") or {}).get("steps") or [] log_agent_event( "browser", "browse", f"Browse: {(meta.get('title') or body.url)[:100]}", body.task or body.url, channel="browser", metadata={ "url": body.url, "final_url": meta.get("final_url"), "session_id": meta.get("id"), "steps": steps[:12], }, ) except Exception: pass return data @admin_router.get("/browser/sessions") async def browser_sessions(limit: int = 20) -> dict[str, Any]: import httpx async with httpx.AsyncClient(timeout=30.0) as client: r = await client.get(f"{BROWSER_AGENT_URL.rstrip('/')}/sessions", params={"limit": limit}) r.raise_for_status() return r.json() @admin_router.get("/browser/sessions/{session_id}") async def browser_session_detail(session_id: int, include_screenshot: bool = True) -> dict[str, Any]: import httpx async with httpx.AsyncClient(timeout=30.0) as client: r = await client.get( f"{BROWSER_AGENT_URL.rstrip('/')}/sessions/{session_id}", params={"include_screenshot": include_screenshot}, ) r.raise_for_status() return r.json() @admin_router.get("/browser/live/screenshot.jpg") async def browser_live_screenshot(): import httpx from fastapi.responses import Response try: async with httpx.AsyncClient(timeout=15.0) as client: r = await client.get(f"{BROWSER_AGENT_URL.rstrip('/')}/live/screenshot.jpg") if r.status_code == 404: return Response(status_code=404) return Response(content=r.content, media_type="image/jpeg") except Exception: return Response(status_code=404) @admin_router.get("/browser/sessions/{session_id}/screenshot.jpg") async def browser_session_screenshot(session_id: int): import httpx from fastapi.responses import Response async with httpx.AsyncClient(timeout=30.0) as client: r = await client.get(f"{BROWSER_AGENT_URL.rstrip('/')}/sessions/{session_id}/screenshot.jpg") if r.status_code != 200: raise HTTPException(status_code=r.status_code, detail="Screenshot not found") return Response(content=r.content, media_type="image/jpeg") class BrowserInstructBody(BaseModel): url: str instruction: str wait_seconds: float = 5.0 @admin_router.post("/browser/instruct") async def browser_instruct(body: BrowserInstructBody) -> dict[str, Any]: import httpx async with httpx.AsyncClient(timeout=180.0) as client: r = await client.post(f"{BROWSER_AGENT_URL.rstrip('/')}/browse/instruct", json=body.model_dump()) r.raise_for_status() data = r.json() try: meta = data.get("session") or data steps = (meta.get("metadata") or {}).get("steps") or [body.instruction[:200]] log_agent_event( "browser", "browse_instruct", f"Instructie: {body.instruction[:80]}", " ยท ".join(steps[:8])[:500], channel="browser", metadata={"url": body.url, "instruction": body.instruction, "session_id": meta.get("id")}, ) except Exception: pass return data @admin_router.get("/browser/sessions/{session_id}/ocr") async def browser_session_ocr(session_id: int) -> dict[str, Any]: import httpx async with httpx.AsyncClient(timeout=60.0) as client: r = await client.get(f"{BROWSER_AGENT_URL.rstrip('/')}/sessions/{session_id}/ocr") r.raise_for_status() return r.json() class BrowserVncBody(BaseModel): url: str instruction: Optional[str] = None wait_seconds: float = 4.0 class BrowserExtractFullBody(BaseModel): url: str instruction: Optional[str] = None wait_seconds: float = 4.0 scroll_pages: int = 6 site_id: Optional[int] = None also_vnc: bool = True class PhotoUploadBody(BaseModel): image_b64: str source: str = "cockpit" filename: Optional[str] = None storage_path: Optional[str] = None session_id: Optional[int] = None @admin_router.post("/browser/vnc-navigate") async def browser_vnc_navigate(body: BrowserVncBody) -> dict[str, Any]: import httpx async with httpx.AsyncClient(timeout=120.0) as client: r = await client.post(f"{BROWSER_AGENT_URL.rstrip('/')}/vnc/navigate", json=body.model_dump()) if r.status_code >= 400: raise HTTPException(status_code=r.status_code, detail=r.text[:500]) data = r.json() try: log_agent_event( "browser", "vnc_navigate", f"VNC โ†’ {body.url[:100]}", body.instruction or "navigate", channel="browser-vnc", metadata={"url": body.url, "instruction": body.instruction}, ) except Exception: pass return data @admin_router.post("/browser/extract-full") async def browser_extract_full(body: BrowserExtractFullBody) -> dict[str, Any]: import httpx async with httpx.AsyncClient(timeout=300.0) as client: r = await client.post(f"{BROWSER_AGENT_URL.rstrip('/')}/browse/extract-full", json=body.model_dump()) if r.status_code >= 400: raise HTTPException(status_code=r.status_code, detail=r.text[:500]) data = r.json() try: meta = data.get("session") or data log_agent_event( "browser", "extract_full", f"Extract: {(meta.get('title') or body.url)[:100]}", body.instruction or body.url, channel="browser", metadata={"url": body.url, "session_id": meta.get("id"), "scroll_pages": body.scroll_pages}, ) except Exception: pass return data @admin_router.get("/photos") async def photos_list(limit: int = 30) -> dict[str, Any]: import httpx async with httpx.AsyncClient(timeout=30.0) as client: r = await client.get(f"{BROWSER_AGENT_URL.rstrip('/')}/photos", params={"limit": limit}) r.raise_for_status() return r.json() @admin_router.get("/photos/{photo_id}") async def photos_detail(photo_id: int) -> dict[str, Any]: import httpx async with httpx.AsyncClient(timeout=30.0) as client: r = await client.get(f"{BROWSER_AGENT_URL.rstrip('/')}/photos/{photo_id}") r.raise_for_status() return r.json() @admin_router.get("/photos/{photo_id}/image.jpg") async def photos_image(photo_id: int): import httpx from fastapi.responses import Response async with httpx.AsyncClient(timeout=60.0) as client: r = await client.get(f"{BROWSER_AGENT_URL.rstrip('/')}/photos/{photo_id}/image.jpg") if r.status_code != 200: raise HTTPException(status_code=r.status_code) return Response(content=r.content, media_type="image/jpeg") @admin_router.get("/photos/{photo_id}/detections") async def photos_detections(photo_id: int) -> dict[str, Any]: import httpx async with httpx.AsyncClient(timeout=30.0) as client: r = await client.get(f"{BROWSER_AGENT_URL.rstrip('/')}/photos/{photo_id}/detections") r.raise_for_status() return r.json() @admin_router.post("/photos/analyze") async def photos_analyze(body: PhotoUploadBody) -> dict[str, Any]: import httpx async with httpx.AsyncClient(timeout=120.0) as client: r = await client.post(f"{BROWSER_AGENT_URL.rstrip('/')}/photos/analyze", json=body.model_dump()) if r.status_code >= 400: raise HTTPException(status_code=r.status_code, detail=r.text[:500]) return r.json() @admin_router.get("/documents/share-files") async def documents_share_files( limit: int = 200, ext: Optional[str] = None, ) -> dict[str, Any]: """Alle bestanden op NAS share โ€” incl. pptx, pdf, etc.""" import httpx doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") limit = max(1, min(limit, 1000)) params: dict[str, Any] = {"limit": limit} if ext: params["ext"] = ext try: async with httpx.AsyncClient(timeout=120.0) as client: r = await client.get(f"{doc_url.rstrip('/')}/nas/files", params=params) r.raise_for_status() data = r.json() return {"ok": True, **data} except Exception as exc: return {"ok": False, "items": [], "error": str(exc)} @admin_router.post("/documents/trigger-scan") async def documents_trigger_scan(force: bool = False) -> dict[str, Any]: """Start doc-ingest scan zodat nieuwe share-bestanden worden herkend.""" import httpx doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") try: async with httpx.AsyncClient(timeout=120.0) as client: r = await client.post( f"{doc_url.rstrip('/')}/ingest/scan", params={"force": "true" if force else "false"}, ) r.raise_for_status() return {"ok": True, **r.json()} except Exception as exc: raise HTTPException(status_code=502, detail=str(exc)) from exc @admin_router.get("/documents/search") async def documents_search( q: str, limit: int = 8, ) -> dict[str, Any]: """Semantisch zoeken in NAS-documenten (Chroma RAG).""" doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") try: async with httpx.AsyncClient(timeout=60.0) as client: r = await client.get( f"{doc_url.rstrip('/')}/search", params={"q": q, "limit": max(1, min(limit, 20))}, ) r.raise_for_status() return {"ok": True, **r.json()} except Exception as exc: return {"ok": False, "results": [], "error": str(exc)} @admin_router.get("/documents/nas-diagnostics") async def documents_nas_diagnostics() -> dict[str, Any]: """NAS mount status + zichtbare vs. ontbrekende bestanden.""" doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") out: dict[str, Any] = {"ok": True} try: async with httpx.AsyncClient(timeout=30.0) as client: h = await client.get(f"{doc_url.rstrip('/')}/health") h.raise_for_status() out["health"] = h.json() f = await client.get(f"{doc_url.rstrip('/')}/nas/files", params={"limit": 1000}) f.raise_for_status() files = f.json() out["visible_files"] = files.get("total", 0) out["sample_paths"] = [i.get("path") for i in (files.get("items") or [])[:8]] except Exception as exc: out["ok"] = False out["error"] = str(exc) out["hint"] = ( "Zie je mappen (CUCINA/Foodlinkk) maar weinig bestanden? " "Geef gebruiker aissa Read/Write op de Synology share + submappen in DSM." ) return out @admin_router.post("/documents/chat") async def documents_chat(body: DocumentsChatBody) -> dict[str, Any]: """Chat met NAS-data via RAG + Ollama (snel) of optioneel Herman.""" doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") message = body.message.strip() sources: list[str] = [] context = "" linked_paths: list[str] = [] if body.client_id: rows = fetch_all( "SELECT storage_path, is_folder FROM document_links WHERE client_id = %s", (body.client_id,), ) linked_paths = [r["storage_path"] for r in rows if r.get("storage_path")] if body.project_id: rows = fetch_all( "SELECT storage_path, is_folder FROM document_links WHERE project_id = %s", (body.project_id,), ) linked_paths.extend([r["storage_path"] for r in rows if r.get("storage_path")]) try: async with httpx.AsyncClient(timeout=45.0) as client: r = await client.get( f"{doc_url.rstrip('/')}/context", params={"q": message, "limit": 8}, ) r.raise_for_status() data = r.json() context = (data.get("context") or "").strip() sources = [s for s in (data.get("sources") or []) if s] except Exception: context = "" sources = [] scope_prefix = (body.path_prefix or "").strip().lstrip("/") if linked_paths or scope_prefix: try: async with httpx.AsyncClient(timeout=45.0) as client: sr = await client.get( f"{doc_url.rstrip('/')}/search", params={"q": message, "limit": 12}, ) sr.raise_for_status() hits = sr.json().get("results") or [] def _in_scope(path: str) -> bool: if not path: return False if scope_prefix and (path == scope_prefix or path.startswith(scope_prefix + "/")): return True for lp in linked_paths: if path == lp or path.startswith(lp.rstrip("/") + "/"): return True return False scoped = [h for h in hits if _in_scope(h.get("path") or "")] if scoped: parts = [f"[{h.get('filename', 'doc')}] {(h.get('text') or '')[:1200]}" for h in scoped[:8]] context = "\n\n".join(parts) sources = [h.get("path") for h in scoped if h.get("path")] except Exception: pass if body.use_herman: try: scoped_msg = message if scope_prefix: scoped_msg = f"[Scope: {scope_prefix}]\n{message}" result = await herman_service.chat(scoped_msg, channel="documents") result["rag_sources"] = result.get("rag_sources") or sources result["ok"] = True result["context_used"] = bool(context) return result except Exception as exc: pass client_note = "" if body.client_id: c = fetch_one("SELECT name FROM clients WHERE id = %s", (body.client_id,)) if c: client_note = f"\nKlantcontext: {c['name']} (ID {body.client_id})" system = ( "Je bent Herman, Foodlinkk document-assistent en second brain. Beantwoord in het Nederlands " "op basis van de NAS-documentfragmenten hieronder. Citeer bronbestanden als [bestandsnaam]. " "Als het antwoord niet in de context staat, zeg dat eerlijk." f"{client_note}\n\nDOCUMENT CONTEXT:\n{context or '(geen documenten gevonden โ€” geef algemeen antwoord)'}" ) messages: list[dict[str, str]] = [{"role": "system", "content": system}] for h in body.history[-6:]: messages.append({"role": h.role, "content": h.content}) messages.append({"role": "user", "content": message}) try: reply = await ollama.chat_messages(messages, timeout=120.0) except Exception as exc: raise HTTPException(502, f"LLM niet bereikbaar: {exc}") from exc return { "ok": True, "agent": "knowledge", "agent_label": "Herman ยท NAS RAG", "reply": reply, "rag_sources": sources, "context_used": bool(context), "delegated_agents": ["knowledge"], } @admin_router.get("/documents/nas-images") async def documents_nas_images(limit: int = 40) -> dict[str, Any]: import httpx doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") exts = (".jpg", ".jpeg", ".png", ".webp", ".gif", ".bmp") try: async with httpx.AsyncClient(timeout=30.0) as client: r = await client.get(f"{doc_url.rstrip('/')}/nas/list", params={"limit": 500}) r.raise_for_status() data = r.json() files = data.get("files") or [] images = [f for f in files if (f.get("path") or "").lower().endswith(exts)] return {"ok": True, "items": images[:limit], "total": len(images)} except Exception as exc: return {"ok": False, "items": [], "error": str(exc)} @admin_router.post("/documents/nas-ocr") async def documents_nas_ocr( path: str | None = None, force: bool = False, ) -> dict[str, Any]: """OCR alle NAS-afbeeldingen (jpg/png) + plaatjes uit pdf/docx via doc-ingest โ†’ browser-agent.""" doc_url = os.getenv("DOC_INGEST_URL", "http://10.4.7.19:8750") params: dict[str, str] = {"force": "true" if force else "false"} if path: params["path"] = path try: async with httpx.AsyncClient(timeout=600.0) as client: r = await client.post( f"{doc_url.rstrip('/')}/ingest/extract-images", params=params, ) r.raise_for_status() data = r.json() try: log_agent_event( "browser", "nas_ocr_batch", f"NAS OCR: {data.get('images_extracted', 0)} afbeeldingen", path or "alle NAS-bestanden", channel="documents", metadata={"path": path, "force": force, **data}, ) except Exception: pass return {"ok": True, **data} except Exception as exc: raise HTTPException(status_code=502, detail=str(exc)) from exc @admin_router.get("/browser/vnc/screenshot.jpg") async def browser_vnc_screenshot(): import httpx from fastapi.responses import Response try: async with httpx.AsyncClient(timeout=15.0) as client: r = await client.get(f"{BROWSER_AGENT_URL.rstrip('/')}/vnc/screenshot.jpg") if r.status_code != 200: return Response(status_code=404) return Response(content=r.content, media_type="image/jpeg") except Exception: return Response(status_code=404) # --- Hermes / Telegram second brain proxies --- class HermesSearchBody(BaseModel): query: str = Field(..., min_length=2) chat_id: Optional[int] = None limit: int = Field(default=15, ge=1, le=30) @admin_router.get("/hermes/stats") async def hermes_stats_proxy() -> dict[str, Any]: import httpx async with httpx.AsyncClient(timeout=30.0) as client: r = await client.get(f"{settings.TOOLS_API_URL.rstrip('/')}/brain/stats") r.raise_for_status() return r.json() @admin_router.get("/hermes/conversations") async def hermes_conversations_proxy(limit: int = 50) -> dict[str, Any]: import httpx async with httpx.AsyncClient(timeout=30.0) as client: r = await client.get(f"{settings.TOOLS_API_URL.rstrip('/')}/brain/conversations", params={"limit": limit}) r.raise_for_status() return r.json() @admin_router.get("/hermes/feed") async def hermes_feed_proxy( chat_id: Optional[int] = None, limit: int = 80, offset: int = 0, ) -> dict[str, Any]: import httpx params = {"limit": limit, "offset": offset} if chat_id is not None: params["chat_id"] = chat_id async with httpx.AsyncClient(timeout=30.0) as client: r = await client.get(f"{settings.TOOLS_API_URL.rstrip('/')}/brain/feed", params=params) r.raise_for_status() return r.json() @admin_router.get("/hermes/graph") async def hermes_graph_global(limit: int = 100) -> dict[str, Any]: import httpx async with httpx.AsyncClient(timeout=30.0) as client: r = await client.get(f"{settings.TOOLS_API_URL.rstrip('/')}/brain/graph", params={"limit": limit}) r.raise_for_status() return r.json() @admin_router.get("/hermes/graph/{chat_id}") async def hermes_graph_chat(chat_id: int, limit: int = 80) -> dict[str, Any]: import httpx async with httpx.AsyncClient(timeout=30.0) as client: r = await client.get(f"{settings.TOOLS_API_URL.rstrip('/')}/brain/graph/{chat_id}", params={"limit": limit}) r.raise_for_status() return r.json() @admin_router.post("/hermes/search") async def hermes_search_proxy(body: HermesSearchBody) -> dict[str, Any]: import httpx async with httpx.AsyncClient(timeout=120.0) as client: r = await client.post( f"{settings.TOOLS_API_URL.rstrip('/')}/brain/search", json=body.model_dump(), ) r.raise_for_status() return r.json() # --- Hermes PA live + users + control (added by deploy) --- HERMES_CONTROL_URL = os.getenv("HERMES_CONTROL_URL", "http://10.4.7.19:8799").rstrip("/") @admin_router.get("/hermes/pa/live") async def hermes_pa_live_proxy() -> dict[str, Any]: import httpx async with httpx.AsyncClient(timeout=15.0) as client: r = await client.get(f"{BROWSER_AGENT_URL.rstrip('/')}/pa/live") r.raise_for_status() live = r.json() try: async with httpx.AsyncClient(timeout=8.0) as client: ctrl = await client.get(f"{HERMES_CONTROL_URL}/status") if ctrl.status_code == 200: live["hermes"] = ctrl.json() except Exception: live["hermes"] = {"online": False} return live @admin_router.get("/hermes/pa/live/{label}/screenshot.jpg") async def hermes_pa_slot_screenshot(label: str): import httpx from fastapi.responses import Response async with httpx.AsyncClient(timeout=20.0) as client: r = await client.get( f"{BROWSER_AGENT_URL.rstrip('/')}/pa/live/{label}/screenshot.jpg" ) if r.status_code != 200: raise HTTPException(status_code=r.status_code, detail="No screenshot") return Response(content=r.content, media_type="image/jpeg") @admin_router.get("/hermes/users") async def hermes_users_proxy() -> dict[str, Any]: import httpx users: list[dict[str, Any]] = [] try: async with httpx.AsyncClient(timeout=8.0) as client: r = await client.get(f"{HERMES_CONTROL_URL}/users") if r.status_code == 200: return r.json() except Exception: pass async with httpx.AsyncClient(timeout=15.0) as client: conv = await client.get(f"{settings.TOOLS_API_URL.rstrip('/')}/brain/conversations", params={"limit": 20}) items = conv.json().get("items") or [] if conv.status_code == 200 else [] default = [ {"chat_id": 8859782446, "name": "Aรฏssa", "role": "CEO", "allowed": True, "pa_mode": False, "online": True}, {"chat_id": 789036463, "name": "Mo", "role": "CTO", "allowed": True, "pa_mode": False, "online": False}, ] by_id = {u["chat_id"]: u for u in default} for c in items: cid = c.get("chat_id") if cid in by_id: by_id[cid]["message_count"] = c.get("message_count") by_id[cid]["last_message_at"] = c.get("last_message_at") by_id[cid]["online"] = True else: by_id[cid] = { "chat_id": cid, "name": c.get("user_name") or str(cid), "role": c.get("user_role") or "user", "allowed": True, "pa_mode": False, "online": True, "message_count": c.get("message_count"), "last_message_at": c.get("last_message_at"), } return {"users": list(by_id.values())} @admin_router.post("/hermes/actions/{action}") async def hermes_action_proxy(action: str) -> dict[str, Any]: import httpx allowed = {"evening-briefing", "morning-briefing", "resume-jobs", "test-briefing"} if action not in allowed: raise HTTPException(status_code=400, detail=f"Unknown action: {action}") async with httpx.AsyncClient(timeout=180.0) as client: r = await client.post(f"{HERMES_CONTROL_URL}/actions/{action}") if r.status_code >= 400: raise HTTPException(status_code=r.status_code, detail=r.text[:300]) return r.json()