diff --git a/acp_adapter/tools.py b/acp_adapter/tools.py index e3ce3b11492c6..3cf32354b5be9 100644 --- a/acp_adapter/tools.py +++ b/acp_adapter/tools.py @@ -78,7 +78,7 @@ "kanban_block", "kanban_request_review", "kanban_request_changes", "kanban_link", "kanban_heartbeat", "yb_query_group_info", "yb_query_group_members", "yb_search_sticker", - "yb_send_dm", "yb_send_sticker", + "yb_send_dm", "yb_send_sticker", "mixture_of_agents", } diff --git a/agent/display.py b/agent/display.py index 2880cecccb84e..234458b9ee831 100644 --- a/agent/display.py +++ b/agent/display.py @@ -460,7 +460,7 @@ def build_tool_preview(tool_name: str, args: dict, max_len: int | None = None) - "search_files": "pattern", "browser_navigate": "url", "browser_click": "ref", "browser_type": "text", "image_generate": "prompt", "text_to_speech": "text", - "vision_analyze": "question", + "vision_analyze": "question", "mixture_of_agents": "user_prompt", "skill_view": "name", "skills_list": "category", "cronjob": "action", "execute_code": "code", "browser_exec": "code", "delegate_task": "goal", @@ -1530,6 +1530,8 @@ def _wrap(line: str) -> str: return _wrap(f"ā”Š šŸ”Š speak {_trunc(args.get('text', ''), 30)} {dur}") if tool_name == "vision_analyze": return _wrap(f"ā”Š šŸ‘ļø vision {_trunc(args.get('question', ''), 30)} {dur}") + if tool_name == "mixture_of_agents": + return _wrap(f"ā”Š 🧠 reason {_trunc(args.get('user_prompt', ''), 30)} {dur}") if tool_name == "send_message": return _wrap(f"ā”Š šŸ“Ø send {args.get('target', '?')}: \"{_trunc(args.get('message', ''), 25)}\" {dur}") if tool_name == "cronjob": diff --git a/hermes_cli/tips.py b/hermes_cli/tips.py index 0a8f89b4df730..0f4eba7d3f767 100644 --- a/hermes_cli/tips.py +++ b/hermes_cli/tips.py @@ -144,6 +144,7 @@ "The todo tool helps the agent track complex multi-step tasks during a session.", "session_search performs full-text search across ALL past conversations.", "The agent automatically saves preferences, corrections, and environment facts to memory.", + "mixture_of_agents routes hard problems through 4 frontier LLMs collaboratively.", "/moa routes one hard prompt through your configured Mixture of Agents model set.", "Terminal commands support background mode with notify_on_complete for long-running tasks.", "Terminal background processes support watch_patterns to alert on specific output lines.", diff --git a/hermes_cli/tools_config.py b/hermes_cli/tools_config.py index 018178d91676b..90d4c3496f2fe 100644 --- a/hermes_cli/tools_config.py +++ b/hermes_cli/tools_config.py @@ -120,6 +120,9 @@ def _post_setup_no_window_flags(*, streams_to_console: bool = False) -> int: ("discord_admin", "šŸ›”ļø Discord Server Admin", "list channels/roles, pin, assign roles"), ("yuanbao", "šŸ¤– Yuanbao", "group info, member queries, DM"), ("computer_use", "šŸ–±ļø Computer Use (macOS/Windows/Linux)", "background desktop control via cua-driver"), + # Fork-custom (Kyzcreig) — preserve across upstream parity merges. + ("messaging", "šŸ“Ø Cross-Platform Messaging", "send_message"), + ("moa", "🧠 Mixture of Agents", "mixture_of_agents"), ] @@ -152,7 +155,7 @@ def gui_toolset_label(label: str) -> str: # `hermes tools` → X (Twitter) Search setup walks users through credential # setup. The tool's check_fn means the schema still won't appear to the # model if the credential later goes missing or expires. -_DEFAULT_OFF_TOOLSETS = {"homeassistant", "spotify", "discord", "discord_admin", "video", "video_gen", "x_search", "a2a"} +_DEFAULT_OFF_TOOLSETS = {"homeassistant", "spotify", "discord", "discord_admin", "video", "video_gen", "x_search", "a2a", "moa"} # Config-only capabilities: they appear in `hermes tools` for provider/API-key @@ -758,6 +761,7 @@ def _checklist_toolset_keys(platform: str) -> Set[str]: # prompted or read for vision; it's purely a presence marker. TOOLSET_ENV_REQUIREMENTS = { "vision": [("OPENROUTER_API_KEY", "https://openrouter.ai/keys")], + "moa": [("OPENROUTER_API_KEY", "https://openrouter.ai/keys")], } diff --git a/model_tools.py b/model_tools.py index 20ce327a2a6d4..32f416cf8cfcf 100644 --- a/model_tools.py +++ b/model_tools.py @@ -271,6 +271,7 @@ def _run_in_worker(): "web_tools": ["web_search", "web_extract"], "terminal_tools": ["terminal"], "vision_tools": ["vision_analyze"], + "moa_tools": ["mixture_of_agents"], "image_tools": ["image_generate"], "skills_tools": ["skills_list", "skill_view", "skill_manage"], "browser_tools": [ diff --git a/tests/agent/test_fork_custom_toolsets.py b/tests/agent/test_fork_custom_toolsets.py new file mode 100644 index 0000000000000..6b4c9c153d4d3 --- /dev/null +++ b/tests/agent/test_fork_custom_toolsets.py @@ -0,0 +1,133 @@ +"""Fork-custom toolset regression guard (Kyzcreig fork). + +`messaging` (agent-callable ``send_message``) and `moa` (``mixture_of_agents``) +are DELIBERATE fork divergences from upstream — upstream has neither (its stance +is that agents do not get an agent-callable send_message, and it has no MoA tool). + +An upstream parity merge that blindly takes upstream's ``toolsets.py`` / +``tools/`` silently DROPS them: ``resolve_toolset()`` then returns ``[]`` for the +name with no error, and every platform that lists the toolset quietly loses the +tool. That exact regression happened in the 2026-06-29 parity sync (PR #119) and +was only caught by the config-migration dry-run's post-migration toolset +validation. This test makes the invariant a HARD CI gate so the next sync can't +reintroduce it. + +If this test fails after an upstream merge, you dropped a fork feature — restore +the toolset entry in ``toolsets.py`` (and, for ``moa``, the +``tools/mixture_of_agents_tool.py`` module that self-registers it), do not delete +the test. +""" + +import importlib + +import pytest + +from toolsets import TOOLSETS, resolve_toolset, validate_toolset + + +class TestForkCustomToolsetsPresent: + """The two fork-custom toolsets must exist and resolve to their tool(s).""" + + def test_messaging_toolset_resolves_send_message(self): + # Upstream deliberately omits an agent-callable send_message; the fork + # re-enables it (origin-routing work #71/#83 depends on it). + assert "messaging" in TOOLSETS, "fork-custom 'messaging' toolset was dropped (likely an upstream merge)" + assert validate_toolset("messaging") is True + tools = resolve_toolset("messaging") + assert tools, "resolve_toolset('messaging') is empty — the toolset lost its tools" + assert "send_message" in tools + + def test_moa_toolset_resolves_mixture_of_agents(self): + assert "moa" in TOOLSETS, "fork-custom 'moa' toolset was dropped (likely an upstream merge)" + assert validate_toolset("moa") is True + tools = resolve_toolset("moa") + assert tools, "resolve_toolset('moa') is empty — the toolset lost its tools" + assert "mixture_of_agents" in tools + + +class TestMixtureOfAgentsToolRegistered: + """The moa toolset is only useful if the tool module self-registers it.""" + + def test_module_imports_and_self_registers(self): + # The module is auto-discovered by tools.registry.discover_builtin_tools() + # (it contains a registry.register() call). Import it directly here so the + # test is order-independent, then assert the registration landed. + # + # Guard against the editable-install live-tree leak (the very TRAP this + # whole PR is about): a deleted module can still import from a sibling + # install. Pin the module's resolved file to THIS repo's tools/ dir, so a + # phantom import from elsewhere is treated as "missing", not a pass. + import os + import toolsets as _toolsets_mod + + repo_root = os.path.dirname(os.path.abspath(_toolsets_mod.__file__)) + expected = os.path.join(repo_root, "tools", "mixture_of_agents_tool.py") + assert os.path.isfile(expected), ( + "tools/mixture_of_agents_tool.py is missing from THIS repo — the moa " + "tool module was dropped (likely an upstream parity merge)" + ) + + mod = importlib.import_module("tools.mixture_of_agents_tool") + mod_file = getattr(mod, "__file__", None) + assert mod_file is not None and os.path.abspath(mod_file) == expected, ( + f"tools.mixture_of_agents_tool imported from {mod_file}, not this " + f"repo ({expected}) — editable-install leak masking a dropped module" + ) + + from tools.registry import registry + + entry = registry._tools.get("mixture_of_agents") + assert entry is not None, "tools/mixture_of_agents_tool.py no longer registers 'mixture_of_agents'" + assert entry.toolset == "moa" + + def test_tool_schema_shape(self): + import os + import toolsets as _toolsets_mod + + repo_root = os.path.dirname(os.path.abspath(_toolsets_mod.__file__)) + expected = os.path.join(repo_root, "tools", "mixture_of_agents_tool.py") + assert os.path.isfile(expected), "tools/mixture_of_agents_tool.py missing from this repo" + + importlib.import_module("tools.mixture_of_agents_tool") + from tools.registry import registry + + entry = registry._tools["mixture_of_agents"] + schema = entry.schema + assert schema.get("name") == "mixture_of_agents" + # contract: takes a single required free-text prompt + params = schema.get("parameters", {}) + assert "user_prompt" in params.get("properties", {}) + assert "user_prompt" in params.get("required", []) + + +class TestForkToolsetsWiredInToolsUI: + """`hermes tools` must surface both toolsets so a user can toggle them, and + moa must declare its OpenRouter credential requirement.""" + + def test_tools_config_rows_and_env(self): + from hermes_cli import tools_config + + # CONFIGURABLE_TOOLSETS rows is a list of (key, label, desc) tuples. + keys = {row[0] for row in tools_config.CONFIGURABLE_TOOLSETS} + assert "messaging" in keys, "'messaging' missing from the `hermes tools` UI list" + assert "moa" in keys, "'moa' missing from the `hermes tools` UI list" + + # moa is an opt-in (default-off) toolset that needs OPENROUTER_API_KEY. + assert "moa" in tools_config._DEFAULT_OFF_TOOLSETS + env_req = tools_config.TOOLSET_ENV_REQUIREMENTS.get("moa", []) + assert any(name == "OPENROUTER_API_KEY" for name, _url in env_req), \ + "moa must declare its OPENROUTER_API_KEY requirement" + + +def test_send_message_and_moa_not_silently_zero(): + """The headline invariant: neither fork toolset may resolve to zero tools. + + A zero-tool resolve is exactly the silent-drop failure mode (resolve_toolset + returns [] for an unknown/empty toolset with no error). Assert it can't happen + for either fork-custom toolset. + """ + for name in ("messaging", "moa"): + assert len(resolve_toolset(name)) >= 1, ( + f"fork-custom toolset {name!r} resolves to 0 tools — it was dropped, " + "probably by an upstream parity merge taking upstream's toolsets.py" + ) diff --git a/tools/mixture_of_agents_tool.py b/tools/mixture_of_agents_tool.py new file mode 100644 index 0000000000000..6399cf01ed8fa --- /dev/null +++ b/tools/mixture_of_agents_tool.py @@ -0,0 +1,562 @@ +#!/usr/bin/env python3 +""" +Mixture-of-Agents Tool Module + +This module implements the Mixture-of-Agents (MoA) methodology that leverages +the collective strengths of multiple LLMs through a layered architecture to +achieve state-of-the-art performance on complex reasoning tasks. + +Based on the research paper: "Mixture-of-Agents Enhances Large Language Model Capabilities" +by Junlin Wang et al. (arXiv:2406.04692v1) + +Key Features: +- Multi-layer LLM collaboration for enhanced reasoning +- Parallel processing of reference models for efficiency +- Intelligent aggregation and synthesis of diverse responses +- Specialized for extremely difficult problems requiring intense reasoning +- Optimized for coding, mathematics, and complex analytical tasks + +Available Tool: +- mixture_of_agents_tool: Process complex queries using multiple frontier models + +Architecture: +1. Reference models generate diverse initial responses in parallel +2. Aggregator model synthesizes responses into a high-quality output +3. Multiple layers can be used for iterative refinement (future enhancement) + +Models Used (via OpenRouter): +- Reference Models: claude-opus-4.6, gemini-3-pro-preview, gpt-5.4-pro, deepseek-v3.2 +- Aggregator Model: claude-opus-4.6 (highest capability for synthesis) + +Configuration: + To customize the MoA setup, modify the configuration constants at the top of this file: + - REFERENCE_MODELS: List of models for generating diverse initial responses + - AGGREGATOR_MODEL: Model used to synthesize the final response + - REFERENCE_TEMPERATURE/AGGREGATOR_TEMPERATURE: Sampling temperatures + - MIN_SUCCESSFUL_REFERENCES: Minimum successful models needed to proceed + +Usage: + from mixture_of_agents_tool import mixture_of_agents_tool + import asyncio + + # Process a complex query + result = await mixture_of_agents_tool( + user_prompt="Solve this complex mathematical proof..." + ) +""" + +import json +import logging +import os +import asyncio +import datetime +from typing import Dict, Any, List, Optional +from tools.openrouter_client import get_async_client as _get_openrouter_client, check_api_key as check_openrouter_api_key +from agent.auxiliary_client import extract_content_or_reasoning +from tools.debug_helpers import DebugSession +import sys + +logger = logging.getLogger(__name__) + +# Configuration for MoA processing +# Reference models - these generate diverse initial responses in parallel. +# Keep this list aligned with current top-tier OpenRouter frontier options. +REFERENCE_MODELS = [ + "anthropic/claude-opus-4.6", + "google/gemini-2.5-pro", + "openai/gpt-5.4-pro", + "deepseek/deepseek-v3.2", +] + +# Aggregator model - synthesizes reference responses into final output. +# Prefer the strongest synthesis model in the current OpenRouter lineup. +AGGREGATOR_MODEL = "anthropic/claude-opus-4.6" + +# Temperature settings optimized for MoA performance +REFERENCE_TEMPERATURE = 0.6 # Balanced creativity for diverse perspectives +AGGREGATOR_TEMPERATURE = 0.4 # Focused synthesis for consistency + +# Failure handling configuration +MIN_SUCCESSFUL_REFERENCES = 1 # Minimum successful reference models needed to proceed + +# System prompt for the aggregator model (from the research paper) +AGGREGATOR_SYSTEM_PROMPT = """You have been provided with a set of responses from various open-source models to the latest user query. Your task is to synthesize these responses into a single, high-quality response. It is crucial to critically evaluate the information provided in these responses, recognizing that some of it may be biased or incorrect. Your response should not simply replicate the given answers but should offer a refined, accurate, and comprehensive reply to the instruction. Ensure your response is well-structured, coherent, and adheres to the highest standards of accuracy and reliability. + +Responses from models:""" + +_debug = DebugSession("moa_tools", env_var="MOA_TOOLS_DEBUG") + + +def _construct_aggregator_prompt(system_prompt: str, responses: List[str]) -> str: + """ + Construct the final system prompt for the aggregator including all model responses. + + Args: + system_prompt (str): Base system prompt for aggregation + responses (List[str]): List of responses from reference models + + Returns: + str: Complete system prompt with enumerated responses + """ + response_text = "\n".join([f"{i+1}. {response}" for i, response in enumerate(responses)]) + return f"{system_prompt}\n\n{response_text}" + + +async def _run_reference_model_safe( + model: str, + user_prompt: str, + temperature: float = REFERENCE_TEMPERATURE, + max_tokens: int = 32000, + max_retries: int = 6 +) -> tuple[str, str, bool]: + """ + Run a single reference model with retry logic and graceful failure handling. + + Args: + model (str): Model identifier to use + user_prompt (str): The user's query + temperature (float): Sampling temperature for response generation + max_tokens (int): Maximum tokens in response + max_retries (int): Maximum number of retry attempts + + Returns: + tuple[str, str, bool]: (model_name, response_content_or_error, success_flag) + """ + for attempt in range(max_retries): + try: + logger.info("Querying %s (attempt %s/%s)", model, attempt + 1, max_retries) + + # Build parameters for the API call + api_params = { + "model": model, + "messages": [{"role": "user", "content": user_prompt}], + "max_tokens": max_tokens, + "extra_body": { + "reasoning": { + "enabled": True, + "effort": "xhigh" + } + } + } + + # GPT models (especially gpt-4o-mini) don't support custom temperature + # values. Match both bare ("gpt-...") and OpenRouter-namespaced + # ("openai/gpt-...") model IDs so the guard actually fires. + _ml = model.lower() + if not (_ml.startswith("gpt-") or _ml.startswith("openai/gpt-")): + api_params["temperature"] = temperature + + response = await _get_openrouter_client().chat.completions.create(**api_params) + + content = extract_content_or_reasoning(response) + if not content: + # Reasoning-only / empty response — let the retry loop handle it. + logger.warning("%s returned empty content (attempt %s/%s), retrying", model, attempt + 1, max_retries) + if attempt < max_retries - 1: + await asyncio.sleep(min(2 ** (attempt + 1), 60)) + continue + # Final attempt still empty: do NOT report success with an empty + # body — an empty string forwarded to the aggregator pollutes the + # synthesis. Treat it as a failed reference. + logger.error("%s returned empty content after %s attempts — marking failed", model, max_retries) + return model, f"{model} returned empty content after {max_retries} attempts", False + logger.info("%s responded (%s characters)", model, len(content)) + return model, content, True + + except Exception as e: + error_str = str(e) + # Keep retry-path logging concise; full tracebacks are reserved for + # terminal failure paths so long-running MoA retries don't flood logs. + if "invalid" in error_str.lower(): + logger.warning("%s invalid request error (attempt %s): %s", model, attempt + 1, error_str) + elif "rate" in error_str.lower() or "limit" in error_str.lower(): + logger.warning("%s rate limit error (attempt %s): %s", model, attempt + 1, error_str) + else: + logger.warning("%s unknown error (attempt %s): %s", model, attempt + 1, error_str) + + if attempt < max_retries - 1: + # Exponential backoff for rate limiting: 2s, 4s, 8s, 16s, 32s, 60s + sleep_time = min(2 ** (attempt + 1), 60) + logger.info("Retrying in %ss...", sleep_time) + await asyncio.sleep(sleep_time) + else: + error_msg = f"{model} failed after {max_retries} attempts: {error_str}" + logger.error("%s", error_msg, exc_info=True) + return model, error_msg, False + + +async def _run_aggregator_model( + system_prompt: str, + user_prompt: str, + temperature: float = AGGREGATOR_TEMPERATURE, + max_tokens: Optional[int] = None, + model: Optional[str] = None, +) -> str: + """ + Run the aggregator model to synthesize the final response. + + Args: + system_prompt (str): System prompt with all reference responses + user_prompt (str): Original user query + temperature (float): Focused temperature for consistent aggregation + max_tokens (int): Maximum tokens in final response + + Returns: + str: Synthesized final response + """ + agg_model = model or AGGREGATOR_MODEL + logger.info("Running aggregator model: %s", agg_model) + + # Build parameters for the API call + api_params = { + "model": agg_model, + "messages": [ + {"role": "system", "content": system_prompt}, + {"role": "user", "content": user_prompt} + ], + "max_tokens": max_tokens, + "extra_body": { + "reasoning": { + "enabled": True, + "effort": "xhigh" + } + } + } + + # GPT models don't support custom temperature values. Match both bare + # ("gpt-...") and OpenRouter-namespaced ("openai/gpt-...") IDs. + _aml = agg_model.lower() + if not (_aml.startswith("gpt-") or _aml.startswith("openai/gpt-")): + api_params["temperature"] = temperature + + response = await _get_openrouter_client().chat.completions.create(**api_params) + + content = extract_content_or_reasoning(response) + + # Retry once on empty content (reasoning-only response) + if not content: + logger.warning("Aggregator returned empty content, retrying once") + response = await _get_openrouter_client().chat.completions.create(**api_params) + content = extract_content_or_reasoning(response) + + # Still empty after the retry: do NOT return an empty success. Raise so the + # caller's failure path reports it instead of handing the agent a well-formed + # {"success": True, "response": ""} with no signal that synthesis produced + # nothing. Symmetric with the final-attempt guard in _run_reference_model_safe. + if not content: + raise RuntimeError( + f"Aggregator model {agg_model} returned empty content after retry" + ) + + logger.info("Aggregation complete (%s characters)", len(content)) + return content + + +async def mixture_of_agents_tool( + user_prompt: str, + reference_models: Optional[List[str]] = None, + aggregator_model: Optional[str] = None +) -> str: + """ + Process a complex query using the Mixture-of-Agents methodology. + + This tool leverages multiple frontier language models to collaboratively solve + extremely difficult problems requiring intense reasoning. It's particularly + effective for: + - Complex mathematical proofs and calculations + - Advanced coding problems and algorithm design + - Multi-step analytical reasoning tasks + - Problems requiring diverse domain expertise + - Tasks where single models show limitations + + The MoA approach uses a fixed 2-layer architecture: + 1. Layer 1: Multiple reference models generate diverse responses in parallel (temp=0.6) + 2. Layer 2: Aggregator model synthesizes the best elements into final response (temp=0.4) + + Args: + user_prompt (str): The complex query or problem to solve + reference_models (Optional[List[str]]): Custom reference models to use + aggregator_model (Optional[str]): Custom aggregator model to use + + Returns: + str: JSON string containing the MoA results with the following structure: + { + "success": bool, + "response": str, + "models_used": { + "reference_models": List[str], + "aggregator_model": str + }, + "processing_time": float + } + + Raises: + Exception: If MoA processing fails or API key is not set + """ + start_time = datetime.datetime.now() + + debug_call_data = { + "parameters": { + "user_prompt": user_prompt[:200] + "..." if len(user_prompt) > 200 else user_prompt, + "reference_models": reference_models or REFERENCE_MODELS, + "aggregator_model": aggregator_model or AGGREGATOR_MODEL, + "reference_temperature": REFERENCE_TEMPERATURE, + "aggregator_temperature": AGGREGATOR_TEMPERATURE, + "min_successful_references": MIN_SUCCESSFUL_REFERENCES + }, + "error": None, + "success": False, + "reference_responses_count": 0, + "failed_models_count": 0, + "failed_models": [], + "final_response_length": 0, + "processing_time_seconds": 0, + "models_used": {} + } + + try: + logger.info("Starting Mixture-of-Agents processing...") + logger.info("Query: %s", user_prompt[:100]) + + # Validate API key availability + if not os.getenv("OPENROUTER_API_KEY"): + raise ValueError("OPENROUTER_API_KEY environment variable not set") + + # Use provided models or defaults + ref_models = reference_models or REFERENCE_MODELS + agg_model = aggregator_model or AGGREGATOR_MODEL + + logger.info("Using %s reference models in 2-layer MoA architecture", len(ref_models)) + + # Layer 1: Generate diverse responses from reference models (with failure handling) + logger.info("Layer 1: Generating reference responses...") + model_results = await asyncio.gather(*[ + _run_reference_model_safe(model, user_prompt, REFERENCE_TEMPERATURE) + for model in ref_models + ]) + + # Separate successful and failed responses + successful_responses = [] + failed_models = [] + + for model_name, content, success in model_results: + if success: + successful_responses.append(content) + else: + failed_models.append(model_name) + + successful_count = len(successful_responses) + failed_count = len(failed_models) + + logger.info("Reference model results: %s successful, %s failed", successful_count, failed_count) + + if failed_models: + logger.warning("Failed models: %s", ', '.join(failed_models)) + + # Check if we have enough successful responses to proceed + if successful_count < MIN_SUCCESSFUL_REFERENCES: + raise ValueError(f"Insufficient successful reference models ({successful_count}/{len(ref_models)}). Need at least {MIN_SUCCESSFUL_REFERENCES} successful responses.") + + debug_call_data["reference_responses_count"] = successful_count + debug_call_data["failed_models_count"] = failed_count + debug_call_data["failed_models"] = failed_models + + # Layer 2: Aggregate responses using the aggregator model + logger.info("Layer 2: Synthesizing final response...") + aggregator_system_prompt = _construct_aggregator_prompt( + AGGREGATOR_SYSTEM_PROMPT, + successful_responses + ) + + final_response = await _run_aggregator_model( + aggregator_system_prompt, + user_prompt, + AGGREGATOR_TEMPERATURE, + model=agg_model, + ) + + # Calculate processing time + end_time = datetime.datetime.now() + processing_time = (end_time - start_time).total_seconds() + + logger.info("MoA processing completed in %.2f seconds", processing_time) + + # Prepare successful response (only final aggregated result, minimal fields) + result = { + "success": True, + "response": final_response, + "models_used": { + "reference_models": ref_models, + "aggregator_model": agg_model + } + } + + debug_call_data["success"] = True + debug_call_data["final_response_length"] = len(final_response) + debug_call_data["processing_time_seconds"] = processing_time + debug_call_data["models_used"] = result["models_used"] + + # Log debug information + _debug.log_call("mixture_of_agents_tool", debug_call_data) + _debug.save() + + return json.dumps(result, indent=2, ensure_ascii=False) + + except Exception as e: + error_msg = f"Error in MoA processing: {str(e)}" + logger.error("%s", error_msg, exc_info=True) + + # Calculate processing time even for errors + end_time = datetime.datetime.now() + processing_time = (end_time - start_time).total_seconds() + + # Prepare error response (minimal fields) + result = { + "success": False, + "response": "MoA processing failed. Please try again or use a single model for this query.", + "models_used": { + "reference_models": reference_models or REFERENCE_MODELS, + "aggregator_model": aggregator_model or AGGREGATOR_MODEL + }, + "error": error_msg + } + + debug_call_data["error"] = error_msg + debug_call_data["processing_time_seconds"] = processing_time + _debug.log_call("mixture_of_agents_tool", debug_call_data) + _debug.save() + + return json.dumps(result, indent=2, ensure_ascii=False) + + +def check_moa_requirements() -> bool: + """ + Check if all requirements for MoA tools are met. + + Returns: + bool: True if requirements are met, False otherwise + """ + return check_openrouter_api_key() + + + +def get_moa_configuration() -> Dict[str, Any]: + """ + Get the current MoA configuration settings. + + Returns: + Dict[str, Any]: Dictionary containing all configuration parameters + """ + return { + "reference_models": REFERENCE_MODELS, + "aggregator_model": AGGREGATOR_MODEL, + "reference_temperature": REFERENCE_TEMPERATURE, + "aggregator_temperature": AGGREGATOR_TEMPERATURE, + "min_successful_references": MIN_SUCCESSFUL_REFERENCES, + "total_reference_models": len(REFERENCE_MODELS), + "failure_tolerance": f"{len(REFERENCE_MODELS) - MIN_SUCCESSFUL_REFERENCES}/{len(REFERENCE_MODELS)} models can fail" + } + + +if __name__ == "__main__": + """ + Simple test/demo when run directly + """ + print("šŸ¤– Mixture-of-Agents Tool Module") + print("=" * 50) + + # Check if API key is available + api_available = check_openrouter_api_key() + + if not api_available: + print("āŒ OPENROUTER_API_KEY environment variable not set") + print("Please set your API key: export OPENROUTER_API_KEY='your-key-here'") + print("Get API key at: https://openrouter.ai/") + sys.exit(1) + else: + print("āœ… OpenRouter API key found") + + print("šŸ› ļø MoA tools ready for use!") + + # Show current configuration + config = get_moa_configuration() + print("\nāš™ļø Current Configuration:") + print(f" šŸ¤– Reference models ({len(config['reference_models'])}): {', '.join(config['reference_models'])}") + print(f" 🧠 Aggregator model: {config['aggregator_model']}") + print(f" šŸŒ”ļø Reference temperature: {config['reference_temperature']}") + print(f" šŸŒ”ļø Aggregator temperature: {config['aggregator_temperature']}") + print(f" šŸ›”ļø Failure tolerance: {config['failure_tolerance']}") + print(f" šŸ“Š Minimum successful models: {config['min_successful_references']}") + + # Show debug mode status + if _debug.active: + print(f"\nšŸ› Debug mode ENABLED - Session ID: {_debug.session_id}") + print(f" Debug logs will be saved to: ./logs/moa_tools_debug_{_debug.session_id}.json") + else: + print("\nšŸ› Debug mode disabled (set MOA_TOOLS_DEBUG=true to enable)") + + print("\nBasic usage:") + print(" from mixture_of_agents_tool import mixture_of_agents_tool") + print(" import asyncio") + print("") + print(" async def main():") + print(" result = await mixture_of_agents_tool(") + print(" user_prompt='Solve this complex mathematical proof...'") + print(" )") + print(" print(result)") + print(" asyncio.run(main())") + + print("\nBest use cases:") + print(" - Complex mathematical proofs and calculations") + print(" - Advanced coding problems and algorithm design") + print(" - Multi-step analytical reasoning tasks") + print(" - Problems requiring diverse domain expertise") + print(" - Tasks where single models show limitations") + + print("\nPerformance characteristics:") + print(" - Higher latency due to multiple model calls") + print(" - Significantly improved quality for complex tasks") + print(" - Parallel processing for efficiency") + print(f" - Optimized temperatures: {REFERENCE_TEMPERATURE} for reference models, {AGGREGATOR_TEMPERATURE} for aggregation") + print(" - Token-efficient: only returns final aggregated response") + print(" - Resilient: continues with partial model failures") + print(" - Configurable: easy to modify models and settings at top of file") + print(" - State-of-the-art results on challenging benchmarks") + + print("\nDebug mode:") + print(" # Enable debug logging") + print(" export MOA_TOOLS_DEBUG=true") + print(" # Debug logs capture all MoA processing steps and metrics") + print(" # Logs saved to: ./logs/moa_tools_debug_UUID.json") + + +# --------------------------------------------------------------------------- +# Registry +# --------------------------------------------------------------------------- +from tools.registry import registry + +MOA_SCHEMA = { + "name": "mixture_of_agents", + "description": "Route a hard problem through multiple frontier LLMs collaboratively. Makes 5 API calls (4 reference models + 1 aggregator) with maximum reasoning effort — use sparingly for genuinely difficult problems. Best for: complex math, advanced algorithms, multi-step analytical reasoning, problems benefiting from diverse perspectives.", + "parameters": { + "type": "object", + "properties": { + "user_prompt": { + "type": "string", + "description": "The complex query or problem to solve using multiple AI models. Should be a challenging problem that benefits from diverse perspectives and collaborative reasoning." + } + }, + "required": ["user_prompt"] + } +} + +registry.register( + name="mixture_of_agents", + toolset="moa", + schema=MOA_SCHEMA, + handler=lambda args, **kw: mixture_of_agents_tool(user_prompt=args.get("user_prompt", "")), + check_fn=check_moa_requirements, + requires_env=["OPENROUTER_API_KEY"], + is_async=True, + emoji="🧠", +) diff --git a/toolsets.py b/toolsets.py index c7b218c445112..8cd217c306380 100644 --- a/toolsets.py +++ b/toolsets.py @@ -367,6 +367,35 @@ "includes": [] }, + # ========================================================================== + # FORK-CUSTOM toolsets (Kyzcreig fork — NOT present upstream). + # + # These are a DELIBERATE divergence from upstream and an upstream parity + # merge MUST preserve them (they are dropped if the merge blindly takes + # upstream's toolsets.py). See tests/agent/test_fork_custom_toolsets.py, + # which fails CI if either toolset stops resolving. + # + # * messaging — upstream's stance (see the note further down) is that + # agents do NOT get an agent-callable send_message; outbound messaging + # is handled outside the agent loop. THIS FORK intentionally re-enables + # an agent-callable send_message: the fleet's orchestrator sends + # mid-turn status/alerts to the originating channel, and the fork's + # origin-routing work (send_message origin-leak fixes) is built on it. + # * moa — the mixture_of_agents tool (tools/mixture_of_agents_tool.py), + # a fork-only multi-LLM reasoning tool. + # ========================================================================== + "messaging": { + "description": "Cross-platform messaging: send messages to Telegram, Discord, Slack, SMS, etc.", + "tools": ["send_message"], + "includes": [] + }, + + "moa": { + "description": "Advanced reasoning and problem-solving tools", + "tools": ["mixture_of_agents"], + "includes": [] + }, + # Scenario-specific toolsets