Add War Room voice meeting, fix chat overlay, and ship 5 cyan UI themes.
Voice Room lets humans talk with agents via Whisper/TTS; chat typing no longer blocked by a full-screen overlay, attachments/avatars are larger, and Settings exposes five vivid mid-dark cyan/blue themes for the whole cockpit.
This commit is contained in:
@@ -15,11 +15,105 @@ 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 = 2
|
||||
MAX_ROUNDS = 4
|
||||
MAX_REPLIES_PER_ROUND = 3
|
||||
HISTORY_LIMIT = 16
|
||||
WARROOM_MAX_TOKENS = 512
|
||||
WARROOM_TIMEOUT = 55.0
|
||||
A2A_CONTINUE_MAX = 3
|
||||
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:
|
||||
@@ -87,14 +181,15 @@ def list_messages(room_id: int, *, limit: int = 100, after_id: int | None = None
|
||||
if after_id:
|
||||
rows = fetch_all(
|
||||
"""SELECT * FROM warroom_messages
|
||||
WHERE room_id = %s AND id > %s
|
||||
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, limit),
|
||||
)
|
||||
else:
|
||||
rows = fetch_all(
|
||||
"""SELECT * FROM (
|
||||
SELECT * FROM warroom_messages WHERE room_id = %s
|
||||
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, limit),
|
||||
@@ -102,6 +197,34 @@ def list_messages(room_id: int, *, limit: int = 100, after_id: int | None = None
|
||||
return [_iso(r) for r in rows] # type: ignore
|
||||
|
||||
|
||||
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 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,
|
||||
*,
|
||||
@@ -195,24 +318,174 @@ async def _hermes_reply(
|
||||
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 "")[:900]
|
||||
resp = (soul.get("responsibilities") or "")[:400]
|
||||
soul_md = (soul.get("soul_md") or "")[:700]
|
||||
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"
|
||||
)
|
||||
extra = (
|
||||
"War Room: antwoord in max 4 korte zinnen (NL). Alleen @mention als écht nodig."
|
||||
accuracy
|
||||
+ "War Room-regels:\n"
|
||||
"- Max 4 korte zinnen, Nederlands.\n"
|
||||
"- Geen @zelf-prefix / geen herhalingen.\n"
|
||||
"- Peer @mentions oké; zeg nooit dat agents uit staan.\n"
|
||||
f"{snap}\n"
|
||||
if room_kind == "warroom"
|
||||
else "1:1 DM: bondig in het Nederlands (max 5 zinnen)."
|
||||
else 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}"
|
||||
)
|
||||
|
||||
|
||||
_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 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"(?<![a-z0-9_]){re.escape(alias)}(?![a-z0-9_])", low):
|
||||
found.append(key)
|
||||
for key in sorted(active):
|
||||
if key in found:
|
||||
continue
|
||||
if re.search(rf"(?<![a-z0-9_]){re.escape(key)}(?![a-z0-9_])", low):
|
||||
found.append(key)
|
||||
explicit = bool(parse_mentions(text))
|
||||
if not found and default_herman and "herman" in active:
|
||||
found = ["herman"]
|
||||
elif (
|
||||
found
|
||||
and not explicit
|
||||
and default_herman
|
||||
and "herman" in active
|
||||
and "herman" not in found
|
||||
):
|
||||
# Alleen zonder @mentions: "meeting met jou en bizdev" → Herman erbij
|
||||
if re.search(r"\b(jou|jullie|samen|meeting|overleg|bellen|call|afspraak)\b", low):
|
||||
found.insert(0, "herman")
|
||||
return found[:MAX_REPLIES_PER_ROUND]
|
||||
|
||||
|
||||
|
||||
def _clean_agent_reply(agent_key: str, reply: str) -> 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,
|
||||
@@ -234,12 +507,32 @@ async def _agent_speak(
|
||||
{"role": "system", "content": _soul_system(soul, room_kind=room.get("kind") or "warroom")},
|
||||
]
|
||||
messages.extend(_history_as_chat(hist, limit=8))
|
||||
cue = trigger
|
||||
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"(Reageer op @{in_reply_to}) {trigger}"
|
||||
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=hermes_user)
|
||||
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(
|
||||
@@ -252,7 +545,7 @@ async def _agent_speak(
|
||||
int(room["id"]),
|
||||
sender_type="agent",
|
||||
agent_key=agent_key,
|
||||
content=reply.strip() or "(geen antwoord)",
|
||||
content=_sanitize_warroom_reply(agent_key, reply),
|
||||
metadata={
|
||||
"llm": meta,
|
||||
"in_reply_to": in_reply_to,
|
||||
@@ -261,6 +554,33 @@ async def _agent_speak(
|
||||
)
|
||||
|
||||
|
||||
|
||||
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,
|
||||
@@ -270,11 +590,7 @@ async def _orchestrate_warroom(
|
||||
) -> list[dict[str, Any]]:
|
||||
"""Herman facilitates; mentioned agents reply; limited A2A rounds."""
|
||||
produced: list[dict[str, Any]] = []
|
||||
mentions = parse_mentions(user_text)
|
||||
speakers = mentions[:]
|
||||
if not speakers:
|
||||
# Herman opens, then may suggest others — still let Herman speak first
|
||||
speakers = ["herman"]
|
||||
speakers = resolve_speakers(user_text, default_herman=True)
|
||||
|
||||
# Round 0: initial speakers in parallel (sneller bij meerdere @mentions)
|
||||
keys = speakers[:MAX_REPLIES_PER_ROUND]
|
||||
@@ -487,10 +803,30 @@ async def _agent_speak_stream(
|
||||
{"role": "system", "content": _soul_system(soul, room_kind=room.get("kind") or "warroom")},
|
||||
]
|
||||
messages.extend(_history_as_chat(hist, limit=8))
|
||||
cue = trigger
|
||||
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"(Reageer op @{in_reply_to}) {trigger}"
|
||||
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] = []
|
||||
@@ -498,20 +834,20 @@ async def _agent_speak_stream(
|
||||
async for delta in llm_router.stream_chat_messages(
|
||||
messages,
|
||||
timeout=WARROOM_TIMEOUT,
|
||||
user=hermes_user,
|
||||
user=agent_user,
|
||||
max_tokens=WARROOM_MAX_TOKENS,
|
||||
temperature=0.35,
|
||||
):
|
||||
chunks.append(delta)
|
||||
yield ("token", {"agent_key": agent_key, "delta": delta})
|
||||
reply = "".join(chunks).strip() or "(geen antwoord)"
|
||||
reply = _sanitize_warroom_reply(agent_key, "".join(chunks))
|
||||
msg = _insert_message(
|
||||
int(room["id"]),
|
||||
sender_type="agent",
|
||||
agent_key=agent_key,
|
||||
content=reply,
|
||||
metadata={
|
||||
"llm": {"provider_type": "hermes", "hermes_user": hermes_user, "streamed": True},
|
||||
"llm": {"provider_type": "hermes", "hermes_user": agent_user, "streamed": True},
|
||||
"in_reply_to": in_reply_to,
|
||||
"model": "hermes-agent/mimo",
|
||||
},
|
||||
@@ -568,22 +904,121 @@ async def post_user_message_stream(
|
||||
if key:
|
||||
speakers = [str(key)]
|
||||
else:
|
||||
speakers = parse_mentions(trigger) or ["herman"]
|
||||
speakers = speakers[:MAX_REPLIES_PER_ROUND]
|
||||
speakers = resolve_speakers(trigger, default_herman=True)
|
||||
|
||||
for key in speakers:
|
||||
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=trigger, hermes_user=hermes_user
|
||||
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}
|
||||
|
||||
for key in speakers:
|
||||
async for ev in _emit_speak(key, trig=trigger):
|
||||
yield ev
|
||||
|
||||
# A2A rounds: agents @mention each other and answer with cockpit data
|
||||
for _round in range(MAX_ROUNDS - 1):
|
||||
nxt = _next_a2a_speakers(produced)
|
||||
if not nxt:
|
||||
# Force a light peer handoff so the room stays collaborative
|
||||
last_agents = [m.get("agent_key") for m in produced if m.get("sender_type") == "agent"]
|
||||
last = str(last_agents[-1]) if last_agents else (speakers[0] if speakers else "herman")
|
||||
peers = [k for k in ("bizdev", "finance", "retail", "marketing", "herman") if k != last and k in _active_agent_keys()]
|
||||
if not peers:
|
||||
break
|
||||
target = peers[_round % len(peers)]
|
||||
nudge = (
|
||||
f"War Room topic van {username or 'user'}: {trigger}\n"
|
||||
f"@{last} heeft gesproken. Jij bent @{target}. Stel 1 korte vraag aan een peer OF geef 1 cockpit-feit + @mention wie moet antwoorden."
|
||||
)
|
||||
async for ev in _emit_speak(target, trig=nudge, reply_to=last):
|
||||
yield ev
|
||||
continue
|
||||
for target, src in nxt:
|
||||
follow = (
|
||||
f"War Room topic van {username or 'user'}: {trigger}\n"
|
||||
f"@{src} noemde je. Antwoord met cockpit-data en @mention zo nodig een andere agent."
|
||||
)
|
||||
async for ev in _emit_speak(target, trig=follow, reply_to=src):
|
||||
yield ev
|
||||
|
||||
yield {"event": "done", "ok": True, "a2a": True}
|
||||
|
||||
|
||||
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:
|
||||
raise ValueError("Geen gesprek om te vervolgen")
|
||||
|
||||
# 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 = "Blijf het Foodlinkk cockpit-dossier bespreken (pipeline, deals, retail, finance)."
|
||||
|
||||
produced = [m for m in recent if m.get("sender_type") == "agent"][-6:]
|
||||
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"
|
||||
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-overleg (user kan meelezen). Topic: {topic}\n"
|
||||
f"Jij bent @{target}. Reageer op @{src} met 1 cockpit-feit en @mention 1 collega met een gerichte vraag."
|
||||
)
|
||||
else:
|
||||
target, src = nxt[0]
|
||||
trig = (
|
||||
f"Autonoom War Room-overleg. Topic: {topic}\n"
|
||||
f"@{src} vroeg jouw input. Antwoord met cockpit-data en @mention desnoods een andere agent."
|
||||
)
|
||||
|
||||
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
|
||||
|
||||
# light A2A: one follow-up round max from last agent reply mentions
|
||||
# (non-streamed for simplicity / speed)
|
||||
yield {"event": "done", "ok": True}
|
||||
yield {"event": "done", "ok": True, "continued": True}
|
||||
|
||||
@@ -126,10 +126,14 @@ def list_souls() -> list[dict[str, Any]]:
|
||||
soul["current_status"] = st.get("current_status")
|
||||
# Working = recent activity (not just DB is_active flag)
|
||||
soul["working"] = events_6h > 0
|
||||
soul["activity"] = (
|
||||
"busy" if (str(soul.get("current_status") or "").lower() in ("needs_approval", "pending") or events_6h >= 5)
|
||||
else ("active" if events_6h > 0 else ("idle" if event_count > 0 else "offline"))
|
||||
)
|
||||
# is_active souls are never "offline" in the registry (War Room / mesh ready)
|
||||
if not soul.get("is_active", True):
|
||||
soul["activity"] = "offline"
|
||||
else:
|
||||
soul["activity"] = (
|
||||
"busy" if (str(soul.get("current_status") or "").lower() in ("needs_approval", "pending") or events_6h >= 5)
|
||||
else ("active" if events_6h > 0 else "idle")
|
||||
)
|
||||
out.append(soul)
|
||||
return out
|
||||
|
||||
|
||||
@@ -243,10 +243,13 @@ def _collect_retail_intel(data: dict[str, Any]) -> None:
|
||||
|
||||
try:
|
||||
data["top_opportunities"] = fetch_all(
|
||||
"""SELECT s.name, s.chain, s.city, ros.halal_opportunity_score
|
||||
"""SELECT s.name, s.chain, s.city, s.address, s.postcode, ros.halal_opportunity_score
|
||||
FROM retail_opportunity_scores ros
|
||||
JOIN supermarkets s ON s.id = ros.supermarket_id
|
||||
ORDER BY ros.halal_opportunity_score DESC LIMIT 5"""
|
||||
WHERE 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 LIMIT 8"""
|
||||
)
|
||||
except Exception:
|
||||
data["top_opportunities"] = []
|
||||
|
||||
@@ -0,0 +1,180 @@
|
||||
"""Daily Herman War Room digest → Telegram (Aïssa + Mo)."""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any
|
||||
|
||||
import httpx
|
||||
|
||||
from app.config import settings
|
||||
from app.db import execute, fetch_all, fetch_one
|
||||
from app.services import llm_router
|
||||
|
||||
log = logging.getLogger("cockpit.warroom_digest")
|
||||
|
||||
|
||||
def _recipients() -> list[tuple[str, str]]:
|
||||
"""(chat_id, display_name)"""
|
||||
return [
|
||||
(str(getattr(settings, "HERMES_CEO_CHAT_ID", None) or "8859782446"), "Aïssa"),
|
||||
(str(getattr(settings, "HERMES_CTO_CHAT_ID", None) or "789036463"), "Mo"),
|
||||
]
|
||||
|
||||
|
||||
def collect_warroom_day(*, hours: int = 24, limit: int = 120) -> list[dict[str, Any]]:
|
||||
hours = max(1, min(int(hours), 72))
|
||||
rows = fetch_all(
|
||||
"""
|
||||
SELECT m.id, m.sender_type, m.agent_key, m.username, m.content, m.created_at, r.kind, r.title
|
||||
FROM warroom_messages m
|
||||
JOIN warroom_rooms r ON r.id = m.room_id
|
||||
WHERE m.created_at >= NOW() - make_interval(hours => %s)
|
||||
AND r.kind = 'warroom'
|
||||
ORDER BY m.created_at ASC
|
||||
LIMIT %s
|
||||
""",
|
||||
(hours, limit),
|
||||
)
|
||||
out = []
|
||||
for r in rows or []:
|
||||
item = dict(r)
|
||||
if item.get("created_at") and hasattr(item["created_at"], "isoformat"):
|
||||
item["created_at"] = item["created_at"].isoformat()
|
||||
out.append(item)
|
||||
return out
|
||||
|
||||
|
||||
def _transcript(messages: list[dict[str, Any]], *, max_chars: int = 6000) -> str:
|
||||
lines: list[str] = []
|
||||
for m in messages:
|
||||
st = m.get("sender_type")
|
||||
if st == "user":
|
||||
who = m.get("username") or "user"
|
||||
elif st == "agent":
|
||||
who = m.get("agent_key") or "agent"
|
||||
else:
|
||||
who = "system"
|
||||
content = (m.get("content") or "").strip().replace("\n", " ")
|
||||
if not content:
|
||||
continue
|
||||
lines.append(f"- {who}: {content[:280]}")
|
||||
text = "\n".join(lines)
|
||||
return text[-max_chars:]
|
||||
|
||||
|
||||
async def generate_digest_text(messages: list[dict[str, Any]]) -> str:
|
||||
day = datetime.now(timezone.utc).strftime("%Y-%m-%d")
|
||||
if not messages:
|
||||
return (
|
||||
f"📋 War Room dagrapport {day}\n\n"
|
||||
"Geen activiteit in #war-room vandaag. Agents stonden klaar; geen overleg gelogd."
|
||||
)
|
||||
transcript = _transcript(messages)
|
||||
system = (
|
||||
"Je bent Herman, Co-CEO van Foodlinkk. Schrijf een beknopte Nederlandse "
|
||||
"Telegram-samenvatting van het War Room-overleg voor Aïssa (CEO) en Mo (CTO).\n"
|
||||
"Structuur:\n"
|
||||
f"1) Kop: 📋 War Room dagrapport {day}\n"
|
||||
"2) 3–6 bullets: wat agents bespraken (pipeline, deals, acties, blockers)\n"
|
||||
"3) Sectie 'Acties / open punten' met concrete next steps\n"
|
||||
"4) 1 zin afsluiting\n"
|
||||
"Max ~350 woorden. Geen @zelf-prefix. Geen 'agents staan uit'."
|
||||
)
|
||||
user = f"War Room transcript (nieuwste onderaan):\n{transcript}"
|
||||
reply, _meta = await llm_router.chat_messages(
|
||||
[{"role": "system", "content": system}, {"role": "user", "content": user}],
|
||||
timeout=90.0,
|
||||
user=str(getattr(settings, "HERMES_CEO_CHAT_ID", None) or "8859782446"),
|
||||
max_tokens=900,
|
||||
temperature=0.3,
|
||||
)
|
||||
text = (reply or "").strip()
|
||||
if not text.startswith("📋"):
|
||||
text = f"📋 War Room dagrapport {day}\n\n{text}"
|
||||
return text[:3500]
|
||||
|
||||
|
||||
async def _deliver_telegram(chat_id: str, name: str, digest: str) -> dict[str, Any]:
|
||||
"""Ask Hermes agent to deliver the digest into that Telegram chat."""
|
||||
system = (
|
||||
f"Je bent Herman. Bezorg het volgende War Room-dagrapport NU via Telegram "
|
||||
f"aan {name} (chat_id={chat_id}). Stuur de tekst zo letterlijk mogelijk. "
|
||||
"Bevestig kort of verzonden."
|
||||
)
|
||||
try:
|
||||
reply, meta = await llm_router.chat_messages(
|
||||
[
|
||||
{"role": "system", "content": system},
|
||||
{"role": "user", "content": digest},
|
||||
],
|
||||
timeout=120.0,
|
||||
user=str(chat_id),
|
||||
max_tokens=200,
|
||||
temperature=0.1,
|
||||
)
|
||||
return {
|
||||
"ok": True,
|
||||
"chat_id": chat_id,
|
||||
"name": name,
|
||||
"confirm": (reply or "")[:240],
|
||||
"provider": (meta or {}).get("provider_type"),
|
||||
}
|
||||
except Exception as exc:
|
||||
log.exception("warroom digest deliver failed %s", chat_id)
|
||||
return {"ok": False, "chat_id": chat_id, "name": name, "error": str(exc)}
|
||||
|
||||
|
||||
def _persist(digest: str, message_count: int, deliveries: list[dict[str, Any]]) -> None:
|
||||
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)
|
||||
""",
|
||||
(
|
||||
"herman",
|
||||
"warroom_digest",
|
||||
"warroom_daily_digest",
|
||||
"War Room dagrapport verstuurd",
|
||||
digest[:4000],
|
||||
"completed",
|
||||
"telegram",
|
||||
__import__("json").dumps(
|
||||
{
|
||||
"message_count": message_count,
|
||||
"deliveries": deliveries,
|
||||
}
|
||||
),
|
||||
),
|
||||
)
|
||||
except Exception:
|
||||
log.exception("persist digest event failed")
|
||||
|
||||
try:
|
||||
from app.services import agent_chat
|
||||
|
||||
agent_chat.post_system_line(
|
||||
"📤 Dagrapport verstuurd naar Aïssa + Mo via Telegram.\n\n" + digest[:1500],
|
||||
agent_key="herman",
|
||||
metadata={"kind": "warroom_digest", "deliveries": deliveries},
|
||||
)
|
||||
except Exception:
|
||||
log.exception("post digest to warroom failed")
|
||||
|
||||
|
||||
async def run_daily_digest(*, hours: int = 24, deliver: bool = True) -> dict[str, Any]:
|
||||
messages = collect_warroom_day(hours=hours)
|
||||
digest = await generate_digest_text(messages)
|
||||
deliveries: list[dict[str, Any]] = []
|
||||
if deliver:
|
||||
for chat_id, name in _recipients():
|
||||
deliveries.append(await _deliver_telegram(chat_id, name, digest))
|
||||
_persist(digest, len(messages), deliveries)
|
||||
return {
|
||||
"ok": True,
|
||||
"message_count": len(messages),
|
||||
"digest": digest,
|
||||
"deliveries": deliveries,
|
||||
"generated_at": datetime.now(timezone.utc).isoformat(),
|
||||
}
|
||||
Reference in New Issue
Block a user