Spaces:
Running on CPU Upgrade
Running on CPU Upgrade
File size: 6,044 Bytes
8bff299 ef7b74a 8bff299 9de209d 8bff299 ecacd30 8bff299 9de209d ecacd30 8bff299 ecacd30 8bff299 9fe493b 8bff299 ef7b74a ecacd30 ef7b74a ecacd30 6b80d78 ef7b74a 6b80d78 ef7b74a 6b80d78 ef7b74a ecacd30 6b80d78 ef7b74a 8bff299 ef7b74a 8bff299 ecacd30 ef7b74a ecacd30 ef7b74a 8bff299 ecacd30 ef7b74a ecacd30 ef7b74a 8bff299 ef7b74a ecacd30 ef7b74a 8bff299 ef7b74a ecacd30 8bff299 ef7b74a 8bff299 ef7b74a 8bff299 ef7b74a e7068c0 ef7b74a 8bff299 ef7b74a 8bff299 ef7b74a 8bff299 ef7b74a | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 | """
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!")
|