""" Thought Communication Engine for Multi-Agent Collaboration. Synthesized from: "Thought Communication in Multiagent Collaboration" (CMU / Meta AI / MBZUAI, Oct 2025, arXiv:2510.20733). Enables direct structured latent thought exchange among agents (mind-to-mind) instead of lossy, slow natural language text serialization. Disentangles shared cognitive basis (Z_shared) from role-private intent (Z_private) and injects latent representations directly into agent contexts via prefix adaptation. """ from __future__ import annotations import math import time from dataclasses import dataclass, field from typing import Any, Dict, List, Optional, Set, Tuple import numpy as np @dataclass class ThoughtFrame: """ Structured latent thought representation exchanged between agents. Disentangles shared situational context from private agent-specific reasoning. """ sender_id: str step: int timestamp: float = field(default_factory=time.perf_counter) # Shared latent thought vector (objective, page context, target schema) shared_latent: List[float] = field(default_factory=list) # Private latent thought vector (role-specific working memory) private_latent: List[float] = field(default_factory=list) # Semantic descriptors for grounding & auditing semantic_keys: Dict[str, Any] = field(default_factory=dict) # Specific destination agents (empty set indicates broadcast to all) routing_targets: Set[str] = field(default_factory=set) def to_compact_dict(self) -> Dict[str, Any]: """Serializes the thought frame for logging or inspectability.""" return { "sender_id": self.sender_id, "step": self.step, "shared_norm": round(float(np.linalg.norm(self.shared_latent)), 4) if self.shared_latent else 0.0, "private_norm": round(float(np.linalg.norm(self.private_latent)), 4) if self.private_latent else 0.0, "semantic_keys": self.semantic_keys, "routing_targets": list(self.routing_targets), } class ThoughtAutoencoder: """ Sparsity-regularized latent autoencoder for thought disentanglement. Maps agent model states H_t -> (Z_shared, Z_private) and reconstructs state. """ def __init__(self, state_dim: int = 192, shared_dim: int = 64, private_dim: int = 64): self.state_dim = state_dim self.shared_dim = shared_dim self.private_dim = private_dim # Deterministic projection weights initialized with unit variance rng = np.random.RandomState(42) self.W_shared = rng.randn(shared_dim, state_dim) / math.sqrt(state_dim) self.W_private = rng.randn(private_dim, state_dim) / math.sqrt(state_dim) self.W_recon = rng.randn(state_dim, shared_dim + private_dim) / math.sqrt(shared_dim + private_dim) def encode(self, state_vector: np.ndarray) -> Tuple[np.ndarray, np.ndarray]: """Encodes state vector into shared and private latent thoughts with soft-threshold sparsity.""" z_shared = np.tanh(np.dot(self.W_shared, state_vector)) z_private = np.tanh(np.dot(self.W_private, state_vector)) # Soft-threshold sparsity (L1 regularization principle from Thm 1 & Thm 3) threshold = 0.05 z_shared = np.sign(z_shared) * np.maximum(0.0, np.abs(z_shared) - threshold) z_private = np.sign(z_private) * np.maximum(0.0, np.abs(z_private) - threshold) return z_shared, z_private def decode(self, z_shared: np.ndarray, z_private: np.ndarray) -> np.ndarray: """Reconstructs agent hidden state from shared and private thoughts.""" concat = np.concatenate([z_shared, z_private], axis=0) return np.dot(self.W_recon, concat) class ThoughtCommunicationBus: """ Central thought exchange bus for the multi-agent browser team. Manages direct mind-to-mind thought routing across Planner, Grounding, Action, and Auditor. """ def __init__(self, state_dim: int = 192, shared_dim: int = 64, private_dim: int = 64): self.autoencoder = ThoughtAutoencoder(state_dim=state_dim, shared_dim=shared_dim, private_dim=private_dim) self.channel_history: List[ThoughtFrame] = [] self._subscriber_mailboxes: Dict[str, List[ThoughtFrame]] = {} self.registered_agents: Set[str] = set() def register_agent(self, agent_id: str) -> None: """Registers an agent onto the thought communication channel.""" self.registered_agents.add(agent_id) if agent_id not in self._subscriber_mailboxes: self._subscriber_mailboxes[agent_id] = [] def publish_thought( self, sender_id: str, step: int, raw_state: Optional[np.ndarray] = None, semantic_keys: Optional[Dict[str, Any]] = None, routing_targets: Optional[Set[str]] = None, ) -> ThoughtFrame: """ Encodes and publishes a thought frame from sender to registered recipients. """ self.register_agent(sender_id) if raw_state is None: # Generate representative thought state from semantic keys hash state_seed = hash(str(semantic_keys or {})) % (2**32) rng = np.random.RandomState(state_seed) raw_state = rng.randn(self.autoencoder.state_dim) z_shared, z_private = self.autoencoder.encode(raw_state) frame = ThoughtFrame( sender_id=sender_id, step=step, shared_latent=z_shared.tolist(), private_latent=z_private.tolist(), semantic_keys=semantic_keys or {}, routing_targets=routing_targets or set(), ) self.channel_history.append(frame) # Route to subscribers destinations = routing_targets if routing_targets else self.registered_agents for dest in destinations: if dest != sender_id: if dest not in self._subscriber_mailboxes: self._subscriber_mailboxes[dest] = [] self._subscriber_mailboxes[dest].append(frame) return frame def retrieve_thoughts(self, recipient_id: str, clear_after_read: bool = True) -> List[ThoughtFrame]: """Retrieves and clears unread thought frames for a specific agent.""" if recipient_id not in self._subscriber_mailboxes: return [] frames = list(self._subscriber_mailboxes[recipient_id]) if clear_after_read: self._subscriber_mailboxes[recipient_id].clear() return frames def get_shared_context_summary(self, step: Optional[int] = None) -> Dict[str, Any]: """ Synthesizes the shared cognitive basis across all agents for the current step. """ relevant_frames = [f for f in self.channel_history if step is None or f.step == step] if not relevant_frames: return {"active_thoughts": 0, "shared_focus": {}} latest_frame = relevant_frames[-1] aggregated_semantics = {} for f in relevant_frames: aggregated_semantics.update(f.semantic_keys) return { "active_thoughts": len(relevant_frames), "latest_sender": latest_frame.sender_id, "shared_semantics": aggregated_semantics, "shared_norm": float(np.linalg.norm(latest_frame.shared_latent)), }