""" Interactive CLI chat with the agent """ import asyncio import os from dataclasses import dataclass from pathlib import Path from typing import Any, Optional import litellm from lmnr import Laminar, LaminarLiteLLMCallback from agent.config import load_config from agent.core.agent_loop import submission_loop from agent.core.session import OpType from agent.core.tools import ToolRouter litellm.drop_params = True lmnr_api_key = os.environ.get("LMNR_API_KEY") if lmnr_api_key: try: Laminar.initialize(project_api_key=lmnr_api_key) litellm.callbacks = [LaminarLiteLLMCallback()] print("āœ… Laminar initialized") except Exception as e: print(f"āš ļø Failed to initialize Laminar: {e}") @dataclass class Operation: """Operation to be executed by the agent""" op_type: OpType data: Optional[dict[str, Any]] = None @dataclass class Submission: """Submission to the agent loop""" id: str operation: Operation async def event_listener( event_queue: asyncio.Queue, turn_complete_event: asyncio.Event, ready_event: asyncio.Event, ) -> None: """Background task that listens for events and displays them""" while True: try: event = await event_queue.get() # Display event if event.event_type == "ready": print("āœ… Agent ready") ready_event.set() elif event.event_type == "assistant_message": content = event.data.get("content", "") if event.data else "" if content: print(f"\nšŸ¤– Assistant: {content}") elif event.event_type == "tool_call": tool_name = event.data.get("tool", "") if event.data else "" if tool_name: print(f"šŸ”§ Calling tool: {tool_name}") elif event.event_type == "tool_output": output = event.data.get("output", "") if event.data else "" success = event.data.get("success", False) if event.data else False status = "āœ…" if success else "āŒ" if output: print(f"{status} Tool output: {output}") elif event.event_type == "turn_complete": print("āœ… Turn complete\n") turn_complete_event.set() elif event.event_type == "error": error = ( event.data.get("error", "Unknown error") if event.data else "Unknown error" ) print(f"āŒ Error: {error}") turn_complete_event.set() elif event.event_type == "shutdown": print("šŸ›‘ Agent shutdown") break elif event.event_type == "processing": print("ā³ Processing...", flush=True) # Silently ignore other events except asyncio.CancelledError: break except Exception as e: print(f"āš ļø Event listener error: {e}") async def get_user_input() -> str: """Get user input asynchronously""" loop = asyncio.get_event_loop() return await loop.run_in_executor(None, input, "You: ") async def main(): """Interactive chat with the agent""" print("=" * 60) print("šŸ¤– Interactive Agent Chat") print("=" * 60) print("Type your messages below. Type 'exit', 'quit', or '/quit' to end.\n") # Create queues for communication submission_queue = asyncio.Queue() event_queue = asyncio.Queue() # Events to signal agent state turn_complete_event = asyncio.Event() turn_complete_event.set() ready_event = asyncio.Event() # Start agent loop in background config_path = Path(__file__).parent / "config_mcp_example.json" config = load_config(config_path) # Create tool router print(f"Config: {config.mcpServers}") tool_router = ToolRouter(config.mcpServers) agent_task = asyncio.create_task( submission_loop( submission_queue, event_queue, config=config, tool_router=tool_router, ) ) # Start event listener in background listener_task = asyncio.create_task( event_listener(event_queue, turn_complete_event, ready_event) ) # Wait for agent to initialize print("ā³ Initializing agent...") await ready_event.wait() submission_id = 0 try: while True: # Wait for previous turn to complete await turn_complete_event.wait() turn_complete_event.clear() # Get user input try: user_input = await get_user_input() except EOFError: break # Check for exit commands if user_input.strip().lower() in ["exit", "quit", "/quit", "/exit"]: break # Skip empty input if not user_input.strip(): turn_complete_event.set() continue # Submit to agent submission_id += 1 submission = Submission( id=f"sub_{submission_id}", operation=Operation( op_type=OpType.USER_INPUT, data={"text": user_input} ), ) print(f"Main submitting: {submission.operation.op_type}") await submission_queue.put(submission) except KeyboardInterrupt: print("\n\nāš ļø Interrupted by user") # Shutdown print("\nšŸ›‘ Shutting down agent...") shutdown_submission = Submission( id="sub_shutdown", operation=Operation(op_type=OpType.SHUTDOWN) ) await submission_queue.put(shutdown_submission) # Wait for tasks to complete await asyncio.wait_for(agent_task, timeout=2.0) listener_task.cancel() print("✨ Goodbye!\n") if __name__ == "__main__": try: asyncio.run(main()) except KeyboardInterrupt: print("\n\n✨ Goodbye!")