From 56e3ab3b56c6bfc567bdee998f3cbf8586e64178 Mon Sep 17 00:00:00 2001 From: Algis Dumbris Date: Fri, 13 Mar 2026 14:25:17 +0200 Subject: [PATCH] feat: E2E agent test with subscription token auth MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Rewrite test to use 3-tier auth fallback (from dialog-engine pattern): ANTHROPIC_API_KEY → CLAUDE_CODE_OAUTH_TOKEN → macOS Keychain. No dedicated API key required — works with Claude subscription. - Auto-start/stop SynapBus server (--auto-server flag) - Dialog-engine tool loop pattern (max rounds + forced text termination) - Token usage and cost tracking - Add spec for E2E agent testing framework (011) Tested: two Claude agents (Alice, Bob) autonomously exchange messages through SynapBus MCP Streamable HTTP. Cost: ~$0.07 per run. Co-Authored-By: Claude Opus 4.6 --- .specify/specs/011-e2e-agent-testing/spec.md | 371 +++++++++++++++ tests/e2e/test_two_agents.py | 472 +++++++++++++------ 2 files changed, 695 insertions(+), 148 deletions(-) create mode 100644 .specify/specs/011-e2e-agent-testing/spec.md diff --git a/.specify/specs/011-e2e-agent-testing/spec.md b/.specify/specs/011-e2e-agent-testing/spec.md new file mode 100644 index 0000000..9c6d193 --- /dev/null +++ b/.specify/specs/011-e2e-agent-testing/spec.md @@ -0,0 +1,371 @@ +# Feature Specification: E2E Agent Testing Framework + +**Feature Branch**: `011-e2e-agent-testing` +**Created**: 2026-03-13 +**Status**: Draft +**Input**: End-to-end testing framework where Claude-powered AI agents communicate through SynapBus MCP, using the Anthropic Python SDK with subscription-based OAuth authentication (macOS Keychain / CLAUDE_CODE_OAUTH_TOKEN fallback). + +## Overview + +A Python test harness that spins up SynapBus locally, registers AI agents, +and runs Claude-powered agents that autonomously interact through SynapBus +MCP tools (Streamable HTTP transport). Each test scenario validates a +different SynapBus feature — direct messaging, channels, task auctions, +agent discovery, attachments, and semantic search. + +### Authentication Strategy (from dialog-engine pattern) + +The test harness authenticates with the Anthropic API using a **3-tier +fallback** — no dedicated API key required: + +1. `ANTHROPIC_API_KEY` env var (if set) +2. `CLAUDE_CODE_OAUTH_TOKEN` env var (OAuth token) +3. **macOS Keychain** — `security find-generic-password -s "claude-code-credentials"` + +This uses the same subscription token as Claude Code / Claude Desktop, +requiring the `anthropic-beta: oauth-2025-04-20` header for OAuth paths. + +### Architecture + +``` +┌─────────────────────────────────────────────────┐ +│ Test Harness (Python) │ +│ │ +│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ +│ │ Agent A │ │ Agent B │ │ Agent C │ │ +│ │ (Claude) │ │ (Claude) │ │ (Claude) │ │ +│ └────┬─────┘ └────┬─────┘ └────┬─────┘ │ +│ │ │ │ │ +│ └───────┬───────┴───────┬───────┘ │ +│ │ │ │ +│ ┌─────┴─────┐ ┌─────┴─────┐ │ +│ │ MCP Client│ │ MCP Client│ │ +│ │ (httpx) │ │ (httpx) │ │ +│ └─────┬─────┘ └─────┬─────┘ │ +└───────────────┼───────────────┼──────────────────┘ + │ │ + ┌──────┴───────────────┴──────┐ + │ SynapBus Server │ + │ localhost:PORT/mcp │ + │ (Streamable HTTP + Auth) │ + └─────────────────────────────┘ +``` + +Each agent is a Python async function that: +1. Holds its own MCP session (Streamable HTTP + Bearer token) +2. Runs a Claude tool-use loop (vanilla `anthropic.AsyncAnthropic`) +3. Has access to SynapBus tools via inline JSON schema definitions +4. Iterates up to N tool rounds, then forces a text-only final response + +## User Scenarios & Testing *(mandatory)* + +### Scenario 1 — Direct Messaging Round-Trip (Priority: P0) + +Two AI agents (Alice, Bob) exchange messages. Alice discovers Bob via +`discover_agents`, sends a research question via `send_message`. Bob +reads inbox via `read_inbox`, claims the message via `claim_messages`, +sends a reply, and marks the original done via `mark_done`. Alice then +reads Bob's reply from her inbox. + +**Why P0**: This validates the core SynapBus value proposition — AI agents +autonomously communicating through MCP tools. If this doesn't work, +nothing else matters. + +**Acceptance Criteria**: + +1. **Given** SynapBus is running and two agents are registered, **When** + Alice's Claude instance calls `discover_agents`, **Then** the result + contains Bob with his capabilities. +2. **Given** Alice sent a message to Bob, **When** Bob's Claude instance + calls `read_inbox`, **Then** the message appears with status `pending`, + correct sender, body, and subject. +3. **Given** Bob claimed and replied to Alice's message, **When** Alice's + Claude instance calls `read_inbox`, **Then** Bob's reply appears. +4. **Given** the full round-trip completes, **Then** the test harness + verifies at least 2 messages exist in the conversation and both agents + produced coherent responses. + +--- + +### Scenario 2 — Channel Group Communication (Priority: P1) + +Three agents join a shared channel. One agent creates a `standard` +channel, invites the others, and broadcasts a question via +`send_channel_message`. The other two agents read their inboxes and +reply on the channel. + +**Why P1**: Channels are the primary group communication mechanism. +Multi-agent collaboration patterns depend on this. + +**Acceptance Criteria**: + +1. **Given** Agent A creates channel "research-team" via `create_channel`, + **When** agents B and C call `join_channel`, **Then** all three appear + in the member list. +2. **Given** Agent A broadcasts a message on the channel, **When** agents + B and C call `read_inbox`, **Then** both receive the channel message. +3. **Given** Agent B replies on the channel, **When** Agent C reads inbox, + **Then** Agent C sees Bob's channel message. + +--- + +### Scenario 3 — Task Auction (Priority: P1) + +An agent posts a task to an `auction` channel. Two other agents bid on +the task. The poster accepts one bid. The winning agent completes the +task. + +**Why P1**: Task auctions are the core swarm intelligence pattern. They +enable dynamic work distribution among agents. + +**Acceptance Criteria**: + +1. **Given** an auction channel exists, **When** Agent A calls `post_task` + with title, description, and requirements, **Then** the task is created + with status `open`. +2. **Given** an open task, **When** agents B and C call `bid_task` with + capabilities and estimates, **Then** both bids are recorded. +3. **Given** two bids exist, **When** Agent A calls `accept_bid` for + Agent B's bid, **Then** the task status becomes `assigned` and Agent B + is the assignee. +4. **Given** Agent B is assigned, **When** Agent B calls `complete_task`, + **Then** the task status becomes `completed`. + +--- + +### Scenario 4 — Agent Discovery (Priority: P2) + +An agent with a specific need uses `discover_agents` to find agents +with matching capabilities, then initiates communication with the +best match. + +**Why P2**: Discovery enables emergent agent collaboration without +hard-coded routing. + +**Acceptance Criteria**: + +1. **Given** agents registered with diverse capabilities (research, + coding, analysis, translation), **When** an agent searches for + "analysis", **Then** agents with analysis capability appear in results. +2. **Given** discovery results, **When** the agent sends a message to a + discovered agent, **Then** the message is delivered successfully. + +--- + +### Scenario 5 — Blackboard Stigmergy (Priority: P2) + +Agents use a `blackboard` channel as shared state. One agent writes +findings, another reads them and builds on top. + +**Why P2**: Stigmergy enables indirect coordination, a key swarm pattern. + +**Acceptance Criteria**: + +1. **Given** a blackboard channel exists, **When** Agent A writes findings + via `send_channel_message` with metadata, **Then** Agent B can read the + findings from the channel. +2. **Given** Agent B reads Agent A's findings, **When** Agent B writes + additional analysis referencing Agent A's work, **Then** Agent A can + read the combined state from the channel. + +--- + +## Technical Design + +### Project Structure + +``` +tests/e2e/ +├── conftest.py # Shared fixtures: server lifecycle, agent registration +├── synapbus_mcp.py # MCP Streamable HTTP client wrapper +├── agent_runner.py # Claude agent loop (Anthropic SDK + tool execution) +├── auth.py # 3-tier auth (API key → OAuth token → Keychain) +├── tools.py # SynapBus tool JSON schemas for Claude +├── test_direct_msg.py # Scenario 1: Direct messaging +├── test_channels.py # Scenario 2: Channel group communication +├── test_task_auction.py # Scenario 3: Task auction +├── test_discovery.py # Scenario 4: Agent discovery +├── test_blackboard.py # Scenario 5: Blackboard stigmergy +└── pyproject.toml # Python deps (anthropic, httpx, pytest, pytest-asyncio) +``` + +### Authentication Module (`auth.py`) + +```python +import os +import subprocess +import anthropic + +def create_client() -> anthropic.AsyncAnthropic: + """Create Anthropic client with 3-tier auth fallback.""" + # Tier 1: Standard API key + api_key = os.environ.get("ANTHROPIC_API_KEY") + if api_key: + return anthropic.AsyncAnthropic(api_key=api_key) + + # Tier 2: OAuth token from env + oauth_token = os.environ.get("CLAUDE_CODE_OAUTH_TOKEN") + + # Tier 3: macOS Keychain (Claude Code stores JSON with OAuth token) + if not oauth_token: + try: + raw = subprocess.check_output( + ["security", "find-generic-password", + "-s", "Claude Code-credentials", "-w"], + text=True, stderr=subprocess.DEVNULL, + ).strip() + if raw: + creds = json.loads(raw) + oauth_token = creds.get("claudeAiOauth", {}).get("accessToken") + except (subprocess.CalledProcessError, FileNotFoundError, json.JSONDecodeError): + pass + + if oauth_token: + return anthropic.AsyncAnthropic( + auth_token=oauth_token, + default_headers={"anthropic-beta": "oauth-2025-04-20"}, + ) + + raise RuntimeError( + "No Anthropic credentials found. Set ANTHROPIC_API_KEY, " + "CLAUDE_CODE_OAUTH_TOKEN, or ensure Claude Code credentials " + "are in macOS Keychain." + ) +``` + +### Agent Runner (`agent_runner.py`) + +Core loop pattern (from dialog-engine): + +```python +async def run_agent( + client: anthropic.AsyncAnthropic, + mcp: SynapBusMCP, + name: str, + system_prompt: str, + user_prompt: str, + tools: list[dict], + model: str = "claude-sonnet-4-6", + max_tool_rounds: int = 5, +) -> AgentResult: + messages = [{"role": "user", "content": user_prompt}] + + for round_num in range(max_tool_rounds + 1): + api_kwargs = dict( + model=model, + max_tokens=2048, + system=system_prompt, + messages=messages, + ) + # Allow tool use except on final round (force text response) + if round_num < max_tool_rounds: + api_kwargs["tools"] = tools + + response = await client.messages.create(**api_kwargs) + + has_tool_use = any(b.type == "tool_use" for b in response.content) + if not has_tool_use or round_num == max_tool_rounds: + # Extract final text + text = "\n".join( + b.text for b in response.content if b.type == "text" + ) + return AgentResult(text=text, ...) + + # Process tool calls → execute via MCP → feed results back + assistant_content = [...] + tool_results = [...] + for block in response.content: + if block.type == "tool_use": + result = mcp.call_tool(block.name, block.input) + tool_results.append({"type": "tool_result", ...}) + + messages.append({"role": "assistant", "content": assistant_content}) + messages.append({"role": "user", "content": tool_results}) +``` + +### MCP Client (`synapbus_mcp.py`) + +Thin wrapper over Streamable HTTP JSON-RPC (already implemented in +`tests/e2e/test_two_agents.py`): + +- `initialize()` → establishes session, captures `Mcp-Session-Id` +- `call_tool(name, args)` → JSON-RPC `tools/call`, returns parsed result +- `list_tools()` → JSON-RPC `tools/list` +- Bearer token auth via `Authorization` header on every request + +### Tool Definitions (`tools.py`) + +Inline JSON schemas for Claude's tool-use API, matching SynapBus MCP tools: + +- `send_message` (to, body, subject, priority) +- `read_inbox` (limit, status_filter, from_agent) +- `claim_messages` (limit) +- `mark_done` (message_id, status) +- `discover_agents` (query) +- `create_channel` (name, description, type) +- `join_channel` (channel_name) +- `send_channel_message` (channel_name, body) +- `post_task` (channel_name, title, description, requirements) +- `bid_task` (task_id, capabilities, time_estimate, message) +- `accept_bid` (task_id, bid_id) +- `complete_task` (task_id) +- `list_tasks` (channel_name, status) +- `search_messages` (query, limit) + +### Test Fixtures (`conftest.py`) + +```python +@pytest.fixture(scope="session") +def synapbus_server(): + """Start SynapBus server, yield URL, stop on teardown.""" + port = find_free_port() + data_dir = tempfile.mkdtemp() + proc = subprocess.Popen( + ["./synapbus", "serve", "--port", str(port), "--data", data_dir], + stdout=subprocess.PIPE, stderr=subprocess.STDOUT, + ) + wait_for_healthy(f"http://localhost:{port}") + yield f"http://localhost:{port}" + proc.terminate() + shutil.rmtree(data_dir) + +@pytest.fixture +async def registered_agents(synapbus_server): + """Register test user + agents, return dict of {name: (mcp, api_key)}.""" + ... +``` + +### Test Model + +Default model: `claude-sonnet-4-6` (fast, cheap for testing). +Override via `--model` flag or `SYNAPBUS_TEST_MODEL` env var. +Each test scenario costs ~$0.02-0.10 depending on conversation length. + +## Non-Functional Requirements + +- **No API key required**: MUST work with macOS Keychain subscription token +- **Self-contained**: Tests start/stop SynapBus automatically +- **Idempotent**: Each test uses a fresh data directory +- **Cost-aware**: Print token usage and estimated cost per scenario +- **Timeout**: Each agent turn has a 30s timeout; full test 5 min max +- **Model override**: Tests accept model parameter for cost control + +## Dependencies + +```toml +[project] +dependencies = [ + "anthropic>=0.49.0", + "httpx>=0.27", + "pytest>=8.0", + "pytest-asyncio>=0.24", +] +``` + +## Out of Scope + +- Streaming responses (batch mode sufficient for testing) +- Attachment upload/download (binary handling adds complexity) +- Semantic search (requires embedding provider configuration) +- Performance/load testing +- CI/CD integration (requires API key management) diff --git a/tests/e2e/test_two_agents.py b/tests/e2e/test_two_agents.py index f7e008e..31f7de8 100644 --- a/tests/e2e/test_two_agents.py +++ b/tests/e2e/test_two_agents.py @@ -2,13 +2,18 @@ """ E2E test: Two Claude-powered agents communicate through SynapBus MCP. +Authentication (3-tier fallback, no API key required): + 1. ANTHROPIC_API_KEY env var + 2. CLAUDE_CODE_OAUTH_TOKEN env var + 3. macOS Keychain (Claude Code subscription credentials) + Prerequisites: - 1. SynapBus running: ./synapbus serve --port 8080 - 2. ANTHROPIC_API_KEY set in environment - 3. pip install anthropic httpx + pip install anthropic httpx Usage: - python tests/e2e/test_two_agents.py [--port 8080] [--model claude-sonnet-4-6] + # Start server first (or use --auto-server): + python tests/e2e/test_two_agents.py --auto-server + python tests/e2e/test_two_agents.py --port 8080 # if server already running """ from __future__ import annotations @@ -16,14 +21,66 @@ from __future__ import annotations import argparse import json import os +import shutil +import signal +import socket +import subprocess import sys +import tempfile +import time +from dataclasses import dataclass, field from typing import Optional import httpx import anthropic -SYNAPBUS_URL = "" -MCP_URL = "" + +# --------------------------------------------------------------------------- +# Authentication (3-tier fallback from dialog-engine) +# --------------------------------------------------------------------------- + +def create_anthropic_client() -> anthropic.Anthropic: + """Create Anthropic client with 3-tier auth fallback. + + 1. ANTHROPIC_API_KEY env var (standard API key) + 2. CLAUDE_CODE_OAUTH_TOKEN env var (OAuth token) + 3. macOS Keychain (Claude Code subscription credentials) + """ + # Tier 1: Standard API key + api_key = os.environ.get("ANTHROPIC_API_KEY") + if api_key: + print(" Auth: using ANTHROPIC_API_KEY") + return anthropic.Anthropic(api_key=api_key) + + # Tier 2: OAuth token from env + oauth_token = os.environ.get("CLAUDE_CODE_OAUTH_TOKEN") + + # Tier 3: macOS Keychain (Claude Code stores credentials as JSON) + if not oauth_token: + try: + raw = subprocess.check_output( + ["security", "find-generic-password", + "-s", "Claude Code-credentials", "-w"], + text=True, stderr=subprocess.DEVNULL, + ).strip() + if raw: + creds = json.loads(raw) + oauth_token = creds.get("claudeAiOauth", {}).get("accessToken") + if oauth_token: + print(" Auth: using macOS Keychain (Claude Code subscription)") + except (subprocess.CalledProcessError, FileNotFoundError, json.JSONDecodeError): + pass + + if oauth_token: + return anthropic.Anthropic( + auth_token=oauth_token, + default_headers={"anthropic-beta": "oauth-2025-04-20"}, + ) + + print("ERROR: No Anthropic credentials found.") + print(" Set ANTHROPIC_API_KEY, CLAUDE_CODE_OAUTH_TOKEN,") + print(" or ensure Claude Code is logged in (macOS Keychain).") + sys.exit(1) # --------------------------------------------------------------------------- @@ -52,7 +109,7 @@ class SynapBusMCP: h["Mcp-Session-Id"] = self.session_id return h - def _rpc(self, method: str, params: dict | None = None) -> dict: + def _rpc(self, method: str, params: Optional[dict] = None) -> dict: body = { "jsonrpc": "2.0", "id": self._next_id(), @@ -64,7 +121,6 @@ class SynapBusMCP: resp = self._client.post(self.url, json=body, headers=self._headers()) resp.raise_for_status() - # Capture session ID from initialize response if sid := resp.headers.get("Mcp-Session-Id"): self.session_id = sid @@ -81,14 +137,13 @@ class SynapBusMCP: result = self._rpc("tools/call", {"name": name, "arguments": arguments}) if "error" in result: return result - # Parse the text content from MCP result content = result.get("result", {}).get("content", []) for block in content: if block.get("type") == "text": return json.loads(block["text"]) return result - def list_tools(self) -> list[dict]: + def list_tools(self) -> list: result = self._rpc("tools/list") return result.get("result", {}).get("tools", []) @@ -97,34 +152,74 @@ class SynapBusMCP: # --------------------------------------------------------------------------- -# Setup: register user + agents, get API keys +# Server lifecycle management # --------------------------------------------------------------------------- -def setup_agents(base_url: str) -> tuple[str, str]: - """Register admin user (if needed) and two agents. Returns (alice_key, bob_key).""" +def find_free_port() -> int: + with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s: + s.bind(("", 0)) + return s.getsockname()[1] + + +def start_server(port: int) -> tuple: + """Start SynapBus server, return (process, data_dir).""" + data_dir = tempfile.mkdtemp(prefix="synapbus-e2e-") + proc = subprocess.Popen( + ["./synapbus", "serve", "--port", str(port), "--data", data_dir], + stdout=subprocess.PIPE, + stderr=subprocess.STDOUT, + text=True, + ) + + # Wait for server to be healthy + for i in range(30): + try: + resp = httpx.get(f"http://localhost:{port}/health", timeout=2) + if resp.status_code == 200: + return proc, data_dir + except httpx.ConnectError: + pass + time.sleep(0.5) + + proc.terminate() + shutil.rmtree(data_dir, ignore_errors=True) + print("ERROR: Server failed to start within 15 seconds") + sys.exit(1) + + +def stop_server(proc: subprocess.Popen, data_dir: str): + """Stop server and clean up.""" + proc.send_signal(signal.SIGTERM) + proc.wait(timeout=10) + shutil.rmtree(data_dir, ignore_errors=True) + + +# --------------------------------------------------------------------------- +# Setup: register user + agents +# --------------------------------------------------------------------------- + +def setup_agents(base_url: str) -> tuple: + """Register test user + two agents. Returns (alice_key, bob_key).""" client = httpx.Client(timeout=10) - # Try to register a test user (might already exist) + # Register test user (ignore if exists) client.post(f"{base_url}/auth/register", json={ - "username": "e2e-tester", + "username": "e2e_tester", "password": "testpass123456", "display_name": "E2E Tester", }) # Login resp = client.post(f"{base_url}/auth/login", json={ - "username": "e2e-tester", + "username": "e2e_tester", "password": "testpass123456", }) if resp.status_code != 200: - # Fall back to admin (auto-created on first run) - print(" [!] Could not login as e2e-tester, trying admin...") - print(" You may need to provide the admin password.") + print(" [!] Login failed. Server may need a fresh data directory.") sys.exit(1) cookies = resp.cookies - # Register agents (ignore errors if already exist) alice_resp = client.post(f"{base_url}/api/agents", json={ "name": "alice", "display_name": "Alice the Researcher", @@ -143,8 +238,9 @@ def setup_agents(base_url: str) -> tuple[str, str]: bob_key = bob_resp.json().get("api_key", "") if not alice_key or not bob_key: - print(" [!] Could not register agents. They may already exist.") - print(" Delete test-data/ and restart the server for a clean test.") + print(" [!] Agent registration failed.") + print(f" Alice: {alice_resp.json()}") + print(f" Bob: {bob_resp.json()}") sys.exit(1) client.close() @@ -152,10 +248,9 @@ def setup_agents(base_url: str) -> tuple[str, str]: # --------------------------------------------------------------------------- -# Claude-powered agent loop +# Tool definitions for Claude # --------------------------------------------------------------------------- -# Tool definitions for Claude (subset of SynapBus MCP tools) TOOLS = [ { "name": "send_message", @@ -182,7 +277,7 @@ TOOLS = [ }, { "name": "discover_agents", - "description": "Discover other agents by capability", + "description": "Discover other agents by capability keyword search", "input_schema": { "type": "object", "properties": { @@ -192,7 +287,7 @@ TOOLS = [ }, { "name": "claim_messages", - "description": "Claim pending messages for processing (atomic lock)", + "description": "Atomically claim pending messages for processing", "input_schema": { "type": "object", "properties": { @@ -202,11 +297,11 @@ TOOLS = [ }, { "name": "mark_done", - "description": "Mark a claimed message as done", + "description": "Mark a previously claimed message as done or failed", "input_schema": { "type": "object", "properties": { - "message_id": {"type": "integer", "description": "Message ID to mark done"}, + "message_id": {"type": "integer", "description": "Message ID"}, "status": {"type": "string", "enum": ["done", "failed"]}, }, "required": ["message_id", "status"], @@ -215,77 +310,110 @@ TOOLS = [ ] +# --------------------------------------------------------------------------- +# Agent runner (dialog-engine pattern: max tool rounds + forced text) +# --------------------------------------------------------------------------- + +@dataclass +class AgentResult: + text: str = "" + tool_calls: list = field(default_factory=list) + input_tokens: int = 0 + output_tokens: int = 0 + + def run_agent( - agent_name: str, + claude: anthropic.Anthropic, mcp: SynapBusMCP, + agent_name: str, system_prompt: str, user_prompt: str, model: str = "claude-sonnet-4-6", - max_turns: int = 5, -) -> str: - """Run a Claude agent that uses SynapBus MCP tools.""" - client = anthropic.Anthropic() + max_tool_rounds: int = 5, +) -> AgentResult: + """Run a Claude agent with SynapBus MCP tools. + + Pattern from dialog-engine: iterate up to max_tool_rounds allowing + tool use. On the final round, omit tools to force a text response. + """ messages = [{"role": "user", "content": user_prompt}] - final_text = "" + result = AgentResult() - for turn in range(max_turns): - print(f" [{agent_name}] Turn {turn + 1}/{max_turns}") - - response = client.messages.create( + for round_num in range(max_tool_rounds + 1): + api_kwargs = dict( model=model, - max_tokens=1024, + max_tokens=2048, system=system_prompt, - tools=TOOLS, messages=messages, ) + # Allow tool use except on final round + if round_num < max_tool_rounds: + api_kwargs["tools"] = TOOLS - # Collect assistant response + response = claude.messages.create(**api_kwargs) + result.input_tokens += response.usage.input_tokens + result.output_tokens += response.usage.output_tokens + + has_tool_use = any(b.type == "tool_use" for b in response.content) + + if not has_tool_use or round_num == max_tool_rounds: + # Extract final text + text_parts = [] + for block in response.content: + if block.type == "text": + text_parts.append(block.text) + result.text = "\n".join(text_parts) if text_parts else "(no response)" + print(f" [{agent_name}] {result.text[:300]}") + return result + + # Process tool calls assistant_content = [] - tool_uses = [] - for block in response.content: if block.type == "text": - final_text += block.text assistant_content.append({"type": "text", "text": block.text}) - print(f" [{agent_name}] {block.text[:200]}") + print(f" [{agent_name}] {block.text[:150]}") elif block.type == "tool_use": - tool_uses.append(block) assistant_content.append({ "type": "tool_use", "id": block.id, "name": block.name, "input": block.input, }) - print(f" [{agent_name}] -> tool: {block.name}({json.dumps(block.input)[:100]})") messages.append({"role": "assistant", "content": assistant_content}) - if response.stop_reason == "end_turn": - break - - # Execute tool calls + # Execute tools via MCP tool_results = [] - for tool_use in tool_uses: - try: - result = mcp.call_tool(tool_use.name, tool_use.input) - result_text = json.dumps(result, indent=2) - print(f" [{agent_name}] <- {tool_use.name}: {result_text[:150]}") - tool_results.append({ - "type": "tool_result", - "tool_use_id": tool_use.id, - "content": result_text, - }) - except Exception as e: - tool_results.append({ - "type": "tool_result", - "tool_use_id": tool_use.id, - "content": f"Error: {e}", - "is_error": True, - }) + for block in response.content: + if block.type == "tool_use": + tool_input_str = json.dumps(block.input)[:80] + print(f" [{agent_name}] -> {block.name}({tool_input_str})") + try: + tool_result = mcp.call_tool(block.name, block.input) + result_str = json.dumps(tool_result) + print(f" [{agent_name}] <- {result_str[:150]}") + result.tool_calls.append({ + "tool": block.name, + "input": block.input, + "output": tool_result, + }) + tool_results.append({ + "type": "tool_result", + "tool_use_id": block.id, + "content": result_str, + }) + except Exception as e: + print(f" [{agent_name}] <- ERROR: {e}") + tool_results.append({ + "type": "tool_result", + "tool_use_id": block.id, + "content": f"Error: {e}", + "is_error": True, + }) messages.append({"role": "user", "content": tool_results}) - return final_text + return result # --------------------------------------------------------------------------- @@ -294,92 +422,140 @@ def run_agent( def main(): parser = argparse.ArgumentParser(description="SynapBus two-agent E2E test") - parser.add_argument("--port", type=int, default=8080, help="SynapBus port") + parser.add_argument("--port", type=int, default=0, help="SynapBus port (0 = auto)") parser.add_argument("--model", default="claude-sonnet-4-6", help="Claude model") + parser.add_argument("--auto-server", action="store_true", + help="Auto-start and stop SynapBus server") args = parser.parse_args() - global SYNAPBUS_URL, MCP_URL - SYNAPBUS_URL = f"http://localhost:{args.port}" - MCP_URL = f"{SYNAPBUS_URL}/mcp" + server_proc = None + data_dir = None - if not os.environ.get("ANTHROPIC_API_KEY"): - print("ERROR: ANTHROPIC_API_KEY not set") - sys.exit(1) - - print(f"=== SynapBus E2E Test ===") - print(f"Server: {SYNAPBUS_URL}") - print(f"Model: {args.model}") - - # 1. Setup - print("\n[1/4] Setting up agents...") - alice_key, bob_key = setup_agents(SYNAPBUS_URL) - print(f" Alice key: {alice_key[:16]}...") - print(f" Bob key: {bob_key[:16]}...") - - # 2. Initialize MCP sessions - print("\n[2/4] Initializing MCP sessions...") - alice_mcp = SynapBusMCP(SYNAPBUS_URL, alice_key) - bob_mcp = SynapBusMCP(SYNAPBUS_URL, bob_key) - alice_mcp.initialize() - bob_mcp.initialize() - print(f" Alice session: {alice_mcp.session_id}") - print(f" Bob session: {bob_mcp.session_id}") - - # 3. Alice sends a message - print("\n[3/4] Alice sends a research request to Bob...") - alice_result = run_agent( - "Alice", - alice_mcp, - system_prompt=( - "You are Alice, a research agent. You communicate with other agents " - "through SynapBus messaging tools. Be concise and direct." - ), - user_prompt=( - "Discover what agents are available, then send a message to 'bob' " - "asking him to analyze the pros and cons of using MCP (Model Context " - "Protocol) for agent-to-agent communication. Include a specific " - "question in your message." - ), - model=args.model, - ) - - # 4. Bob reads inbox and replies - print("\n[4/4] Bob reads inbox and replies...") - bob_result = run_agent( - "Bob", - bob_mcp, - system_prompt=( - "You are Bob, a data analyst agent. You communicate with other agents " - "through SynapBus messaging tools. When you receive messages, read them, " - "process the request, and send a thoughtful reply. Be concise." - ), - user_prompt=( - "Check your inbox for new messages. Read any pending messages, then " - "reply to each sender with a helpful response to their question. " - "After replying, mark the original messages as done." - ), - model=args.model, - ) - - # 5. Verify Alice got the reply - print("\n=== Verification ===") - print("Checking Alice's inbox for Bob's reply...") - alice_inbox = alice_mcp.call_tool("read_inbox", {"limit": 10}) - msgs = alice_inbox.get("messages", []) - if msgs: - for msg in msgs: - print(f" From: {msg['from_agent']}") - print(f" Body: {msg['body'][:200]}") - print(f" Status: {msg['status']}") - print("\n SUCCESS: Alice received Bob's reply") + if args.auto_server or args.port == 0: + port = find_free_port() if args.port == 0 else args.port + print(f"Starting SynapBus on port {port}...") + server_proc, data_dir = start_server(port) + args.auto_server = True else: - print(" WARNING: No messages in Alice's inbox yet") + port = args.port - # Cleanup - alice_mcp.close() - bob_mcp.close() + base_url = f"http://localhost:{port}" - print("\n=== Test Complete ===") + try: + print(f"\n{'=' * 50}") + print(f" SynapBus E2E Agent Test") + print(f" Server: {base_url}") + print(f" Model: {args.model}") + print(f"{'=' * 50}") + + # 1. Authenticate with Anthropic + print("\n[1/5] Authenticating with Anthropic...") + claude = create_anthropic_client() + + # 2. Setup agents + print("\n[2/5] Registering agents...") + alice_key, bob_key = setup_agents(base_url) + print(f" Alice: {alice_key[:16]}...") + print(f" Bob: {bob_key[:16]}...") + + # 3. Initialize MCP sessions + print("\n[3/5] Initializing MCP sessions...") + alice_mcp = SynapBusMCP(base_url, alice_key) + bob_mcp = SynapBusMCP(base_url, bob_key) + alice_mcp.initialize() + bob_mcp.initialize() + print(f" Alice session: {alice_mcp.session_id}") + print(f" Bob session: {bob_mcp.session_id}") + + # List available tools + tools = alice_mcp.list_tools() + print(f" Available tools: {len(tools)}") + for t in tools[:5]: + print(f" - {t['name']}: {t.get('description', '')[:60]}") + if len(tools) > 5: + print(f" ... and {len(tools) - 5} more") + + # 4. Alice discovers agents and sends message + print("\n[4/5] Alice discovers agents and sends research request...") + alice_result = run_agent( + claude, alice_mcp, "Alice", + system_prompt=( + "You are Alice, a research agent on the SynapBus messaging platform. " + "You communicate with other agents using the provided tools. " + "Be concise and direct. Complete your task in as few tool calls as possible." + ), + user_prompt=( + "First, discover what other agents are available using discover_agents. " + "Then send a message to 'bob' asking him to analyze the trade-offs " + "of using MCP (Model Context Protocol) vs REST APIs for agent-to-agent " + "communication. Ask a specific question." + ), + model=args.model, + ) + + # 5. Bob reads inbox, processes, and replies + print("\n[5/5] Bob reads inbox and replies...") + bob_result = run_agent( + claude, bob_mcp, "Bob", + system_prompt=( + "You are Bob, a data analyst agent on the SynapBus messaging platform. " + "You communicate with other agents using the provided tools. " + "When you receive messages, process the request and send a thoughtful reply. " + "Be concise and direct." + ), + user_prompt=( + "Check your inbox for new messages using read_inbox. " + "For each message: claim it with claim_messages, send a reply to the " + "sender using send_message, then mark the original as done with mark_done." + ), + model=args.model, + ) + + # Verification + print(f"\n{'=' * 50}") + print(" VERIFICATION") + print(f"{'=' * 50}") + + # Check Alice's inbox for Bob's reply + alice_inbox = alice_mcp.call_tool("read_inbox", {"limit": 10}) + msgs = alice_inbox.get("messages", []) + bob_replies = [m for m in msgs if m.get("from_agent") == "bob"] + + if bob_replies: + print(f"\n PASS: Alice received {len(bob_replies)} reply(ies) from Bob") + for msg in bob_replies: + body_preview = msg["body"][:200] + print(f" Subject: {msg.get('subject', 'N/A')}") + print(f" Body: {body_preview}") + else: + print("\n FAIL: Alice did not receive a reply from Bob") + print(f" Alice inbox: {len(msgs)} total messages") + + # Cost summary + total_input = alice_result.input_tokens + bob_result.input_tokens + total_output = alice_result.output_tokens + bob_result.output_tokens + # Sonnet pricing: $3/M input, $15/M output + est_cost = (total_input * 3 + total_output * 15) / 1_000_000 + print(f"\n Token usage:") + print(f" Alice: {alice_result.input_tokens} in / {alice_result.output_tokens} out") + print(f" Bob: {bob_result.input_tokens} in / {bob_result.output_tokens} out") + print(f" Total: {total_input} in / {total_output} out") + print(f" Est. cost: ${est_cost:.4f}") + + print(f"\n Tool calls: Alice={len(alice_result.tool_calls)}, Bob={len(bob_result.tool_calls)}") + + # Cleanup + alice_mcp.close() + bob_mcp.close() + + print(f"\n{'=' * 50}") + print(" TEST COMPLETE") + print(f"{'=' * 50}") + + finally: + if server_proc: + print("\nStopping server...") + stop_server(server_proc, data_dir) if __name__ == "__main__":