canopy-258m-r3 / miniswardbower /agents /thought_communication.py
psikosen's picture
Update to v6: Tri-Engine Swarm Coordinator, URL Invariance, Thought Bus & State Verification
ed79a7f verified
Raw History Blame Contribute Delete
7.36 kB
"""
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)),
}