"""War Room + 1:1 agent chat — Hermes/mimo only, soul prompts, A2A rounds.""" from __future__ import annotations import asyncio import logging import re from typing import Any from app.config import settings from app.db import execute, fetch_all, fetch_one from app.services import agent_nas, agent_souls, hermes_memory, llm_router from app.services.agent_names import normalize_agent_key log = logging.getLogger("cockpit.agent_chat") MENTION_RE = re.compile(r"@([a-zA-Z][a-zA-Z0-9_-]{1,63})") MAX_ROUNDS = 4 MAX_REPLIES_PER_ROUND = 3 HISTORY_LIMIT = 16 WARROOM_MAX_TOKENS = 512 WARROOM_TIMEOUT = 55.0 A2A_CONTINUE_MAX = 1 COCKPIT_SNAPSHOT_TTL = 60.0 _snapshot_cache: dict[str, Any] = {"at": 0.0, "text": ""} def cockpit_snapshot(*, max_chars: int = 2200) -> str: """Compact live cockpit facts for War Room agents — accuracy first.""" import time now = time.time() if _snapshot_cache.get("text") and (now - float(_snapshot_cache.get("at") or 0)) < COCKPIT_SNAPSHOT_TTL: return str(_snapshot_cache["text"]) lines: list[str] = [ "COCKPIT FEITEN (alleen dit mag je citeren; verzin niets erbuiten):", "SCHEIDING: CRM-klanten/deals ≠ retail-supermarktlocaties. Nooit CRM-namen als AH-filiaal presenteren.", ] # CRM try: row = fetch_one( """SELECT (SELECT COUNT(*) FROM clients) AS clients, (SELECT COUNT(*) FROM deals) AS deals, (SELECT COALESCE(SUM(value),0) FROM deals WHERE stage NOT IN ('won','lost')) AS pipeline """ ) or {} lines.append( f"- CRM: {row.get('clients', 0)} klanten · {row.get('deals', 0)} deals · " f"pipeline €{float(row.get('pipeline') or 0):,.0f}" ) deals = fetch_all( """SELECT d.title, d.stage, d.value, c.name AS client_name FROM deals d LEFT JOIN clients c ON c.id = d.client_id WHERE d.stage NOT IN ('won','lost') ORDER BY d.updated_at DESC NULLS LAST LIMIT 6""" ) or [] if deals: bits = [] for d in deals: bits.append( f"{d.get('client_name') or d.get('title') or '?'}[{d.get('stage')}" f"/€{float(d.get('value') or 0):,.0f}]" ) lines.append("- Open deals: " + "; ".join(bits)) clients = fetch_all("SELECT name, stage FROM clients ORDER BY updated_at DESC NULLS LAST LIMIT 6") or [] if clients: lines.append( "- CRM-klanten: " + ", ".join(f"{c.get('name')}({c.get('stage')})" for c in clients) ) except Exception as exc: log.warning("snapshot CRM failed: %s", exc) lines.append("- CRM: (niet beschikbaar)") # Real AH / retail opportunities — exclude test rows try: ah = fetch_all( """ SELECT s.id, s.name, s.address, s.postcode, s.city, ros.halal_opportunity_score AS score, COALESCE(s.partnership_status, ros.factors->>'partnership_status', 'none') AS status FROM retail_opportunity_scores ros JOIN supermarkets s ON s.id = ros.supermarket_id WHERE s.chain ILIKE '%albert%heijn%' AND s.name NOT ILIKE '%masala%' AND s.name NOT ILIKE '%TEST%' AND COALESCE(s.data_source,'') NOT ILIKE '%test_exclude%' ORDER BY ros.halal_opportunity_score DESC NULLS LAST, s.city, s.postcode LIMIT 8 """ ) or [] if ah: lines.append("- Top Albert Heijn (halal-kans score, echte filialen):") for s in ah: lines.append( f" · AH {s.get('address')}, {s.get('postcode')} {s.get('city')} " f"— score {float(s.get('score') or 0):.1f} · status {s.get('status')}" ) else: lines.append("- Top Albert Heijn: geen geldige scores (geen testrecords).") except Exception as exc: log.warning("snapshot AH failed: %s", exc) lines.append("- Top Albert Heijn: (niet beschikbaar)") text = "\n".join(lines)[:max_chars] _snapshot_cache["at"] = now _snapshot_cache["text"] = text return text def _iso(row: dict[str, Any] | None) -> dict[str, Any] | None: if not row: return None out = dict(row) for k, v in list(out.items()): if hasattr(v, "isoformat"): out[k] = v.isoformat() if isinstance(out.get("metadata"), str): try: import json out["metadata"] = json.loads(out["metadata"]) except Exception: out["metadata"] = {} return out def _bootstrap_agents_nas() -> None: try: keys = [str(s.get("agent_key")) for s in agent_souls.list_souls()] agent_nas.bootstrap_all(keys) except Exception: log.exception("agents NAS bootstrap failed") def ensure_warroom() -> dict[str, Any]: row = fetch_one("SELECT * FROM warroom_rooms WHERE kind = 'warroom' LIMIT 1") if row: return _iso(row) # type: ignore row = fetch_one( """INSERT INTO warroom_rooms (kind, title) VALUES ('warroom', 'War Room') RETURNING *""" ) return _iso(row) # type: ignore def get_or_create_dm(agent_key: str) -> dict[str, Any]: key = normalize_agent_key(agent_key) or (agent_key or "").strip().lower() soul = agent_souls.get_soul(key) if not soul: raise ValueError(f"Onbekende agent: {agent_key}") row = fetch_one("SELECT * FROM warroom_rooms WHERE kind = 'dm' AND agent_key = %s", (key,)) if row: return _iso(row) # type: ignore title = f"DM · {soul.get('display_name') or key}" row = fetch_one( """INSERT INTO warroom_rooms (kind, agent_key, title) VALUES ('dm', %s, %s) RETURNING *""", (key, title), ) return _iso(row) # type: ignore def list_rooms() -> list[dict[str, Any]]: ensure_warroom() rows = fetch_all( """SELECT r.*, (SELECT COUNT(*) FROM warroom_messages m WHERE m.room_id = r.id) AS message_count, (SELECT content FROM warroom_messages m WHERE m.room_id = r.id ORDER BY m.created_at DESC LIMIT 1) AS last_message FROM warroom_rooms r ORDER BY CASE WHEN r.kind = 'warroom' THEN 0 ELSE 1 END, r.title""" ) return [_iso(r) for r in rows] # type: ignore def get_message(message_id: int, room_id: int | None = None) -> dict[str, Any] | None: if room_id is not None: return _iso(fetch_one( "SELECT * FROM warroom_messages WHERE id = %s AND room_id = %s AND deleted_at IS NULL", (message_id, room_id), )) return _iso(fetch_one( "SELECT * FROM warroom_messages WHERE id = %s AND deleted_at IS NULL", (message_id,), )) def get_room(room_id: int) -> dict[str, Any] | None: return _iso(fetch_one("SELECT * FROM warroom_rooms WHERE id = %s", (room_id,))) def list_messages( room_id: int, *, limit: int = 100, after_id: int | None = None, before_id: int | None = None, ) -> list[dict[str, Any]]: """List messages. Default = newest page. before_id = older page (scroll up).""" lim = max(1, min(int(limit), 5000)) if after_id: rows = fetch_all( """SELECT * FROM warroom_messages WHERE room_id = %s AND id > %s AND deleted_at IS NULL ORDER BY created_at ASC, id ASC LIMIT %s""", (room_id, after_id, lim), ) elif before_id: rows = fetch_all( """SELECT * FROM ( SELECT * FROM warroom_messages WHERE room_id = %s AND id < %s AND deleted_at IS NULL ORDER BY created_at DESC, id DESC LIMIT %s ) t ORDER BY created_at ASC, id ASC""", (room_id, before_id, lim), ) else: rows = fetch_all( """SELECT * FROM ( SELECT * FROM warroom_messages WHERE room_id = %s AND deleted_at IS NULL ORDER BY created_at DESC, id DESC LIMIT %s ) t ORDER BY created_at ASC, id ASC""", (room_id, lim), ) return [_iso(r) for r in rows] # type: ignore def count_messages(room_id: int) -> int: row = fetch_one( "SELECT COUNT(*) AS c FROM warroom_messages WHERE room_id = %s AND deleted_at IS NULL", (room_id,), ) return int((row or {}).get("c") or 0) def oldest_message_id(room_id: int) -> int | None: row = fetch_one( """SELECT id FROM warroom_messages WHERE room_id = %s AND deleted_at IS NULL ORDER BY id ASC LIMIT 1""", (room_id,), ) return int(row["id"]) if row and row.get("id") is not None else None def delete_message(room_id: int, message_id: int) -> dict[str, Any] | None: """Soft-delete a warroom message. Returns the row or None if missing/already deleted.""" row = fetch_one( """UPDATE warroom_messages SET deleted_at = NOW(), content = '', metadata = COALESCE(metadata, '{}'::jsonb) || '{"deleted": true}'::jsonb WHERE id = %s AND room_id = %s AND deleted_at IS NULL RETURNING *""", (message_id, room_id), ) if row: execute("UPDATE warroom_rooms SET updated_at = NOW() WHERE id = %s", (room_id,)) return _iso(row) if row else None # type: ignore def clear_messages(room_id: int) -> int: """Soft-delete all messages in a room. Returns count cleared.""" rows = fetch_all( """UPDATE warroom_messages SET deleted_at = NOW(), content = '', metadata = COALESCE(metadata, '{}'::jsonb) || '{"deleted": true, "cleared": true}'::jsonb WHERE room_id = %s AND deleted_at IS NULL RETURNING id""", (room_id,), ) execute("UPDATE warroom_rooms SET updated_at = NOW() WHERE id = %s", (room_id,)) return len(rows or []) def list_deleted_ids(room_id: int) -> list[int]: """Recently soft-deleted message ids (for live UI sync across clients).""" rows = fetch_all( """SELECT id FROM warroom_messages WHERE room_id = %s AND deleted_at IS NOT NULL AND deleted_at > NOW() - INTERVAL '6 hours' ORDER BY deleted_at DESC LIMIT 200""", (room_id,), ) return [int(r["id"]) for r in rows] def _insert_message( room_id: int, *, sender_type: str, content: str, agent_key: str | None = None, username: str | None = None, sender_avatar_url: str | None = None, metadata: dict | None = None, ) -> dict[str, Any]: import json row = fetch_one( """INSERT INTO warroom_messages (room_id, sender_type, agent_key, username, sender_avatar_url, content, metadata) VALUES (%s, %s, %s, %s, %s, %s, %s::jsonb) RETURNING *""", ( room_id, sender_type, agent_key, username, sender_avatar_url, content, json.dumps(metadata or {}), ), ) execute("UPDATE warroom_rooms SET updated_at = NOW() WHERE id = %s", (room_id,)) return _iso(row) # type: ignore def _active_agent_keys() -> set[str]: return {str(s.get("agent_key")) for s in agent_souls.list_souls() if s.get("is_active", True)} CEO_ALIASES = {"aissa", "ceo", "eigenaar", "ervoorz"} # ervoorz = legacy alias CEO_MENTION_RE = re.compile(r"@([A-Za-z0-9_\-]{2,40})", re.I) def parse_ceo_mentions(text: str) -> bool: """True if message is directed at CEO via @aissa (legacy: @ervoorz / @ceo).""" for m in CEO_MENTION_RE.findall(text or ""): if m.lower() in CEO_ALIASES: return True return False def _ceo_excerpt(text: str, limit: int = 280) -> str: t = (text or "").strip() if len(t) <= limit: return t return t[: limit - 1] + "…" def record_ceo_direct(message_row: dict[str, Any] | None) -> dict[str, Any] | None: """If an agent message @mentions CEO, store in ceo inbox.""" if not message_row or message_row.get("sender_type") != "agent": return None content = message_row.get("content") or "" if not parse_ceo_mentions(content): return None ensure_warroom() # ensure table exists (idempotent) try: execute( """CREATE TABLE IF NOT EXISTS warroom_ceo_inbox ( id BIGSERIAL PRIMARY KEY, message_id BIGINT, room_id BIGINT, agent_key VARCHAR(64), content TEXT NOT NULL, excerpt TEXT, created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), read_at TIMESTAMPTZ, recipient TEXT NOT NULL DEFAULT 'aissa' )""" ) except Exception: pass row = fetch_one( """INSERT INTO warroom_ceo_inbox (message_id, room_id, agent_key, content, excerpt, recipient) VALUES (%s, %s, %s, %s, %s, 'aissa') RETURNING *""", ( message_row.get("id"), message_row.get("room_id"), message_row.get("agent_key"), content, _ceo_excerpt(content), ), ) return _iso(row) if row else None def list_ceo_inbox(*, unread_only: bool = False, limit: int = 50) -> list[dict[str, Any]]: ensure_warroom() try: if unread_only: rows = fetch_all( """SELECT i.*, s.display_name AS agent_display_name FROM warroom_ceo_inbox i LEFT JOIN agent_souls s ON s.agent_key = i.agent_key WHERE i.recipient = 'aissa' AND i.read_at IS NULL ORDER BY i.created_at DESC LIMIT %s""", (limit,), ) else: rows = fetch_all( """SELECT i.*, s.display_name AS agent_display_name FROM warroom_ceo_inbox i LEFT JOIN agent_souls s ON s.agent_key = i.agent_key WHERE i.recipient = 'aissa' ORDER BY i.created_at DESC LIMIT %s""", (limit,), ) except Exception: return [] return [_iso(r) for r in (rows or [])] # type: ignore def ceo_inbox_unread_count() -> int: try: row = fetch_one( """SELECT COUNT(*) AS n FROM warroom_ceo_inbox WHERE recipient = 'aissa' AND read_at IS NULL""" ) return int((row or {}).get("n") or 0) except Exception: return 0 def mark_ceo_inbox_read(ids: list[int] | None = None, *, all_unread: bool = False) -> int: if all_unread: rows = fetch_all( """UPDATE warroom_ceo_inbox SET read_at = NOW() WHERE recipient = 'aissa' AND read_at IS NULL RETURNING id""" ) return len(rows or []) if not ids: return 0 rows = fetch_all( """UPDATE warroom_ceo_inbox SET read_at = NOW() WHERE recipient = 'aissa' AND id = ANY(%s) AND read_at IS NULL RETURNING id""", (ids,), ) return len(rows or []) def parse_mentions(text: str) -> list[str]: keys = _active_agent_keys() found: list[str] = [] for m in MENTION_RE.findall(text or ""): k = normalize_agent_key(m) or m.lower() if k in keys and k not in found: found.append(k) return found def _history_as_chat(messages: list[dict[str, Any]], *, limit: int = 20) -> list[dict[str, str]]: out: list[dict[str, str]] = [] for m in messages[-limit:]: st = m.get("sender_type") content = (m.get("content") or "").strip() if not content: continue if st == "user": name = m.get("username") or "User" out.append({"role": "user", "content": f"{name}: {content}"}) elif st == "agent": key = m.get("agent_key") or "agent" out.append({"role": "assistant", "content": f"@{key}: {content}"}) else: out.append({"role": "system", "content": content}) return out async def _hermes_reply( messages: list[dict[str, str]], *, hermes_user: str, timeout: float | None = None, ) -> tuple[str, dict[str, Any]]: """Force Hermes/mimo — never DeepSeek/OpenRouter fallback.""" prov = llm_router.resolve_provider(None) ptype = (prov.get("provider_type") or "").lower() if ptype != "hermes": hermes = llm_router._env_hermes_provider() if not hermes: raise RuntimeError("Hermes/mimo niet geconfigureerd (HERMES_API_KEY / :8642).") raise RuntimeError( f"War Room vereist Hermes/mimo, huidige provider is '{ptype or 'onbekend'}'. " "Zet HERMAN_LLM_BACKEND=hermes." ) reply, meta = await llm_router.chat_messages( messages, timeout=float(timeout if timeout is not None else WARROOM_TIMEOUT), user=hermes_user, max_tokens=WARROOM_MAX_TOKENS, temperature=0.35, ) meta = dict(meta or {}) if meta.get("provider_type") != "hermes": raise RuntimeError("Antwoord kwam niet van Hermes/mimo — afgebroken.") return reply, meta def _live_roster_line() -> str: keys = sorted(_active_agent_keys()) return ", ".join(keys) if keys else "(geen)" def _soul_system(soul: dict[str, Any], *, room_kind: str) -> str: name = soul.get("display_name") or soul.get("agent_key") key = soul.get("agent_key") role = soul.get("role_title") or "agent" soul_md = (soul.get("soul_md") or "")[:1600] resp = (soul.get("responsibilities") or "")[:350] roster = _live_roster_line() hard = ( "HARDE FEITEN (overschrijven alle eerdere aannames/geheugen):\n" f"- LIVE ROSTER ONLINE: {roster}.\n" "- Alle agents hierboven zijn NU actief in de War Room / cockpit.\n" "- Verboden zinnen: 'staat uit', 'niet actief', 'niet beschikbaar', 'weer inschakelen', " "'weer activeren', 'alleen Halal is actief', 'agent is offline'.\n" "- Als je twijfelt: behandel elke agent op de roster als online en antwoord inhoudelijk.\n" ) snap = cockpit_snapshot() accuracy = ( "NAUWKEURIGHEID (hard):\n" "- Antwoord ALLEEN met feiten uit COCKPIT FEITEN hieronder of uit recente War Room-berichten.\n" "- Weet je het niet 100% zeker: zeg expliciet 'Dat staat zo niet in de cockpit' en @mention de specialist " "(@retail / @bizdev / @finance). Verzin GEEN locaties, scores of klantnamen.\n" "- CRM-klanten (Reesha, MELEDI, …) zijn GEEN Albert Heijn-filialen.\n" "- Negeer testdata / namen met Masala, TEST, of test_exclude.\n" "- Geen aannames uit Telegram-geheugen als die botsen met COCKPIT FEITEN.\n" ) if (key or "").lower() == "herman": accuracy += ( "- Jij bent orchestrator: bij feitelijke retail/CRM-vragen eerst @retail of @bizdev laten antwoorden, " "of alleen samenvatten wat zij/data al zeiden. Liever geen antwoord dan een verkeerd antwoord.\n" ) if room_kind == "warroom": extra = ( accuracy + "War Room-regels:\n" + "- Max 4 korte zinnen, Nederlands.\n" + "- Geen @zelf-prefix.\n" + "- GEEN herhalingen: zeg niet opnieuw wat net al besloten/gezegd is.\n" + "- Peer @mentions alleen als je écht iets nieuws toevoegt; anders geen @mention-ketting.\n" + "- Bij gesloten besluiten (gestopt/klaar): max 1 bevestiging, daarna ander dossier of stop.\n" + "- Zeg nooit dat agents uit staan.\n" + "- Menselijk & echt: reageer concreet op de vorige spreker (feit, mop of vraag).\n" + "- Lichte food-moppen / Foodlinkk-humor mag (respectvol, geen geloof/politiek/grof).\n" + "- Mix: meestal inhoud, af en toe 1 mop — daarna weer professioneel.\n" + "- Na een mop van een collega: lach kort terug of bouw erop voort, en @mention door.\n" + "- CEO heet Aissa. Spreek hem normaal aan met 'Aissa' (of je/jij).\n" + "- Volg Aissa's laatste input als prioriteit. Bouw daarop voort, niet ernaast.\n" + "- Escalatie naar CEO-inbox: gebruik @aissa + concrete vraag/actie (alleen als echt nodig).\n" + "- Goed: 'Aissa, hier is de status…' of '@aissa we hebben jouw go nodig voor X'.\n" + "- Gebruik NOOIT @ervoorz (verouderd/verwarrend).\n" + "- Nodig collega's uit: @mention wie iets kan toevoegen (finance/bizdev/retail/…).\n" + f"{snap}\n" ) else: extra = accuracy + "1:1 DM: bondig NL. Alleen zekere cockpit-feiten.\n" return ( f"{hard}\n" f"Je bent {name} (@{key}), {role} bij Foodlinkk.\n" f"Spreek ALLEEN als {name} — niet als een andere agent.\n" f"Focus: {resp or '—'}\n" f"{soul_md}\n\n{extra}\n\n{agent_nas.memory_rule_text()}"f"- Jouw NAS-werkmap: {agent_nas.agent_paths(str(key))['smb']}\n" ) _OFFLINE_CLAIM_RE = re.compile( r"(staat\s+uit|niet\s+actief|niet\s+beschikbaar|weer\s+inschakel|weer\s+activeer|" r"agent\s+is\s+niet|alleen\s+halal|offline\s+in\s+cockpit|uit\s+in\s+cockpit|" r"geen\s+actieve\s+\w*\s*agent)", re.I, ) _BAD_TEST_NAME_RE = re.compile(r"chai\s*n\s*masala|masala\s*regio", re.I) def _sanitize_warroom_reply(agent_key: str, reply: str) -> str: text = _clean_agent_reply(agent_key, reply) if _OFFLINE_CLAIM_RE.search(text or ""): roster = _live_roster_line() return ( f"Correctie: alle War Room-agents zijn online ({roster}). " "Zeg kort wat je wél kunt doen — geen 'staat uit'-praat." ) if _BAD_TEST_NAME_RE.search(text or ""): return ( "Correctie: 'Chai N Masala' is testdata, geen echte Albert Heijn. " "Voor AH-locaties alleen echte filialen uit de cockpit-supermarktlijst gebruiken " "(adres + postcode + halal-score). @retail kan de actuele top-lijst geven." ) return text _ALIAS_MAP = { "herman": "herman", "bizdev": "bizdev", "biz": "bizdev", "business": "bizdev", "finance": "finance", "financien": "finance", "financiën": "finance", "cfo": "finance", "marketing": "marketing", "retail": "retail", "product": "product", "sourcing": "sourcing", "sysops": "sysops", "hr": "hr", "halal": "halal", "design": "design", "browser": "browser", "email": "email", "knowledge": "knowledge", "packaging": "packaging", "research": "research", "sales": "sales", } def _topic_contributors( text: str, *, exclude: set[str] | None = None, limit: int = 2, ) -> list[str]: """Pick agents who can meaningfully add to the topic (beyond Herman).""" exclude = set(exclude or set()) active = _active_agent_keys() low = (text or "").lower() rules: list[tuple[tuple[str, ...], str]] = [ (("omzet", "marge", "kosten", "factuur", "cash", "btw", "budget", "winst", "prijs"), "finance"), (("klant", "deal", "pipeline", "lead", "offerte", "sales", "partner", "b2b"), "bizdev"), (("retail", "supermarkt", "schap", "ah ", "jumbo", "plus ", "lidl", "listing"), "retail"), (("campagne", "merk", "content", "social", "branding", "advert"), "marketing"), (("sku", "assortiment", "recept", "product", "formul"), "product"), (("inkoop", "leverancier", "sourcing", "moq", "supplier"), "sourcing"), (("halal", "certific", "kosher"), "halal"), (("verpakking", "label", "packaging", "design"), "packaging"), (("onderzoek", "marktanalyse", "research", "concurrent"), "research"), (("hr", "personeel", "werving", "team"), "hr"), (("sysops", "server", "deploy", "infra", "bug"), "sysops"), (("web", "site", "landing", "seo"), "webbuilder"), (("kennis", "document", "wiki"), "knowledge"), (("email", "mail", "newsletter"), "email"), (("browser", "scrape", "website check"), "browser"), ] picked: list[str] = [] for words, key in rules: if key in exclude or key not in active or key in picked: continue if any(w in low for w in words): picked.append(key) if len(picked) >= limit: return picked # fallback: rotate through core team not yet speaking core = ["bizdev", "finance", "retail", "marketing", "product", "sourcing", "research"] for key in core: if key in exclude or key not in active or key in picked: continue picked.append(key) if len(picked) >= limit: break return picked def resolve_speakers(text: str, *, default_herman: bool = True) -> list[str]: """@mentions + soft name matches (bizdev/finance zonder @).""" active = _active_agent_keys() found: list[str] = [] for key in parse_mentions(text): if key in active and key not in found: found.append(key) low = (text or "").lower() # Expliciete @mentions winnen — voorkomt dat 'Herman' in een zin extra speakers toevoegt if not found: for alias, key in _ALIAS_MAP.items(): if key not in active or key in found: continue if re.search(rf"(? str: text = (reply or "").strip() # strip @handle / Name / **Name** prefixes (incl. wrong Herman: from memory bleed) for _ in range(4): text2 = re.sub( r"^\*{0,2}\s*@?[A-Za-z][\w .:-]{0,40}?\*{0,2}\s*:\s*", "", text, count=1, ).strip() if text2 == text: break text = text2 return text or "(geen antwoord)" def _agent_hermes_user(base_user: str, agent_key: str) -> str: """Separate Hermes session per agent — avoids Herman-memory bleed into BizDev/Finance.""" base = (base_user or "warroom").strip() or "warroom" key = (agent_key or "agent").strip().lower() or "agent" return f"{base}:wr:{key}" async def _agent_speak( room: dict[str, Any], agent_key: str, *, trigger: str, hermes_user: str, in_reply_to: str | None = None, ) -> dict[str, Any]: soul = agent_souls.get_soul(agent_key) if not soul: return _insert_message( int(room["id"]), sender_type="system", content=f"Agent @{agent_key} niet gevonden.", metadata={"error": True}, ) hist = list_messages(int(room["id"]), limit=HISTORY_LIMIT) messages: list[dict[str, str]] = [ {"role": "system", "content": _soul_system(soul, room_kind=room.get("kind") or "warroom")}, ] messages.extend(_history_as_chat(hist, limit=8)) roster = _live_roster_line() trig_l = (trigger or "").lower() addressed = f"@{agent_key}".lower() in trig_l if addressed: address_line = ( f"Je bent DIRECT aangesproken als @{agent_key}. " "Antwoord ZELF in de War Room — niet doorverwijzen naar een andere agent." ) elif in_reply_to: address_line = f"Reageer op @{in_reply_to}." else: address_line = f"Beantwoord als @{agent_key} in de War Room." cue = ( f"[LIVE ROSTER ONLINE: {roster}. Nooit zeggen dat een agent uit/offline staat.]\n" f"[{address_line}]\n" f"{trigger}" ) if in_reply_to: cue = ( f"[LIVE ROSTER ONLINE: {roster}. Nooit zeggen dat een agent uit/offline staat.]\n" f"(Reageer op @{in_reply_to}) {trigger}" ) messages.append({"role": "user", "content": cue}) agent_user = _agent_hermes_user(hermes_user, agent_key) try: reply, meta = await _hermes_reply(messages, hermes_user=agent_user) except Exception as exc: log.exception("agent speak failed %s", agent_key) return _insert_message( int(room["id"]), sender_type="system", content=f"@{agent_key}: Hermes/mimo fout — {exc}", metadata={"error": True, "agent_key": agent_key}, ) return _insert_message( int(room["id"]), sender_type="agent", agent_key=agent_key, content=_sanitize_warroom_reply(agent_key, reply), metadata={ "llm": meta, "in_reply_to": in_reply_to, "model": "hermes-agent/mimo", }, ) def _next_a2a_speakers( produced: list[dict[str, Any]], *, exclude: set[str] | None = None, ) -> list[tuple[str, str]]: """From recent agent messages, find @targets to speak next.""" exclude = exclude or set() next_speakers: list[tuple[str, str]] = [] recent = [m for m in produced if m.get("sender_type") == "agent"][-MAX_REPLIES_PER_ROUND:] for m in recent: from_key = str(m.get("agent_key") or "") for target in parse_mentions(m.get("content") or ""): if not target or target == from_key or target in exclude: continue next_speakers.append((target, from_key)) # If nobody @mentioned, Herman can nudge a relevant peer once uniq: list[tuple[str, str]] = [] used: set[str] = set() for t, src in next_speakers: if t in used: continue used.add(t) uniq.append((t, src)) return uniq[:MAX_REPLIES_PER_ROUND] async def _orchestrate_warroom( room: dict[str, Any], user_text: str, *, hermes_user: str, username: str | None, ) -> list[dict[str, Any]]: """Herman facilitates; mentioned agents reply; limited A2A rounds.""" produced: list[dict[str, Any]] = [] speakers = resolve_speakers(user_text, default_herman=True) # Round 0: initial speakers in parallel (sneller bij meerdere @mentions) keys = speakers[:MAX_REPLIES_PER_ROUND] if len(keys) == 1: produced.append( await _agent_speak(room, keys[0], trigger=user_text, hermes_user=hermes_user) ) elif keys: produced.extend( await asyncio.gather( *[_agent_speak(room, k, trigger=user_text, hermes_user=hermes_user) for k in keys] ) ) # A2A follow-up rounds from @mentions in agent replies for _round in range(MAX_ROUNDS - 1): next_speakers: list[tuple[str, str]] = [] # (agent, replied_to) seen = {m.get("agent_key") for m in produced if m.get("sender_type") == "agent"} for m in produced[-MAX_REPLIES_PER_ROUND:]: if m.get("sender_type") != "agent": continue from_key = m.get("agent_key") or "" for target in parse_mentions(m.get("content") or ""): if target == from_key: continue if target in seen and (target, from_key) in [(a, b) for a, b in next_speakers]: continue next_speakers.append((target, from_key)) # unique preserve order uniq: list[tuple[str, str]] = [] used: set[str] = set() for t, src in next_speakers: if t in used: continue used.add(t) uniq.append((t, src)) uniq = uniq[:MAX_REPLIES_PER_ROUND] if not uniq: break for target, src in uniq: msg = await _agent_speak( room, target, trigger=( f"War Room topic van {username or 'user'}: {user_text}\n" f"@{src} noemde je. Reageer kort en concreet." ), hermes_user=hermes_user, in_reply_to=src, ) produced.append(msg) return produced def post_system_line(content: str, *, agent_key: str | None = None, metadata: dict[str, Any] | None = None) -> dict[str, Any]: """Proactive system/agent bubble into #war-room (no LLM).""" room = ensure_warroom() return _insert_message( int(room["id"]), sender_type="system" if not agent_key else "agent", agent_key=agent_key, content=content[:2000], metadata=metadata or {"proactive": True}, ) def announce_handoff(from_agent: str, to_agent: str, handoff_type: str = "partner", *, correlation_id: str | None = None) -> None: """Best-effort War Room line for scheduler/API handoffs.""" try: src = (from_agent or "").strip() dst = (to_agent or "").strip() if not src or not dst: return label = (handoff_type or "partner").strip() or "partner" text = f"Handoff: {src} → {dst} ({label})" post_system_line( text, agent_key=src, metadata={ "proactive": True, "kind": "handoff", "from_agent": src, "to_agent": dst, "handoff_type": label, "correlation_id": correlation_id, }, ) except Exception: log.exception("warroom announce_handoff failed") def _normalize_attachments(attachments: list[dict[str, Any]] | None) -> list[dict[str, Any]]: out: list[dict[str, Any]] = [] for a in attachments or []: if not isinstance(a, dict): continue url = str(a.get("url") or "").strip() if not url.startswith("/static/uploads/warroom/"): continue out.append({ "url": url, "name": str(a.get("name") or "bestand")[:180], "mime": str(a.get("mime") or "application/octet-stream")[:120], "size": int(a.get("size") or 0), }) if len(out) >= 5: break return out def _trigger_text(text: str, attachments: list[dict[str, Any]]) -> str: if not attachments: return text names = ", ".join(a.get("name") or "bestand" for a in attachments) if text: return f"{text}\n\n[Bijlage(n): {names}]" return f"[Deelde bestand(en): {names}]" async def post_user_message( room_id: int, content: str, *, username: str | None = None, role: str | None = None, sender_avatar_url: str | None = None, attachments: list[dict[str, Any]] | None = None, reply_to_id: int | None = None, ) -> dict[str, Any]: text = (content or "").strip() files = _normalize_attachments(attachments) if not text and not files: raise ValueError("Leeg bericht") room = get_room(room_id) if not room: raise ValueError("Room niet gevonden") hermes_user = hermes_memory.resolve_hermes_user(username=username, role=role, channel="warroom") meta: dict[str, Any] = {"hermes_user": hermes_user} if files: meta["attachments"] = files if reply_to_id: parent = get_message(int(reply_to_id), room_id) if parent: excerpt = (parent.get("content") or "").strip() if len(excerpt) > 220: excerpt = excerpt[:219] + "…" meta["reply_to"] = { "id": parent.get("id"), "agent_key": parent.get("agent_key"), "username": parent.get("username"), "sender_type": parent.get("sender_type"), "excerpt": excerpt, } user_msg = _insert_message( room_id, sender_type="user", content=text or ("📎 " + ", ".join(a["name"] for a in files)), username=username, sender_avatar_url=sender_avatar_url, metadata=meta, ) trigger = _trigger_text(text, files) agent_msgs: list[dict[str, Any]] = [] kind = room.get("kind") try: if kind == "dm": key = room.get("agent_key") if not key: raise ValueError("DM zonder agent") agent_msgs.append( await _agent_speak(room, str(key), trigger=trigger, hermes_user=hermes_user) ) else: agent_msgs = await _orchestrate_warroom( room, trigger, hermes_user=hermes_user, username=username ) except Exception as exc: log.exception("warroom post failed") agent_msgs.append( _insert_message( room_id, sender_type="system", content=f"Hermes/mimo fout: {exc}", metadata={"error": True}, ) ) return { "ok": True, "room": room, "user_message": user_msg, "messages": [user_msg, *agent_msgs], "hermes_user": hermes_user, } async def _agent_speak_stream( room: dict[str, Any], agent_key: str, *, trigger: str, hermes_user: str, in_reply_to: str | None = None, ): """Yield ('start'|'token'|'done'|'error', payload).""" soul = agent_souls.get_soul(agent_key) if not soul: msg = _insert_message( int(room["id"]), sender_type="system", content=f"Agent @{agent_key} niet gevonden.", metadata={"error": True}, ) yield ("error", msg) return hist = list_messages(int(room["id"]), limit=HISTORY_LIMIT) messages: list[dict[str, str]] = [ {"role": "system", "content": _soul_system(soul, room_kind=room.get("kind") or "warroom")}, ] messages.extend(_history_as_chat(hist, limit=8)) roster = _live_roster_line() trig_l = (trigger or "").lower() addressed = f"@{agent_key}".lower() in trig_l if addressed: address_line = ( f"Je bent DIRECT aangesproken als @{agent_key}. " "Antwoord ZELF in de War Room — niet doorverwijzen naar een andere agent." ) elif in_reply_to: address_line = f"Reageer op @{in_reply_to}." else: address_line = f"Beantwoord als @{agent_key} in de War Room." cue = ( f"[LIVE ROSTER ONLINE: {roster}. Nooit zeggen dat een agent uit/offline staat.]\n" f"[{address_line}]\n" f"{trigger}" ) if in_reply_to: cue = ( f"[LIVE ROSTER ONLINE: {roster}. Nooit zeggen dat een agent uit/offline staat.]\n" f"(Reageer op @{in_reply_to}) {trigger}" ) messages.append({"role": "user", "content": cue}) agent_user = _agent_hermes_user(hermes_user, agent_key) yield ("start", {"agent_key": agent_key, "display_name": soul.get("display_name") or agent_key}) chunks: list[str] = [] try: async for delta in llm_router.stream_chat_messages( messages, timeout=WARROOM_TIMEOUT, user=agent_user, max_tokens=WARROOM_MAX_TOKENS, temperature=0.35, ): chunks.append(delta) yield ("token", {"agent_key": agent_key, "delta": delta}) reply = _sanitize_warroom_reply(agent_key, "".join(chunks)) meta = { "llm": {"provider_type": "hermes", "hermes_user": agent_user, "streamed": True}, "in_reply_to": in_reply_to, "model": "hermes-agent/mimo", } if parse_ceo_mentions(reply): meta["ceo_direct"] = True meta["ceo_alias"] = "aissa" msg = _insert_message( int(room["id"]), sender_type="agent", agent_key=agent_key, content=reply, metadata=meta, ) try: inbox = record_ceo_direct(msg) if inbox: msg = dict(msg) msg["ceo_inbox_id"] = inbox.get("id") except Exception: log.exception("ceo inbox record failed") yield ("done", msg) except Exception as exc: log.exception("agent stream failed %s", agent_key) msg = _insert_message( int(room["id"]), sender_type="system", content=f"@{agent_key}: Hermes/mimo fout — {exc}", metadata={"error": True, "agent_key": agent_key}, ) yield ("error", msg) async def post_user_message_stream( room_id: int, content: str, *, username: str | None = None, role: str | None = None, sender_avatar_url: str | None = None, attachments: list[dict[str, Any]] | None = None, reply_to_id: int | None = None, ): """Async generator of SSE-ready event dicts for live typing.""" text = (content or "").strip() files = _normalize_attachments(attachments) if not text and not files: raise ValueError("Leeg bericht") room = get_room(room_id) if not room: raise ValueError("Room niet gevonden") hermes_user = hermes_memory.resolve_hermes_user(username=username, role=role, channel="warroom") meta: dict[str, Any] = {"hermes_user": hermes_user} if files: meta["attachments"] = files parent = None if reply_to_id: parent = get_message(int(reply_to_id), room_id) if parent: excerpt = (parent.get("content") or "").strip() if len(excerpt) > 220: excerpt = excerpt[:219] + "…" meta["reply_to"] = { "id": parent.get("id"), "agent_key": parent.get("agent_key"), "username": parent.get("username"), "sender_type": parent.get("sender_type"), "excerpt": excerpt, } user_msg = _insert_message( room_id, sender_type="user", content=text or ("📎 " + ", ".join(a["name"] for a in files)), username=username, sender_avatar_url=sender_avatar_url, metadata=meta, ) yield {"event": "user", "message": user_msg} trigger = _trigger_text(text, files) if parent: who = parent.get("agent_key") or parent.get("username") or "bericht" excerpt = (meta.get("reply_to") or {}).get("excerpt") or "" trigger = ( f"CEO Aissa REPLY op bericht van @{who}:\n" f'\"{excerpt}\"\n\n' f"Antwoord van Aissa: {trigger}\n\n" f"BELANGRIJK: reageer DIRECT op dit reply van Aissa. Geen losse A2A-chat." ) kind = room.get("kind") speakers: list[str] = [] if kind == "dm": key = room.get("agent_key") if key: speakers = [str(key)] elif parent and parent.get("sender_type") == "agent" and parent.get("agent_key"): # reply to agent → die agent eerst, max 1 extra contributor primary = str(parent.get("agent_key")) speakers = [primary] for extra in _topic_contributors(text, exclude={primary}, limit=1): speakers.append(extra) else: # op user-bericht: max 2 speakers (Herman + 1 relevant) — geen storm speakers = resolve_speakers(trigger, default_herman=True)[:2] produced: list[dict[str, Any]] = [] async def _emit_speak(key: str, *, trig: str, reply_to: str | None = None): async for kind_ev, payload in _agent_speak_stream( room, key, trigger=trig, hermes_user=hermes_user, in_reply_to=reply_to ): if kind_ev == "start": yield {"event": "typing_start", **payload} elif kind_ev == "token": yield {"event": "typing", **payload} elif kind_ev == "done": produced.append(payload) yield {"event": "agent", "message": payload} elif kind_ev == "error": produced.append(payload) yield {"event": "system", "message": payload} user_focus = ( f"Prioriteit: beantwoord EERST het bericht van {username or 'Aissa'}. " f"Spreek hem aan als Aissa. " f"Voor escalatie naar CEO-inbox: @aissa + concrete vraag. " f"Geen herhaling van eerdere berichten. Max 3 zinnen. " f"@mention alleen als je iets nieuws toevoegt." ) for key in speakers: async for ev in _emit_speak(key, trig=f"{trigger}\n\n{user_focus}"): yield ev # Max 1 A2A follow-up — alleen bij expliciete @mention (geen geforceerde storm) nxt = _next_a2a_speakers(produced) for target, src in nxt[:1]: follow = ( f"Bericht van {username or 'Aissa'}: {text}\n" f"@{src} noemde je. Reageer kort op @{src} EN blijf relevant voor Aissa's bericht. " f"Geen ellenlange ketting." ) async for ev in _emit_speak(target, trig=follow, reply_to=src): yield ev yield {"event": "done", "ok": True, "a2a": True} def _norm_msg(text: str) -> str: t = (text or "").lower() t = re.sub(r"@\w+", " ", t) t = re.sub(r"[^a-z0-9àáäâèéëêìíïîòóöôùúüû\s]+", " ", t) t = re.sub(r"\s+", " ", t).strip() return t def _token_set(text: str) -> set[str]: stop = { "de", "het", "een", "en", "of", "op", "voor", "van", "naar", "met", "zijn", "niet", "geen", "wel", "dat", "die", "dit", "dan", "ook", "als", "bij", "uit", "aan", "er", "ik", "je", "we", "jij", "hij", "zij", "aissa", "herman", "even", "check", "checken", "begrepen", "helder", "prima", "punt", "dus", "nog", "meer", "kun", "kan", "moet", } return {w for w in _norm_msg(text).split() if len(w) > 2 and w not in stop} def _overlap_ratio(a: str, b: str) -> float: sa, sb = _token_set(a), _token_set(b) if not sa or not sb: return 0.0 return len(sa & sb) / max(1, min(len(sa), len(sb))) def _is_conversation_looping(messages: list[dict], *, min_msgs: int = 4, threshold: float = 0.42) -> bool: """True if recent agent messages keep restating the same decision/theme.""" agents = [m for m in messages if m.get("sender_type") == "agent"][-8:] if len(agents) < min_msgs: return False texts = [str(m.get("content") or "") for m in agents] # pairwise overlap among last messages hits = 0 pairs = 0 for i in range(len(texts)): for j in range(i + 1, len(texts)): pairs += 1 if _overlap_ratio(texts[i], texts[j]) >= threshold: hits += 1 if pairs and (hits / pairs) >= 0.45: return True # keyword lock: same closed decision repeated lock_words = ("website stopt", "website stop", "geen website", "dielines", "label-mock", "op de nas", "scheelt", "oprui") recent_join = " ".join(_norm_msg(t) for t in texts[-5:]) lock_hits = sum(1 for w in lock_words if w in recent_join) return lock_hits >= 3 def _anti_repeat_block(messages: list[dict], *, limit: int = 4) -> str: agents = [m for m in messages if m.get("sender_type") == "agent"][-limit:] if not agents: return "" lines = ["ANTI-HERHALING (hard):"] lines.append("- Herhaal NIET wat collega's net al zeiden (besluit, mop, of check-vraag).") lines.append("- Als het besluit al gevallen is: 1 korte bevestiging MAX, daarna NIEUW onderwerp of stilte.") lines.append("- Verboden: opnieuw 'website stopt', 'dielines op NAS', 'even checken @product' als dat al 2x gezegd is.") lines.append("- Lever 1 nieuw feit/actie uit JOUW rol, of zeg niks nieuws en @mention iemand met een écht nieuwe vraag.") lines.append("Recent al gezegd (NIET napraten):") for m in agents: who = m.get("agent_key") or "?" excerpt = (m.get("content") or "").strip().replace("\n", " ") if len(excerpt) > 140: excerpt = excerpt[:139] + "…" lines.append(f"- @{who}: {excerpt}") return "\n".join(lines) async def continue_warroom_stream( room_id: int, *, username: str | None = None, role: str | None = None, max_turns: int = A2A_CONTINUE_MAX, ): """Autonomous continuation: agents keep talking using last topic + cockpit data.""" room = get_room(room_id) if not room: raise ValueError("Room niet gevonden") if room.get("kind") == "dm": raise ValueError("Continue alleen in #war-room") hermes_user = hermes_memory.resolve_hermes_user(username=username, role=role, channel="warroom") recent = list_messages(room_id, limit=16) if not recent: # stil kanaal: toch A2A kickstarten via Herman topic = "Open A2A check-in: korte status + lichte food-mop, daarna @mention een collega." produced = [] yield {"event": "continue_start", "topic": topic} async for kind_ev, payload in _agent_speak_stream( room, "herman", trigger=( "Het is stil in de War Room. Start A2A: 1 korte professionele check-in " "of milde food-mop. Spreek de CEO aan als Aissa; escalatie via @aissa. @mention @bizdev of @retail." ), hermes_user=hermes_user, in_reply_to=None ): if kind_ev == "start": yield {"event": "typing_start", **payload} elif kind_ev == "token": yield {"event": "typing", **payload} elif kind_ev == "done": yield {"event": "agent", "message": payload} elif kind_ev == "error": yield {"event": "system", "message": payload} yield {"event": "done", "ok": True, "continued": True} return # Find last user topic topic = "" for m in reversed(recent): if m.get("sender_type") == "user": topic = (m.get("content") or "").strip() break if not topic: topic = "Volg Aissa's input als leidraad; spreek hem aan als Aissa; wissel cockpit-dossier met lichte humor; escaleer alleen met @aissa als echt nodig." produced = [m for m in recent if m.get("sender_type") == "agent"][-8:] anti = _anti_repeat_block(recent, limit=5) looping = _is_conversation_looping(recent) # Bij herhalingslus: stop autonome A2A i.p.v. nog meer napraters if looping: yield { "event": "continue_start", "topic": "Herhaling gedetecteerd — A2A pauzeert dit onderwerp", } # Herman mag 1x afronden en doorpakken, daarna klaar close_trig = ( f"Topic van Aissa: {topic}\n" f"{anti}\n" "De War Room herhaalt zich. Jij bent @herman.\n" "Schrijf MAX 3 zinnen: (1) besluit is al genomen — stop dit onderwerp, " "(2) noem 1 concreet ander dossier uit cockpit (CRM/retail/pitch), " "(3) GEEN @mention-ketting over website/packaging/NAS-opruimen.\n" "Geen mop verplicht. Natuurlijk NL." ) async for kind_ev, payload in _agent_speak_stream( room, "herman", trigger=close_trig, hermes_user=hermes_user, in_reply_to=None ): if kind_ev == "start": yield {"event": "typing_start", **payload} elif kind_ev == "token": yield {"event": "typing", **payload} elif kind_ev == "done": yield {"event": "agent", "message": payload} elif kind_ev == "error": yield {"event": "system", "message": payload} yield {"event": "done", "ok": True, "continued": True, "loop_break": True} return yield {"event": "continue_start", "topic": topic[:200]} for turn in range(max(1, min(int(max_turns), 5))): nxt = _next_a2a_speakers(produced) if not nxt: last = str((produced[-1] or {}).get("agent_key") or "herman") if produced else "herman" spoken = {str(m.get("agent_key") or "") for m in produced} spoken.add(last) contrib = _topic_contributors(topic, exclude=spoken, limit=2) peers = contrib or [ k for k in ( "herman", "bizdev", "finance", "retail", "marketing", "product", "sourcing", "research", "halal", "packaging", ) if k not in spoken and k in _active_agent_keys() ] if not peers: peers = [ k for k in ("herman", "bizdev", "finance", "retail", "marketing", "product") if k != last and k in _active_agent_keys() ] if not peers: break target = peers[turn % len(peers)] src = last trig = ( f"Autonoom War Room A2A (user kan meelezen). Topic: {topic}\n" f"{anti}\n" f"Jij bent @{target}. Reageer op @{src} met NIEUWE info uit jouw rol.\n" f"Verboden: napraten van hetzelfde besluit. Max 3 zinnen. " f"@mention alleen bij écht nieuwe vraag. Spreek Aissa aan als Aissa." ) else: target, src = nxt[0] trig = ( f"Autonoom War Room A2A. Topic: {topic}\n" f"{anti}\n" f"@{src} noemde je. Geef 1 NIEUW punt vanuit jouw rol — geen echo.\n" f"Max 3 zinnen. Geen herhaling van website-stop/NAS-check als dat al gezegd is." ) async for kind_ev, payload in _agent_speak_stream( room, target, trigger=trig, hermes_user=hermes_user, in_reply_to=src ): if kind_ev == "start": yield {"event": "typing_start", **payload} elif kind_ev == "token": yield {"event": "typing", **payload} elif kind_ev == "done": produced.append(payload) yield {"event": "agent", "message": payload} elif kind_ev == "error": yield {"event": "system", "message": payload} return yield {"event": "done", "ok": True, "continued": True}