Spaces:
Sleeping
Sleeping
Cyber Catalyst Team
Support both goal and problem_statement parameter inputs in EternityInitRequest endpoint
be32809 | """ | |
| Claude Code Backend — Agentic coding backend powered by NVIDIA NIM models. | |
| Exposes an OpenAI-compatible /v1/chat/completions endpoint with built-in | |
| tools for file operations and bash execution. | |
| Architecture: | |
| Space 1 (better-chatbot) --> this backend --> NVIDIA NIM API | |
| The agentic loop: | |
| 1. Receive user message from Space 1 | |
| 2. Send to NIM model with tool definitions | |
| 3. If model returns tool_calls, execute them and loop | |
| 4. If model returns text, stream it back to Space 1 | |
| 5. Persist conversation in Postgres | |
| """ | |
| import os | |
| import shutil | |
| import threading | |
| import json | |
| import uuid | |
| import subprocess | |
| import asyncio | |
| import time | |
| import re | |
| import collections | |
| from pathlib import Path | |
| from typing import AsyncIterator, Optional, List, Dict, Any | |
| from pydantic import BaseModel | |
| from fastapi import FastAPI, Request, Header, HTTPException | |
| from fastapi.responses import StreamingResponse, JSONResponse, HTMLResponse | |
| from fastapi.middleware.cors import CORSMiddleware | |
| from openai import AsyncOpenAI | |
| import anyio | |
| import asyncpg | |
| # --------------------------------------------------------------------------- | |
| # Globals & Activity Logs | |
| # --------------------------------------------------------------------------- | |
| activity_logs = collections.deque(maxlen=100) | |
| MODEL_STATUSES = {} | |
| ACTIVE_SESSIONS = set() | |
| def log_activity(msg: str): | |
| timestamp = time.strftime("%H:%M:%S") | |
| log_line = f"[{timestamp}] {msg}" | |
| activity_logs.append(log_line) | |
| print(log_line) | |
| # --------------------------------------------------------------------------- | |
| # Configuration | |
| # --------------------------------------------------------------------------- | |
| NIM_API_KEY = os.environ.get("NVIDIA_NIM_API_KEY", "") | |
| BACKEND_API_KEY = os.environ.get("BACKEND_API_KEY", "") | |
| DATABASE_URL = os.environ.get("DATABASE_URL", "") | |
| WORKSPACE_DIR = os.environ.get("WORKSPACE_DIR", "/tmp/workspace") | |
| BACKUP_GIT_REPO = os.environ.get("BACKUP_GIT_REPO", "") | |
| MAX_TOOL_ROUNDS = int(os.environ.get("MAX_TOOL_ROUNDS", "10")) | |
| # NIM models that reliably support tool/function calling | |
| TOOL_CAPABLE_MODELS = { | |
| "nvidia/nemotron-3-ultra-550b-a55b": "Nemotron 3 Ultra 550B (Agentic)", | |
| "z-ai/glm-5.1": "GLM 5.1 (Agentic)", | |
| "moonshotai/kimi-k2.6": "Kimi K2.6 (Agentic)", | |
| "minimaxai/minimax-m3": "MiniMax M3 (Agentic)", | |
| "stepfun-ai/step-3.7-flash": "Step 3.7 Flash (Agentic)", | |
| "minimaxai/minimax-m2.7": "MiniMax M2.7 (Agentic)", | |
| "meta/llama-3.1-70b-instruct": "Llama 3.1 70B (Agentic)", | |
| "meta/llama-3.1-405b-instruct": "Llama 3.1 405B (Agentic)", | |
| "qwen/qwen2.5-coder-32b-instruct": "Qwen 2.5 Coder 32B (Agentic)", | |
| "nvidia/llama-3.1-nemotron-70b-instruct": "Nemotron 70B (Agentic)", | |
| "meta/llama-3.3-70b-instruct": "Llama 3.3 70B (Agentic)", | |
| } | |
| # All models (tool-capable get agentic mode, others get plain chat) | |
| ALL_MODELS = { | |
| **TOOL_CAPABLE_MODELS, | |
| "deepseek-ai/deepseek-r1": "DeepSeek R1 (Chat only)", | |
| "mistralai/mistral-large-2-instruct": "Mistral Large 2 (Chat only)", | |
| } | |
| RECOMMENDED_MODEL = "nvidia/llama-3.1-nemotron-70b-instruct" | |
| # Ensure workspace exists | |
| Path(WORKSPACE_DIR).mkdir(parents=True, exist_ok=True) | |
| # --------------------------------------------------------------------------- | |
| # NIM Client | |
| # --------------------------------------------------------------------------- | |
| nim_client = AsyncOpenAI( | |
| base_url="https://integrate.api.nvidia.com/v1", | |
| api_key=NIM_API_KEY, | |
| ) | |
| # --------------------------------------------------------------------------- | |
| # Rate Limiting & Multi-Provider Setup | |
| # --------------------------------------------------------------------------- | |
| MISTRAL_API_KEY = os.environ.get("MISTRAL_API_KEY", "") | |
| mistral_client = AsyncOpenAI( | |
| base_url="https://api.mistral.ai/v1", | |
| api_key=MISTRAL_API_KEY if MISTRAL_API_KEY else "dummy_key", | |
| ) if MISTRAL_API_KEY else None | |
| class MultiProviderRateLimiter: | |
| def __init__(self): | |
| self.nim_limit = 40 | |
| self.nim_window = 60 | |
| self.nim_calls = [] | |
| self.mistral_last_call = 0.0 | |
| self.lock = asyncio.Lock() | |
| async def wait_for_mistral(self): | |
| async with self.lock: | |
| now = time.time() | |
| elapsed = now - self.mistral_last_call | |
| if elapsed < 1.0: | |
| await asyncio.sleep(1.0 - elapsed) | |
| self.mistral_last_call = time.time() | |
| async def wait_for_nim(self): | |
| async with self.lock: | |
| now = time.time() | |
| self.nim_calls = [t for t in self.nim_calls if now - t < self.nim_window] | |
| if len(self.nim_calls) >= self.nim_limit - 2: | |
| sleep_time = self.nim_window - (now - self.nim_calls[0]) | |
| print(f"[RateLimiter] Approaching NIM rate limit (40 RPM). Sleeping {sleep_time:.2f}s...") | |
| await asyncio.sleep(sleep_time) | |
| self.nim_calls.append(time.time()) | |
| rate_limiter = MultiProviderRateLimiter() | |
| # --------------------------------------------------------------------------- | |
| # Tool Definitions (OpenAI function calling format) | |
| # --------------------------------------------------------------------------- | |
| TOOLS = [ | |
| { | |
| "type": "function", | |
| "function": { | |
| "name": "read_file", | |
| "description": "Read the contents of a file. Use this to inspect existing code, configs, or any text file.", | |
| "parameters": { | |
| "type": "object", | |
| "properties": { | |
| "path": { | |
| "type": "string", | |
| "description": "Relative path to the file from the workspace root" | |
| } | |
| }, | |
| "required": ["path"] | |
| } | |
| } | |
| }, | |
| { | |
| "type": "function", | |
| "function": { | |
| "name": "write_file", | |
| "description": "Write content to a file. Creates the file if it doesn't exist, overwrites if it does. Creates parent directories automatically.", | |
| "parameters": { | |
| "type": "object", | |
| "properties": { | |
| "path": { | |
| "type": "string", | |
| "description": "Relative path to the file from the workspace root" | |
| }, | |
| "content": { | |
| "type": "string", | |
| "description": "The full content to write to the file" | |
| } | |
| }, | |
| "required": ["path", "content"] | |
| } | |
| } | |
| }, | |
| { | |
| "type": "function", | |
| "function": { | |
| "name": "run_bash", | |
| "description": "Execute a bash command in the workspace directory. Use for installing packages, running scripts, git operations, etc. Commands run with a 30 second timeout.", | |
| "parameters": { | |
| "type": "object", | |
| "properties": { | |
| "command": { | |
| "type": "string", | |
| "description": "The bash command to execute" | |
| } | |
| }, | |
| "required": ["command"] | |
| } | |
| } | |
| }, | |
| { | |
| "type": "function", | |
| "function": { | |
| "name": "list_directory", | |
| "description": "List files and directories in a given path. Shows file sizes and directory markers.", | |
| "parameters": { | |
| "type": "object", | |
| "properties": { | |
| "path": { | |
| "type": "string", | |
| "description": "Relative path to the directory from workspace root. Use '.' for the workspace root." | |
| } | |
| }, | |
| "required": ["path"] | |
| } | |
| } | |
| }, | |
| { | |
| "type": "function", | |
| "function": { | |
| "name": "grep_search", | |
| "description": "Search for a pattern in files within the workspace. Returns matching lines with file paths and line numbers.", | |
| "parameters": { | |
| "type": "object", | |
| "properties": { | |
| "pattern": { | |
| "type": "string", | |
| "description": "The search pattern (supports basic regex)" | |
| }, | |
| "path": { | |
| "type": "string", | |
| "description": "Directory or file to search in, relative to workspace root. Defaults to '.'", | |
| } | |
| }, | |
| "required": ["pattern"] | |
| } | |
| } | |
| }, | |
| ] | |
| # --------------------------------------------------------------------------- | |
| # Tool Execution | |
| # --------------------------------------------------------------------------- | |
| def _safe_path(rel_path: str) -> Path: | |
| """Resolve a relative path safely within the workspace.""" | |
| workspace = Path(WORKSPACE_DIR).resolve() | |
| target = (workspace / rel_path).resolve() | |
| # Prevent path traversal | |
| if not str(target).startswith(str(workspace)): | |
| raise ValueError(f"Path traversal detected: {rel_path}") | |
| return target | |
| def repair_arguments(func_name: str, args: dict) -> tuple[dict, list[str]]: | |
| notes = [] | |
| repaired_args = dict(args) | |
| # 1. Nesting extraction (e.g. {"path": {"path": "file.txt"}}) | |
| for key in list(repaired_args.keys()): | |
| val = repaired_args[key] | |
| if isinstance(val, dict) and key in val: | |
| repaired_args[key] = val[key] | |
| notes.append(f"Flattened nested parameter '{key}'") | |
| # 2. Markdown stripping from bash command | |
| if func_name == "run_bash" and "command" in repaired_args: | |
| cmd = repaired_args["command"] | |
| if isinstance(cmd, str): | |
| pattern = r"```(?:bash)?\s*(.*?)\s*```" | |
| match = re.search(pattern, cmd, re.DOTALL) | |
| if match: | |
| repaired_args["command"] = match.group(1).strip() | |
| notes.append("Stripped markdown code blocks from bash command") | |
| # 3. Stringified array conversion | |
| for key, val in repaired_args.items(): | |
| if isinstance(val, str) and val.strip().startswith("[") and val.strip().endswith("]"): | |
| try: | |
| parsed_arr = json.loads(val) | |
| if isinstance(parsed_arr, list): | |
| repaired_args[key] = parsed_arr | |
| notes.append(f"Converted stringified array for parameter '{key}' to native array") | |
| except: | |
| pass | |
| # 4. Optional empty objects replacing Null | |
| for key in list(repaired_args.keys()): | |
| if repaired_args[key] == {}: | |
| repaired_args[key] = None | |
| notes.append(f"Replaced empty object for parameter '{key}' with null") | |
| return repaired_args, notes | |
| async def execute_tool(name: str, arguments: dict) -> str: | |
| """Execute a tool and return its output as a string asynchronously.""" | |
| try: | |
| if name == "read_file": | |
| path = _safe_path(arguments["path"]) | |
| if not path.exists(): | |
| return f"Error: File not found: {arguments['path']}" | |
| if not path.is_file(): | |
| return f"Error: Not a file: {arguments['path']}" | |
| content = path.read_text(encoding="utf-8", errors="replace") | |
| if len(content) > 50000: | |
| return content[:50000] + f"\n\n[Truncated — file is {len(content)} chars]" | |
| return content | |
| elif name == "write_file": | |
| path = _safe_path(arguments["path"]) | |
| path.parent.mkdir(parents=True, exist_ok=True) | |
| path.write_text(arguments["content"], encoding="utf-8") | |
| return f"Successfully wrote {len(arguments['content'])} chars to {arguments['path']}" | |
| elif name == "run_bash": | |
| command = arguments["command"] | |
| # Safety: block dangerous commands | |
| blocked = ["rm -rf /", "mkfs", "dd if=", ":(){", "fork bomb"] | |
| if any(b in command.lower() for b in blocked): | |
| return "Error: Command blocked for safety reasons" | |
| # ASYNC SUBPROCESS - This prevents the FastAPI server from freezing! | |
| process = await asyncio.create_subprocess_shell( | |
| command, | |
| stdout=asyncio.subprocess.PIPE, | |
| stderr=asyncio.subprocess.PIPE, | |
| cwd=WORKSPACE_DIR, | |
| env={**os.environ, "HOME": "/tmp", "PATH": os.environ.get("PATH", "/usr/local/bin:/usr/bin:/bin")}, | |
| ) | |
| try: | |
| stdout, stderr = await asyncio.wait_for(process.communicate(), timeout=30) | |
| output = "" | |
| if stdout: | |
| output += stdout.decode('utf-8', errors='replace') | |
| if stderr: | |
| output += ("\n" if output else "") + f"[stderr] {stderr.decode('utf-8', errors='replace')}" | |
| if process.returncode != 0: | |
| output += f"\n[exit code: {process.returncode}]" | |
| if not output: | |
| output = "[command completed with no output]" | |
| except asyncio.TimeoutError: | |
| try: | |
| process.kill() | |
| except Exception: | |
| pass | |
| await process.communicate() | |
| return "Error: Command timed out after 30 seconds" | |
| if len(output) > 20000: | |
| output = output[:20000] + f"\n\n[Truncated — output is {len(output)} chars]" | |
| return output | |
| elif name == "list_directory": | |
| path = _safe_path(arguments.get("path", ".")) | |
| if not path.exists(): | |
| return f"Error: Directory not found: {arguments.get('path', '.')}" | |
| if not path.is_dir(): | |
| return f"Error: Not a directory: {arguments.get('path', '.')}" | |
| entries = [] | |
| for item in sorted(path.iterdir()): | |
| if item.is_dir(): | |
| entries.append(f" 📁 {item.name}/") | |
| else: | |
| size = item.stat().st_size | |
| if size < 1024: | |
| size_str = f"{size}B" | |
| elif size < 1024 * 1024: | |
| size_str = f"{size/1024:.1f}KB" | |
| else: | |
| size_str = f"{size/(1024*1024):.1f}MB" | |
| entries.append(f" 📄 {item.name} ({size_str})") | |
| return f"Contents of {arguments.get('path', '.')}:\n" + "\n".join(entries) if entries else "Empty directory" | |
| elif name == "grep_search": | |
| pattern = arguments["pattern"] | |
| search_path = arguments.get("path", ".") | |
| path = _safe_path(search_path) | |
| # Escape single quotes in pattern for safety | |
| escaped_pattern = pattern.replace("'", "'\\''") | |
| process = await asyncio.create_subprocess_shell( | |
| f"grep -rn --include=* '{escaped_pattern}' '{path}'", | |
| stdout=asyncio.subprocess.PIPE, | |
| stderr=asyncio.subprocess.PIPE, | |
| cwd=WORKSPACE_DIR, | |
| ) | |
| try: | |
| stdout, stderr = await asyncio.wait_for(process.communicate(), timeout=10) | |
| output = stdout.decode('utf-8', errors='replace') if stdout else "No matches found" | |
| except asyncio.TimeoutError: | |
| try: | |
| process.kill() | |
| except Exception: | |
| pass | |
| await process.communicate() | |
| return "Error: Grep search timed out" | |
| if len(output) > 10000: | |
| output = output[:10000] + "\n\n[Truncated]" | |
| return output | |
| else: | |
| return f"Error: Unknown tool: {name}" | |
| except Exception as e: | |
| return f"Error executing {name}: {str(e)}" | |
| # --------------------------------------------------------------------------- | |
| # Database (Session Persistence) | |
| # --------------------------------------------------------------------------- | |
| db_pool: Optional[asyncpg.Pool] = None | |
| async def init_db(): | |
| """Initialize database connection pool and create tables.""" | |
| global db_pool | |
| if not DATABASE_URL: | |
| return | |
| try: | |
| db_pool = await asyncpg.create_pool( | |
| DATABASE_URL, | |
| ssl="require", | |
| min_size=1, | |
| max_size=3, | |
| max_inactive_connection_lifetime=300 | |
| ) | |
| async with db_pool.acquire() as conn: | |
| await conn.execute(""" | |
| CREATE TABLE IF NOT EXISTS agent_session_entries ( | |
| id BIGSERIAL PRIMARY KEY, | |
| project_key TEXT NOT NULL, | |
| session_id TEXT NOT NULL, | |
| subpath TEXT, | |
| entry JSONB NOT NULL, | |
| created_at TIMESTAMPTZ NOT NULL DEFAULT now() | |
| ); | |
| CREATE INDEX IF NOT EXISTS idx_session_key ON agent_session_entries (project_key, session_id, subpath, id); | |
| CREATE INDEX IF NOT EXISTS idx_project_session ON agent_session_entries (project_key, session_id); | |
| CREATE TABLE IF NOT EXISTS eternity_system_state ( | |
| id SERIAL PRIMARY KEY, | |
| goal TEXT NOT NULL, | |
| deadline TIMESTAMPTZ NOT NULL, | |
| current_mode VARCHAR(20) NOT NULL DEFAULT 'build', | |
| roadmap JSONB DEFAULT '[]', | |
| latest_brief TEXT, | |
| created_at TIMESTAMPTZ NOT NULL DEFAULT now() | |
| ); | |
| """) | |
| except Exception as e: | |
| print(f"[DB] Warning: Could not initialize database: {e}") | |
| db_pool = None | |
| async def save_message(session_id: str, role: str, content: str = None, | |
| tool_calls: list = None, tool_call_id: str = None): | |
| """Save a message to the session store using the unified schema.""" | |
| if not db_pool: | |
| return | |
| msg = {"role": role} | |
| if content is not None: | |
| msg["content"] = content | |
| if tool_calls: | |
| msg["tool_calls"] = tool_calls | |
| if tool_call_id: | |
| msg["tool_call_id"] = tool_call_id | |
| try: | |
| async with db_pool.acquire() as conn: | |
| await conn.execute( | |
| "INSERT INTO agent_session_entries (project_key, session_id, subpath, entry) VALUES ($1, $2, $3, $4)", | |
| "fastapi-completions", | |
| session_id, | |
| None, | |
| json.dumps(msg) | |
| ) | |
| except Exception as e: | |
| print(f"[DB] Warning: Could not save message: {e}") | |
| async def load_session(session_id: str) -> list: | |
| """Load conversation history from the session store using the unified schema.""" | |
| if not db_pool: | |
| return [] | |
| try: | |
| async with db_pool.acquire() as conn: | |
| rows = await conn.fetch( | |
| "SELECT entry FROM agent_session_entries WHERE project_key = $1 AND session_id = $2 AND subpath IS NOT DISTINCT FROM $3 ORDER BY id", | |
| "fastapi-completions", | |
| session_id, | |
| None | |
| ) | |
| return [json.loads(row["entry"]) for row in rows] | |
| except Exception as e: | |
| print(f"[DB] Warning: Could not load session: {e}") | |
| return [] | |
| # --------------------------------------------------------------------------- | |
| # SSE Chunk Formatting (OpenAI delta format) | |
| # --------------------------------------------------------------------------- | |
| def make_chunk(request_id: str, model: str, content: str = "", finish_reason: str = None) -> str: | |
| """Create an OpenAI-compatible SSE chunk.""" | |
| delta = {} | |
| if content: | |
| delta["content"] = content | |
| if finish_reason and not content: | |
| delta = {} | |
| chunk = { | |
| "id": f"chatcmpl-{request_id}", | |
| "object": "chat.completion.chunk", | |
| "created": int(time.time()), | |
| "model": model, | |
| "choices": [{ | |
| "index": 0, | |
| "delta": delta, | |
| "finish_reason": finish_reason, | |
| }], | |
| } | |
| return f"data: {json.dumps(chunk)}\n\n" | |
| # --------------------------------------------------------------------------- | |
| # FastAPI Application | |
| # --------------------------------------------------------------------------- | |
| app = FastAPI(title="Claude Code Backend", version="1.0.0") | |
| app.add_middleware( | |
| CORSMiddleware, | |
| allow_origins=["*"], | |
| allow_methods=["*"], | |
| allow_headers=["*"], | |
| ) | |
| def auth(authorization: str = None): | |
| """Verify bearer token.""" | |
| if not BACKEND_API_KEY: | |
| return # No auth configured | |
| expected = f"Bearer {BACKEND_API_KEY}" | |
| if authorization != expected: | |
| raise HTTPException(status_code=401, detail="Unauthorized") | |
| async def check_models_health(): | |
| global RECOMMENDED_MODEL | |
| # Test only the unstable frontier models (the rest). | |
| # The stable ones (Step 3.7 Flash, Nemotron 3 Ultra, Qwen 2.5 Coder) are always free/working. | |
| models_to_test = [ | |
| "moonshotai/kimi-k2.6", | |
| "z-ai/glm-5.1", | |
| "minimaxai/minimax-m3", | |
| "minimaxai/minimax-m2.7", | |
| "meta/llama-3.1-405b-instruct", | |
| ] | |
| best_model = None | |
| best_latency = 999.0 | |
| # Mark stable models as permanently ONLINE in the status map | |
| stable_models = [ | |
| "stepfun-ai/step-3.7-flash", | |
| "nvidia/nemotron-3-ultra-550b-a55b", | |
| "qwen/qwen2.5-coder-32b-instruct" | |
| ] | |
| for model in stable_models: | |
| MODEL_STATUSES[model] = {"status": "ONLINE (Stable)", "latency": "Fast", "raw_latency": 0.1} | |
| log_activity("Periodic health check started: verifying unstable frontier NIM models...") | |
| for model in models_to_test: | |
| start_time = time.time() | |
| try: | |
| # Send a fast test prompt | |
| async with anyio.fail_after(15.0): # 15 seconds max timeout | |
| await nim_client.chat.completions.create( | |
| model=model, | |
| messages=[{"role": "user", "content": "1+1="}], | |
| max_tokens=3, | |
| ) | |
| latency = time.time() - start_time | |
| MODEL_STATUSES[model] = {"status": "ONLINE", "latency": f"{latency:.2f}s", "raw_latency": latency} | |
| log_activity(f"Model checked: {model} is ONLINE ({latency:.2f}s)") | |
| # Choose the fastest online unstable model | |
| if latency < best_latency: | |
| best_latency = latency | |
| best_model = model | |
| except Exception as e: | |
| MODEL_STATUSES[model] = {"status": "OFFLINE", "latency": "N/A", "raw_latency": 999.0} | |
| log_activity(f"Model checked: {model} is OFFLINE / TIMEOUT: {e}") | |
| if best_model: | |
| RECOMMENDED_MODEL = best_model | |
| log_activity(f"Best frontier model selected: {RECOMMENDED_MODEL} ({best_latency:.2f}s)") | |
| else: | |
| # Fallback to the stable Step 3.7 Flash if all frontier models are offline/throttled | |
| RECOMMENDED_MODEL = "stepfun-ai/step-3.7-flash" | |
| log_activity(f"All frontier models offline. Falling back to stable recommended model: {RECOMMENDED_MODEL}") | |
| async def periodic_health_check_loop(): | |
| # Wait 10 seconds after startup before the first check to let the space boot fully | |
| await asyncio.sleep(10) | |
| while True: | |
| try: | |
| await check_models_health() | |
| except Exception as e: | |
| log_activity(f"Health check loop error: {e}") | |
| await asyncio.sleep(900) # every 15 minutes (reduce frequency to save quota) | |
| async def startup(): | |
| await init_db() | |
| Path(WORKSPACE_DIR).mkdir(parents=True, exist_ok=True) | |
| # Initialize statuses for all models | |
| for model_id, display_name in ALL_MODELS.items(): | |
| MODEL_STATUSES[model_id] = {"status": "UNCHECKED", "latency": "N/A", "raw_latency": 999.0} | |
| # Start background health checking | |
| asyncio.create_task(periodic_health_check_loop()) | |
| log_activity(f"FastAPI backend started. Workspace: {WORKSPACE_DIR}") | |
| # --------------------------------------------------------------------------- | |
| # /v1/chat/completions — Main endpoint | |
| # --------------------------------------------------------------------------- | |
| AGENTIC_SYSTEM_PROMPT = """You are an expert coding assistant with access to tools for file operations and command execution. | |
| When the user asks you to create, edit, or debug code: | |
| 1. Use `list_directory` and `read_file` to understand the current state | |
| 2. Use `write_file` to create or modify files | |
| 3. Use `run_bash` to execute commands (install packages, run scripts, test code) | |
| 4. Use `grep_search` to find patterns in code | |
| IMPORTANT RULES: | |
| - Always use tools to take action. Do NOT just describe what to do — actually DO it. | |
| - After writing code, run it to verify it works. | |
| - If a command fails, read the error and fix it. | |
| - Work in the /tmp/workspace directory. | |
| - Be concise in your explanations, but thorough in your tool usage. | |
| """ | |
| def compact_history(messages: list) -> list: | |
| """ | |
| Compact conversation history to prevent context window overflow. | |
| Replaces massive tool call outputs with concise summaries if history is long. | |
| """ | |
| # Only compact if messages count exceeds 15 (to maintain normal conversation) | |
| if len(messages) <= 15: | |
| return messages | |
| compacted = [] | |
| # Always keep the system prompt (typically the first message) | |
| if messages and messages[0].get("role") == "system": | |
| compacted.append(messages[0]) | |
| start_idx = 1 | |
| else: | |
| start_idx = 0 | |
| # Keep the last 4 messages exactly as they are to preserve immediate context | |
| recent_count = 4 | |
| mid_messages = messages[start_idx:-recent_count] | |
| recent_messages = messages[-recent_count:] | |
| for msg in mid_messages: | |
| role = msg.get("role") | |
| content = msg.get("content") or "" | |
| if role == "tool": | |
| # Compress massive tool outputs (like bash stdout or file reads) | |
| if len(content) > 1000: | |
| summary = f"[Tool output compacted: {content[:200]}... (Total {len(content)} chars truncated for context preservation)]" | |
| compacted.append({ | |
| "role": "tool", | |
| "tool_call_id": msg.get("tool_call_id"), | |
| "content": summary | |
| }) | |
| continue | |
| elif role == "assistant" and msg.get("tool_calls"): | |
| # Keep tool calls metadata so the model's message-tool call mapping DAG doesn't break | |
| pass | |
| # Keep general messages, but truncate if they are too long | |
| if len(content) > 2000: | |
| msg = dict(msg) | |
| msg["content"] = content[:2000] + "\n[Content truncated for compaction]" | |
| compacted.append(msg) | |
| compacted.extend(recent_messages) | |
| log_activity(f"[Auto-Compaction] Compressed message history from {len(messages)} down to {len(compacted)}") | |
| return compacted | |
| async def chat_completions(request: Request, authorization: str = Header(None)): | |
| auth(authorization) | |
| body = await request.json() | |
| requested_model = body.get("model", "meta/llama-3.1-70b-instruct") | |
| messages = body.get("messages", []) | |
| stream = body.get("stream", False) | |
| session_id = body.get("session_id") or str(uuid.uuid4()) | |
| is_agentic = requested_model in TOOL_CAPABLE_MODELS | |
| request_id = str(uuid.uuid4())[:8] | |
| ACTIVE_SESSIONS.add(session_id) | |
| log_activity(f"Session [{session_id[:6]}] connected. Model: {requested_model}") | |
| # Build message history | |
| final_messages = [] | |
| # Add agentic system prompt for tool-capable models | |
| if is_agentic: | |
| # Check if there's already a system message | |
| has_system = any(m.get("role") == "system" for m in messages) | |
| if has_system: | |
| # Prepend agentic prompt to existing system message | |
| for m in messages: | |
| if m["role"] == "system": | |
| final_messages.append({ | |
| "role": "system", | |
| "content": AGENTIC_SYSTEM_PROMPT + "\n\nAdditional instructions:\n" + m["content"] | |
| }) | |
| else: | |
| final_messages.append(m) | |
| else: | |
| final_messages.append({"role": "system", "content": AGENTIC_SYSTEM_PROMPT}) | |
| final_messages.extend(messages) | |
| else: | |
| final_messages = list(messages) | |
| # Perform auto-compaction before executing agent loops | |
| final_messages = compact_history(final_messages) | |
| # Save the user's message to DB | |
| user_msg = next((m for m in reversed(messages) if m.get("role") == "user"), None) | |
| if user_msg: | |
| await save_message(session_id, "user", user_msg.get("content", "")) | |
| if not stream: | |
| # Non-streaming: simple completion | |
| try: | |
| kwargs = {"model": requested_model, "messages": final_messages} | |
| if is_agentic: | |
| kwargs["tools"] = TOOLS | |
| kwargs["tool_choice"] = "auto" | |
| async with completions_semaphore: | |
| response = await nim_client.chat.completions.create(**kwargs) | |
| content = response.choices[0].message.content or "" | |
| await save_message(session_id, "assistant", content) | |
| ACTIVE_SESSIONS.discard(session_id) | |
| log_activity(f"Session [{session_id[:6]}] finished (non-streaming)") | |
| return JSONResponse({ | |
| "id": f"chatcmpl-{request_id}", | |
| "object": "chat.completion", | |
| "created": int(time.time()), | |
| "model": requested_model, | |
| "choices": [{"index": 0, "message": {"role": "assistant", "content": content}, "finish_reason": "stop"}], | |
| }) | |
| except Exception as e: | |
| ACTIVE_SESSIONS.discard(session_id) | |
| return JSONResponse({"error": {"message": str(e), "type": "internal_error"}}, status_code=500) | |
| # Streaming + agentic loop | |
| async def generate() -> AsyncIterator[str]: | |
| nonlocal final_messages | |
| async with completions_semaphore: | |
| try: | |
| for round_num in range(MAX_TOOL_ROUNDS + 1): | |
| # Perform auto-compaction before calling NIM API | |
| final_messages = compact_history(final_messages) | |
| kwargs = {"model": requested_model, "messages": final_messages, "stream": True} | |
| if is_agentic: | |
| kwargs["tools"] = TOOLS | |
| kwargs["tool_choice"] = "auto" | |
| # Collect streamed response | |
| full_content = "" | |
| tool_calls_raw = {} # index -> {id, name, arguments_str} | |
| async for chunk in await nim_client.chat.completions.create(**kwargs): | |
| choice = chunk.choices[0] if chunk.choices else None | |
| if not choice: | |
| continue | |
| delta = choice.delta | |
| # Stream text content to client | |
| if delta and delta.content: | |
| full_content += delta.content | |
| yield make_chunk(request_id, requested_model, delta.content) | |
| # Collect tool calls | |
| if delta and delta.tool_calls: | |
| for tc in delta.tool_calls: | |
| idx = tc.index | |
| if idx not in tool_calls_raw: | |
| tool_calls_raw[idx] = { | |
| "id": tc.id or f"call_{uuid.uuid4().hex[:8]}", | |
| "name": tc.function.name if tc.function and tc.function.name else "", | |
| "arguments": "" | |
| } | |
| if tc.function and tc.function.name: | |
| tool_calls_raw[idx]["name"] = tc.function.name | |
| if tc.id: | |
| tool_calls_raw[idx]["id"] = tc.id | |
| if tc.function and tc.function.arguments: | |
| tool_calls_raw[idx]["arguments"] += tc.function.arguments | |
| # Check for finish | |
| if choice.finish_reason == "stop": | |
| break | |
| if choice.finish_reason == "tool_calls": | |
| break | |
| # If no tool calls, we're done | |
| if not tool_calls_raw: | |
| await save_message(session_id, "assistant", full_content) | |
| yield make_chunk(request_id, requested_model, finish_reason="stop") | |
| yield "data: [DONE]\n\n" | |
| return | |
| # Execute tool calls | |
| tool_calls_list = [] | |
| for idx in sorted(tool_calls_raw.keys()): | |
| tc = tool_calls_raw[idx] | |
| tool_calls_list.append({ | |
| "id": tc["id"], | |
| "type": "function", | |
| "function": {"name": tc["name"], "arguments": tc["arguments"]} | |
| }) | |
| # Add assistant message with tool calls to history | |
| assistant_msg = {"role": "assistant", "content": full_content or None, "tool_calls": tool_calls_list} | |
| final_messages.append(assistant_msg) | |
| # Execute each tool and add results | |
| for tc in tool_calls_list: | |
| func_name = tc["function"]["name"] | |
| raw_args_str = tc["function"]["arguments"] | |
| try: | |
| func_args = json.loads(raw_args_str) | |
| except json.JSONDecodeError: | |
| # Attempt raw JSON repair | |
| repaired_str = raw_args_str.strip() | |
| if not repaired_str.startswith("{"): | |
| repaired_str = "{" + repaired_str | |
| if not repaired_str.endswith("}"): | |
| repaired_str = repaired_str + "}" | |
| try: | |
| func_args = json.loads(repaired_str) | |
| log_activity(f"Auto-fixed invalid JSON string for tool: {func_name}") | |
| except: | |
| func_args = {} | |
| # Perform semantic repairs | |
| repaired_args, repair_notes = repair_arguments(func_name, func_args) | |
| # Log activity | |
| log_activity(f"Tool execution: {func_name} args={repaired_args}") | |
| if repair_notes: | |
| for note in repair_notes: | |
| log_activity(f"[Tool Repair] {note}") | |
| # Show tool execution to user | |
| yield make_chunk(request_id, requested_model, f"\n\n🔧 **{func_name}**") | |
| if repair_notes: | |
| yield make_chunk(request_id, requested_model, " *(Auto-Repaired)*") | |
| if func_name == "run_bash" and "command" in repaired_args: | |
| yield make_chunk(request_id, requested_model, f": `{repaired_args['command']}`\n") | |
| elif func_name == "read_file" and "path" in repaired_args: | |
| yield make_chunk(request_id, requested_model, f": `{repaired_args['path']}`\n") | |
| elif func_name == "write_file" and "path" in repaired_args: | |
| yield make_chunk(request_id, requested_model, f": `{repaired_args['path']}`\n") | |
| elif func_name == "list_directory": | |
| yield make_chunk(request_id, requested_model, f": `{repaired_args.get('path', '.')}`\n") | |
| elif func_name == "grep_search": | |
| yield make_chunk(request_id, requested_model, f": `{repaired_args.get('pattern', '')}`\n") | |
| else: | |
| yield make_chunk(request_id, requested_model, "\n") | |
| # Execute the tool | |
| result = await execute_tool(func_name, repaired_args) | |
| # Append teaching note if repaired | |
| if repair_notes: | |
| result += f"\n\n[SYSTEM REPAIR NOTE: The harness automatically fixed formatting issues: {', '.join(repair_notes)}. Please strictly follow the tool's JSON schema in subsequent calls without these wrapping/formatting errors.]" | |
| # Show truncated result to user | |
| preview = result[:500] + ("..." if len(result) > 500 else "") | |
| yield make_chunk(request_id, requested_model, f"```\n{preview}\n```\n") | |
| # Add tool result to message history | |
| final_messages.append({ | |
| "role": "tool", | |
| "tool_call_id": tc["id"], | |
| "content": result, | |
| }) | |
| await save_message(session_id, "tool", result, tool_call_id=tc["id"]) | |
| # Continue the agentic loop (model processes tool results) | |
| # If we hit max rounds, finish | |
| yield make_chunk(request_id, requested_model, "\n\n⚠️ Reached maximum tool call rounds.") | |
| yield make_chunk(request_id, requested_model, finish_reason="stop") | |
| yield "data: [DONE]\n\n" | |
| except Exception as e: | |
| error_msg = f"\n\n❌ Error: {str(e)}" | |
| yield make_chunk(request_id, requested_model, error_msg) | |
| yield make_chunk(request_id, requested_model, finish_reason="stop") | |
| yield "data: [DONE]\n\n" | |
| return StreamingResponse( | |
| generate(), | |
| media_type="text/event-stream", | |
| headers={ | |
| "Cache-Control": "no-cache", | |
| "X-Accel-Buffering": "no", | |
| "Connection": "keep-alive", | |
| }, | |
| ) | |
| # --------------------------------------------------------------------------- | |
| # /v1/models — Model listing | |
| # --------------------------------------------------------------------------- | |
| async def list_models(authorization: str = Header(None)): | |
| auth(authorization) | |
| models = [] | |
| for model_id, display_name in ALL_MODELS.items(): | |
| models.append({ | |
| "id": model_id, | |
| "object": "model", | |
| "created": 1700000000, | |
| "owned_by": "nvidia-nim", | |
| "permission": [], | |
| "root": model_id, | |
| "parent": None, | |
| }) | |
| return {"object": "list", "data": models} | |
| # --------------------------------------------------------------------------- | |
| # /health — Health check | |
| # --------------------------------------------------------------------------- | |
| async def get_workspace_tree(): | |
| def build_tree(current_path: Path, relative_to: Path) -> dict: | |
| name = current_path.name | |
| try: | |
| rel_path = str(current_path.relative_to(relative_to)).replace("\\", "/") | |
| except ValueError: | |
| rel_path = "" | |
| if rel_path == ".": | |
| rel_path = "" | |
| if current_path.is_dir(): | |
| children = [] | |
| try: | |
| for child in sorted(current_path.iterdir(), key=lambda x: (not x.is_dir(), x.name)): | |
| if child.name in [".git", "node_modules", ".next", "__pycache__", ".agents", ".gemini"]: | |
| continue | |
| children.append(build_tree(child, relative_to)) | |
| except Exception: | |
| pass | |
| return { | |
| "name": name or "workspace", | |
| "path": rel_path, | |
| "type": "directory", | |
| "children": children | |
| } | |
| else: | |
| return { | |
| "name": name, | |
| "path": rel_path, | |
| "type": "file", | |
| "size": current_path.stat().st_size if current_path.exists() else 0 | |
| } | |
| try: | |
| w_path = Path(WORKSPACE_DIR).resolve() | |
| if not w_path.exists(): | |
| w_path.mkdir(parents=True, exist_ok=True) | |
| return build_tree(w_path, w_path) | |
| except Exception as e: | |
| return {"error": str(e)} | |
| async def get_workspace_file(path: str): | |
| try: | |
| safe_p = _safe_path(path) | |
| if not safe_p.exists() or not safe_p.is_file(): | |
| raise HTTPException(status_code=404, detail="File not found") | |
| content = safe_p.read_text(encoding="utf-8", errors="replace") | |
| return {"path": path, "content": content} | |
| except Exception as e: | |
| raise HTTPException(status_code=500, detail=str(e)) | |
| # --------------------------------------------------------------------------- | |
| # Dashboard and Status API | |
| # --------------------------------------------------------------------------- | |
| DASHBOARD_HTML = """ | |
| <!DOCTYPE html> | |
| <html lang="en"> | |
| <head> | |
| <meta charset="UTF-8"> | |
| <meta name="viewport" content="width=device-width, initial-scale=1.0"> | |
| <title>Claude Code Agent Console</title> | |
| <script src="https://cdn.tailwindcss.com"></script> | |
| <style> | |
| @import url('https://fonts.googleapis.com/css2?family=Fira+Code:wght@400;500;700&family=Outfit:wght@400;600;800&display=swap'); | |
| body { | |
| font-family: 'Outfit', sans-serif; | |
| background-color: #0b0c10; | |
| } | |
| .code-font { | |
| font-family: 'Fira Code', monospace; | |
| } | |
| .glow-amber { | |
| box-shadow: 0 0 15px rgba(245, 158, 11, 0.2); | |
| } | |
| </style> | |
| </head> | |
| <body class="text-gray-100 min-h-screen flex flex-col pb-10"> | |
| <header class="border-b border-gray-800 bg-gray-950/80 backdrop-blur px-6 py-4 flex items-center justify-between sticky top-0 z-50"> | |
| <div class="flex items-center space-x-3"> | |
| <span class="text-2xl font-extrabold tracking-tight bg-gradient-to-r from-blue-400 via-indigo-400 to-purple-400 bg-clip-text text-transparent"> | |
| Claude Code Agent Console | |
| </span> | |
| <span class="px-2 py-0.5 text-xs rounded bg-blue-500/10 text-blue-400 border border-blue-500/20 font-semibold animate-pulse"> | |
| LIVE | |
| </span> | |
| </div> | |
| <div class="flex items-center space-x-4 text-sm text-gray-400"> | |
| <div>Workspace: <span class="text-gray-200 code-font">/tmp/workspace</span></div> | |
| <div class="h-4 w-px bg-gray-800"></div> | |
| <div>Active Sessions: <span id="active-sessions-count" class="text-blue-400 font-bold code-font">0</span></div> | |
| </div> | |
| </header> | |
| <!-- Navigation Tabs --> | |
| <div class="border-b border-gray-800 max-w-7xl w-full mx-auto px-6 mt-6 flex space-x-6 text-sm"> | |
| <button onclick="switchTab('models')" id="tab-btn-models" class="pb-3 border-b-2 border-blue-500 font-semibold text-blue-400 transition-all">NIM Models</button> | |
| <button onclick="switchTab('logs')" id="tab-btn-logs" class="pb-3 border-b-2 border-transparent text-gray-400 hover:text-gray-200 font-semibold transition-all">Live Logs</button> | |
| <button onclick="switchTab('explorer')" id="tab-btn-explorer" class="pb-3 border-b-2 border-transparent text-gray-400 hover:text-gray-200 font-semibold flex items-center space-x-1 transition-all"> | |
| <span>Workspace Explorer (IDE)</span> | |
| <span class="px-1.5 py-0.5 rounded bg-blue-500/10 text-blue-400 border border-blue-500/20 text-[10px] font-bold">VS Code View</span> | |
| </button> | |
| </div> | |
| <!-- MAIN SECTIONS --> | |
| <main class="max-w-7xl w-full mx-auto px-6 mt-8 flex-1"> | |
| <!-- SECTION: Models --> | |
| <div id="section-models" class="space-y-6"> | |
| <div class="flex items-center justify-between"> | |
| <h2 class="text-lg font-bold tracking-tight text-gray-300">Nvidia NIM Models & Health Status</h2> | |
| <span class="text-xs text-gray-500">Checked every 15 mins</span> | |
| </div> | |
| <div id="models-container" class="grid grid-cols-1 md:grid-cols-2 lg:grid-cols-3 gap-4"> | |
| <!-- Dynamically loaded models go here --> | |
| </div> | |
| </div> | |
| <!-- SECTION: Logs --> | |
| <div id="section-logs" class="hidden space-y-6"> | |
| <h2 class="text-lg font-bold tracking-tight text-gray-300">System Activity Logs</h2> | |
| <div class="border border-gray-800 rounded-lg overflow-hidden bg-gray-950 flex flex-col min-h-[500px]"> | |
| <div class="bg-gray-900 px-4 py-2 border-b border-gray-800 flex items-center justify-between"> | |
| <span class="text-xs text-gray-400 font-semibold code-font">agent-stdout.log</span> | |
| <div class="flex space-x-1.5"> | |
| <span class="w-2.5 h-2.5 rounded-full bg-red-500/30"></span> | |
| <span class="w-2.5 h-2.5 rounded-full bg-yellow-500/30"></span> | |
| <span class="w-2.5 h-2.5 rounded-full bg-green-500/30"></span> | |
| </div> | |
| </div> | |
| <div id="terminal-content" class="p-4 flex-1 overflow-y-auto code-font text-xs text-green-400 bg-black/90 space-y-1 select-all h-[450px]"> | |
| <!-- Logs go here --> | |
| </div> | |
| </div> | |
| </div> | |
| <!-- SECTION: Workspace Explorer --> | |
| <div id="section-explorer" class="hidden space-y-6"> | |
| <div class="flex items-center justify-between"> | |
| <h2 class="text-lg font-bold tracking-tight text-gray-300">Visual Workspace IDE</h2> | |
| <button onclick="refreshFileTree()" class="text-xs px-2.5 py-1 rounded bg-blue-500/10 text-blue-400 border border-blue-500/20 hover:bg-blue-500/20 transition-all font-semibold"> | |
| 🔄 Refresh Tree | |
| </button> | |
| </div> | |
| <div class="grid grid-cols-1 md:grid-cols-3 gap-6 border border-gray-800 rounded-xl bg-gray-950 overflow-hidden h-[600px]"> | |
| <!-- File Tree Sidebar --> | |
| <div class="border-r border-gray-800 flex flex-col bg-gray-950 h-full"> | |
| <div class="px-4 py-2 border-b border-gray-800 bg-gray-900 text-xs font-semibold tracking-wider text-gray-400 code-font"> | |
| 📁 EXPLORER: WORKSPACE | |
| </div> | |
| <div id="file-tree" class="p-3 flex-1 overflow-y-auto space-y-0.5 select-none"> | |
| <!-- Tree will be loaded here --> | |
| <span class="text-xs text-gray-500 italic px-2">Loading directory tree...</span> | |
| </div> | |
| </div> | |
| <!-- Editor panel --> | |
| <div class="md:col-span-2 flex flex-col bg-black/40 h-full"> | |
| <div class="px-4 py-2 border-b border-gray-800 bg-gray-900 flex items-center justify-between"> | |
| <span id="editor-title" class="text-xs font-semibold text-gray-400 code-font">📄 Welcome screen</span> | |
| <div class="flex space-x-1.5"> | |
| <span class="w-2 h-2 rounded-full bg-gray-700"></span> | |
| <span class="w-2 h-2 rounded-full bg-gray-700"></span> | |
| </div> | |
| </div> | |
| <div class="flex-1 p-4 overflow-auto code-font text-xs text-gray-200"> | |
| <pre id="editor-content" class="whitespace-pre overflow-x-auto select-text h-[500px]"> | |
| Welcome to Claude Code Workspace Explorer. | |
| Select a file from the sidebar explorer on the left to read its code contents in real-time. | |
| </pre> | |
| </div> | |
| </div> | |
| </div> | |
| </div> | |
| </main> | |
| <script> | |
| let currentTab = 'models'; | |
| function switchTab(tabId) { | |
| currentTab = tabId; | |
| // Toggle sections | |
| document.getElementById('section-models').classList.add('hidden'); | |
| document.getElementById('section-logs').classList.add('hidden'); | |
| document.getElementById('section-explorer').classList.add('hidden'); | |
| document.getElementById('tab-btn-models').className = 'pb-3 border-b-2 border-transparent text-gray-400 hover:text-gray-200 font-semibold transition-all'; | |
| document.getElementById('tab-btn-logs').className = 'pb-3 border-b-2 border-transparent text-gray-400 hover:text-gray-200 font-semibold transition-all'; | |
| document.getElementById('tab-btn-explorer').className = 'pb-3 border-b-2 border-transparent text-gray-400 hover:text-gray-200 font-semibold flex items-center space-x-1 transition-all'; | |
| if (tabId === 'models') { | |
| document.getElementById('section-models').classList.remove('hidden'); | |
| document.getElementById('tab-btn-models').className = 'pb-3 border-b-2 border-blue-500 font-semibold text-blue-400 transition-all'; | |
| } else if (tabId === 'logs') { | |
| document.getElementById('section-logs').classList.remove('hidden'); | |
| document.getElementById('tab-btn-logs').className = 'pb-3 border-b-2 border-blue-500 font-semibold text-blue-400 transition-all'; | |
| } else if (tabId === 'explorer') { | |
| document.getElementById('section-explorer').classList.remove('hidden'); | |
| document.getElementById('tab-btn-explorer').className = 'pb-3 border-b-2 border-blue-500 font-semibold text-blue-400 flex items-center space-x-1 transition-all'; | |
| refreshFileTree(); | |
| } | |
| } | |
| async function fetchSystemData() { | |
| try { | |
| const res = await fetch('/health'); | |
| if (!res.ok) return; | |
| const data = await res.json(); | |
| document.getElementById('active-sessions-count').innerText = data.active_sessions || 0; | |
| } catch (e) { | |
| console.error(e); | |
| } | |
| } | |
| async function fetchModels() { | |
| try { | |
| const res = await fetch('/api/models-status'); | |
| if (!res.ok) return; | |
| const models = await res.json(); | |
| const container = document.getElementById('models-container'); | |
| container.innerHTML = ''; | |
| models.forEach(model => { | |
| const isRec = model.is_recommended; | |
| const isOnline = model.status.includes('ONLINE'); | |
| const card = document.createElement('div'); | |
| card.className = `p-4 border rounded-xl bg-gray-950 transition-all ${ | |
| isRec ? 'border-amber-500/50 glow-amber bg-amber-500/5' : 'border-gray-800 bg-gray-950' | |
| }`; | |
| card.innerHTML = ` | |
| <div class="flex items-center justify-between mb-3"> | |
| <span class="text-xs text-gray-500 code-font truncate max-w-[200px]" title="${model.id}">${model.id}</span> | |
| <div class="flex items-center space-x-2"> | |
| ${isRec ? '<span class="text-[10px] px-1.5 py-0.5 rounded bg-amber-500/10 text-amber-400 border border-amber-500/20 font-bold">★ Recommended</span>' : ''} | |
| <span class="h-2 w-2 rounded-full ${isOnline ? 'bg-green-500 animate-pulse' : 'bg-red-500'}"></span> | |
| <span class="text-[10px] font-bold ${isOnline ? 'text-green-400' : 'text-red-400'}">${model.status}</span> | |
| </div> | |
| </div> | |
| <h3 class="text-sm font-bold text-gray-200 mb-2 truncate">${model.name}</h3> | |
| <div class="flex items-center justify-between text-xs text-gray-400 border-t border-gray-900 pt-2"> | |
| <span>Type: <strong class="text-gray-300 font-medium">${model.type}</strong></span> | |
| <span>Latency: <strong class="text-blue-400 code-font">${model.latency}</strong></span> | |
| </div> | |
| `; | |
| container.appendChild(card); | |
| }); | |
| } catch (e) { | |
| console.error(e); | |
| } | |
| } | |
| async function fetchLogs() { | |
| if (currentTab !== 'logs') return; | |
| try { | |
| const res = await fetch('/api/logs'); | |
| if (!res.ok) return; | |
| const logs = await res.json(); | |
| const term = document.getElementById('terminal-content'); | |
| const shouldScroll = term.scrollHeight - term.clientHeight <= term.scrollTop + 50; | |
| term.innerHTML = logs.map(line => `<div>${line}</div>`).join(''); | |
| if (shouldScroll) { | |
| term.scrollTop = term.scrollHeight; | |
| } | |
| } catch (e) { | |
| console.error(e); | |
| } | |
| } | |
| // File Explorer Logic | |
| async function refreshFileTree() { | |
| try { | |
| const res = await fetch('/api/workspace/tree'); | |
| if (!res.ok) return; | |
| const root = await res.json(); | |
| const container = document.getElementById('file-tree'); | |
| container.innerHTML = renderNode(root); | |
| } catch (e) { | |
| console.error(e); | |
| } | |
| } | |
| function renderNode(node, depth = 0) { | |
| const isDir = node.type === 'directory'; | |
| const icon = isDir ? '📁' : '📄'; | |
| const indent = depth * 12; | |
| let html = ` | |
| <div class="flex items-center py-1 px-2 hover:bg-gray-800 rounded cursor-pointer transition-all text-xs" | |
| style="padding-left: ${indent}px" | |
| onclick="${isDir ? `toggleDir('${node.path}')` : `openFile('${node.path}')`}"> | |
| <span class="mr-2">${icon}</span> | |
| <span class="truncate ${isDir ? 'text-gray-300 font-medium' : 'text-gray-400'}">${node.name}</span> | |
| </div> | |
| `; | |
| if (isDir && node.children && node.children.length > 0) { | |
| html += `<div id="dir-${node.path.replace(/\\/g, '-').replace(/\\//g, '-')}" class="space-y-0.5">`; | |
| node.children.forEach(child => { | |
| html += renderNode(child, depth + 1); | |
| }); | |
| html += `</div>`; | |
| } else if (isDir && (!node.children || node.children.length === 0)) { | |
| html += `<div class="text-[10px] text-gray-600 italic" style="padding-left: ${indent + 16}px">(empty)</div>`; | |
| } | |
| return html; | |
| } | |
| async function openFile(path) { | |
| document.getElementById('editor-title').innerText = `📄 ${path}`; | |
| document.getElementById('editor-content').innerText = "Loading file content..."; | |
| try { | |
| const res = await fetch(`/api/workspace/file?path=${encodeURIComponent(path)}`); | |
| if (!res.ok) { | |
| document.getElementById('editor-content').innerText = "Error: Failed to fetch file content."; | |
| return; | |
| } | |
| const data = await res.json(); | |
| document.getElementById('editor-content').innerText = data.content; | |
| } catch (e) { | |
| document.getElementById('editor-content').innerText = `Error: ${e.message}`; | |
| } | |
| } | |
| function toggleDir(path) { | |
| const safeId = `dir-${path.replace(/\\/g, '-').replace(/\\//g, '-')}`; | |
| const elem = document.getElementById(safeId); | |
| if (elem) { | |
| elem.classList.toggle('hidden'); | |
| } | |
| } | |
| setInterval(fetchSystemData, 3000); | |
| setInterval(fetchModels, 3000); | |
| setInterval(fetchLogs, 2000); | |
| fetchSystemData(); | |
| fetchModels(); | |
| fetchLogs(); | |
| </script> | |
| </body> | |
| </html> | |
| """ | |
| # --------------------------------------------------------------------------- | |
| # Concurrency Settings | |
| # --------------------------------------------------------------------------- | |
| completions_semaphore = asyncio.Semaphore(2) | |
| # --------------------------------------------------------------------------- | |
| # SessionStore API (Claude Agent SDK / Claude Code Compatible) | |
| # --------------------------------------------------------------------------- | |
| class SessionAppendRequest(BaseModel): | |
| project_key: str | |
| session_id: str | |
| subpath: Optional[str] = None | |
| entries: List[Dict[str, Any]] | |
| async def append_session_entries(req: SessionAppendRequest, authorization: str = Header(None)): | |
| auth(authorization) | |
| if not db_pool: | |
| raise HTTPException(status_code=500, detail="Database not connected") | |
| try: | |
| # Check for mirror_error and log/alert if present | |
| for e in req.entries: | |
| if e.get("type") == "system" and e.get("subtype") == "mirror_error": | |
| log_activity(f"[ALERT] Claude Agent SDK reported mirror_error: {e.get('message')}") | |
| rows = [(req.project_key, req.session_id, req.subpath, json.dumps(e)) for e in req.entries] | |
| async with db_pool.acquire() as conn: | |
| await conn.executemany( | |
| "INSERT INTO agent_session_entries (project_key, session_id, subpath, entry) VALUES ($1, $2, $3, $4)", | |
| rows | |
| ) | |
| return {"status": "success"} | |
| except Exception as e: | |
| log_activity(f"[SessionStore Error] Failed append: {e}") | |
| raise HTTPException(status_code=500, detail=str(e)) | |
| class SessionLoadRequest(BaseModel): | |
| project_key: str | |
| session_id: str | |
| subpath: Optional[str] = None | |
| async def load_session_entries(req: SessionLoadRequest, authorization: str = Header(None)): | |
| auth(authorization) | |
| if not db_pool: | |
| raise HTTPException(status_code=500, detail="Database not connected") | |
| try: | |
| async with db_pool.acquire() as conn: | |
| rows = await conn.fetch( | |
| "SELECT entry FROM agent_session_entries WHERE project_key=$1 AND session_id=$2 AND subpath IS NOT DISTINCT FROM $3 ORDER BY id", | |
| req.project_key, req.session_id, req.subpath | |
| ) | |
| return {"entries": [json.loads(r["entry"]) for r in rows]} | |
| except Exception as e: | |
| log_activity(f"[SessionStore Error] Failed load: {e}") | |
| raise HTTPException(status_code=500, detail=str(e)) | |
| class SessionListRequest(BaseModel): | |
| project_key: str | |
| async def list_sessions(req: SessionListRequest, authorization: str = Header(None)): | |
| auth(authorization) | |
| if not db_pool: | |
| raise HTTPException(status_code=500, detail="Database not connected") | |
| try: | |
| async with db_pool.acquire() as conn: | |
| rows = await conn.fetch( | |
| "SELECT DISTINCT session_id FROM agent_session_entries WHERE project_key=$1", | |
| req.project_key | |
| ) | |
| return {"sessions": [r["session_id"] for r in rows]} | |
| except Exception as e: | |
| raise HTTPException(status_code=500, detail=str(e)) | |
| class SessionDeleteRequest(BaseModel): | |
| project_key: str | |
| session_id: str | |
| async def delete_session(req: SessionDeleteRequest, authorization: str = Header(None)): | |
| auth(authorization) | |
| if not db_pool: | |
| raise HTTPException(status_code=500, detail="Database not connected") | |
| try: | |
| async with db_pool.acquire() as conn: | |
| await conn.execute( | |
| "DELETE FROM agent_session_entries WHERE project_key=$1 AND session_id=$2", | |
| req.project_key, req.session_id | |
| ) | |
| return {"status": "success"} | |
| except Exception as e: | |
| raise HTTPException(status_code=500, detail=str(e)) | |
| from datetime import datetime, timedelta, timezone | |
| class EternityInitRequest(BaseModel): | |
| goal: Optional[str] = None | |
| problem_statement: Optional[str] = None | |
| deadline_hours: float | |
| async def init_eternity_system(req: EternityInitRequest, authorization: str = Header(None)): | |
| auth(authorization) | |
| if not db_pool: | |
| raise HTTPException(status_code=500, detail="Database not connected") | |
| try: | |
| goal_text = req.goal or req.problem_statement or "Build a calculator" | |
| deadline = datetime.now(timezone.utc) + timedelta(hours=req.deadline_hours) | |
| async with db_pool.acquire() as conn: | |
| # Delete any old state | |
| await conn.execute("DELETE FROM eternity_system_state") | |
| # Insert new state | |
| await conn.execute( | |
| "INSERT INTO eternity_system_state (goal, deadline, current_mode, roadmap) VALUES ($1, $2, 'build', $3)", | |
| goal_text, deadline, "[]" | |
| ) | |
| log_activity(f"[Eternity Loop] System initialized with goal: '{goal_text}' | Deadline: {deadline}") | |
| return {"status": "success", "deadline": deadline.isoformat()} | |
| except Exception as e: | |
| log_activity(f"[Eternity Loop Error] Init failed: {e}") | |
| raise HTTPException(status_code=500, detail=str(e)) | |
| async def get_eternity_status(authorization: str = Header(None)): | |
| auth(authorization) | |
| if not db_pool: | |
| raise HTTPException(status_code=500, detail="Database not connected") | |
| try: | |
| async with db_pool.acquire() as conn: | |
| row = await conn.fetchrow("SELECT goal, deadline, current_mode, roadmap, latest_brief FROM eternity_system_state ORDER BY id DESC LIMIT 1") | |
| if not row: | |
| return {"active": False} | |
| deadline = row["deadline"] | |
| now = datetime.now(timezone.utc) | |
| remaining = max(0.0, (deadline - now).total_seconds()) | |
| return { | |
| "active": True, | |
| "goal": row["goal"], | |
| "deadline": deadline.isoformat(), | |
| "current_mode": row["current_mode"], | |
| "roadmap": json.loads(row["roadmap"]) if isinstance(row["roadmap"], str) else row["roadmap"], | |
| "latest_brief": row["latest_brief"], | |
| "time_remaining_seconds": remaining | |
| } | |
| except Exception as e: | |
| raise HTTPException(status_code=500, detail=str(e)) | |
| class ForgeExecuteRequest(BaseModel): | |
| task_id: str | |
| action: str | |
| prompt: str | |
| context_rules: Optional[str] = "" | |
| async def forge_execute(req: ForgeExecuteRequest, authorization: str = Header(None)): | |
| auth(authorization) | |
| log_activity(f"[Forge] Received execution request for task {req.task_id}: '{req.prompt}'") | |
| try: | |
| messages = [ | |
| {"role": "system", "content": AGENTIC_SYSTEM_PROMPT + f"\nContext rules: {req.context_rules}"}, | |
| {"role": "user", "content": req.prompt} | |
| ] | |
| final_messages = list(messages) | |
| success = False | |
| summary = "" | |
| error = None | |
| for round_num in range(MAX_TOOL_ROUNDS + 1): | |
| await rate_limiter.wait_for_nim() | |
| response = await nim_client.chat.completions.create( | |
| model=RECOMMENDED_MODEL, | |
| messages=final_messages, | |
| tools=TOOLS, | |
| tool_choice="auto" | |
| ) | |
| msg = response.choices[0].message | |
| tool_calls_payload = None | |
| if msg.tool_calls: | |
| tool_calls_payload = [] | |
| for tc in msg.tool_calls: | |
| tc_dict = { | |
| "id": tc.id, | |
| "type": tc.type, | |
| "function": { | |
| "name": tc.function.name, | |
| "arguments": tc.function.arguments | |
| } | |
| } | |
| tool_calls_payload.append(tc_dict) | |
| final_messages.append({ | |
| "role": "assistant", | |
| "content": msg.content, | |
| "tool_calls": tool_calls_payload | |
| }) | |
| if not msg.tool_calls: | |
| success = True | |
| summary = msg.content or "Task completed." | |
| break | |
| for tc in msg.tool_calls: | |
| func_name = tc.function.name | |
| func_args = json.loads(tc.function.arguments) | |
| result = await execute_tool(func_name, func_args) | |
| final_messages.append({ | |
| "role": "tool", | |
| "tool_call_id": tc.id, | |
| "content": result | |
| }) | |
| else: | |
| error = "Max tool call rounds exceeded." | |
| if success: | |
| return { | |
| "status": "success", | |
| "summary": summary, | |
| "error": None | |
| } | |
| else: | |
| return { | |
| "status": "error", | |
| "error": error or "Task execution failed." | |
| } | |
| except Exception as e: | |
| log_activity(f"[Forge Error] Task execution failed: {e}") | |
| return { | |
| "status": "error", | |
| "error": str(e) | |
| } | |
| # --------------------------------------------------------------------------- | |
| # Database Retention & Heartbeat loops | |
| # --------------------------------------------------------------------------- | |
| async def db_heartbeat_loop(): | |
| log_activity("Database Heartbeat task started") | |
| while True: | |
| try: | |
| if db_pool: | |
| async with db_pool.acquire() as conn: | |
| await conn.execute("SELECT 1") | |
| log_activity("[Heartbeat] Pinged Aiven PostgreSQL successfully") | |
| except Exception as e: | |
| log_activity(f"[Heartbeat Warning] Failed to ping database: {e}") | |
| await asyncio.sleep(240) # Every 4 minutes | |
| SPACE3_URL = os.environ.get("SPACE3_URL", "https://augment17-claude-code-backend.hf.space") | |
| SPACE4_URL = os.environ.get("SPACE4_URL", "https://shyota-mcp-cloud-host.hf.space") | |
| SPACE5_URL = os.environ.get("SPACE5_URL", "https://augment17-better-chatbot.hf.space") | |
| SPACE6_URL = os.environ.get("SPACE6_URL", "https://augment17-mcp-cloud-host.hf.space") | |
| async def get_db_state(): | |
| async with db_pool.acquire() as conn: | |
| return await conn.fetchrow("SELECT goal, deadline, current_mode FROM eternity_system_state ORDER BY id DESC LIMIT 1") | |
| async def update_db_mode(mode: str): | |
| async with db_pool.acquire() as conn: | |
| await conn.execute("UPDATE eternity_system_state SET current_mode = $1", mode) | |
| async def update_db_brief(brief: str): | |
| async with db_pool.acquire() as conn: | |
| await conn.execute("UPDATE eternity_system_state SET latest_brief = $1", brief) | |
| async def update_db_roadmap(roadmap: list): | |
| async with db_pool.acquire() as conn: | |
| await conn.execute("UPDATE eternity_system_state SET roadmap = $1", json.dumps(roadmap)) | |
| def post_json(url: str, payload: dict) -> dict: | |
| headers = { | |
| "Authorization": f"Bearer {BACKEND_API_KEY}", | |
| "Content-Type": "application/json" | |
| } | |
| req = urllib.request.Request( | |
| url, | |
| data=json.dumps(payload).encode('utf-8'), | |
| headers=headers, | |
| method="POST" | |
| ) | |
| try: | |
| with urllib.request.urlopen(req, timeout=120) as r: | |
| return json.loads(r.read().decode('utf-8')) | |
| except Exception as e: | |
| print(f"[HTTP Error] POST to {url} failed: {e}") | |
| return {"status": "error", "error": str(e)} | |
| async def execute_build_cycle(goal: str): | |
| log_activity(f"[Build Mode] Initiating build cycle for goal: '{goal}'") | |
| await rate_limiter.wait_for_nim() | |
| prompt = f"We are building: '{goal}'. Write a JSON instruction for Space 3 (The Forge) to code the next milestone. Respond ONLY with JSON matching the contract: " + '{"prompt": "task description", "context_rules": "rules"}' | |
| try: | |
| res = await nim_client.chat.completions.create( | |
| model="nvidia/llama-3.1-nemotron-70b-instruct", | |
| messages=[{"role": "user", "content": prompt}], | |
| max_tokens=300 | |
| ) | |
| task = json.loads(res.choices[0].message.content.strip()) | |
| task_prompt = task.get("prompt") | |
| context_rules = task.get("context_rules", "") | |
| except Exception as e: | |
| log_activity(f"[Build Mode Error] NIM planning failed: {e}") | |
| return | |
| log_activity(f"[Build Mode] Dispatching task to Space 3: '{task_prompt}'") | |
| forge_res = post_json(f"{SPACE3_URL}/api/forge/execute", { | |
| "task_id": f"build_{int(time.time())}", | |
| "action": "execute_code", | |
| "prompt": task_prompt, | |
| "context_rules": context_rules | |
| }) | |
| if forge_res.get("status") == "success": | |
| log_activity(f"[Build Mode] Space 3 success: {forge_res.get('summary')}") | |
| log_activity("[Build Mode] Triggering Space 6 (The Sandbox) UI verification...") | |
| test_res = post_json(f"{SPACE6_URL}/api/sandbox/test", { | |
| "test_cmd": "verify_ui", | |
| "url": f"{SPACE3_URL}" | |
| }) | |
| log_activity(f"[Build Mode] Space 6 Test Verdict: {test_res.get('verdict')} | Reason: {test_res.get('reason')}") | |
| log_activity("[Build Mode] Triggering Space 5 (The Vault) backup commit...") | |
| post_json(f"{SPACE5_URL}/api/vault/push", {}) | |
| else: | |
| log_activity(f"[Build Mode Warning] Space 3 reported failure: {forge_res.get('error')}") | |
| async def execute_eternity_cycle(goal: str): | |
| log_activity(f"[Eternity Mode] Initiating autonomous R&D cycle for goal: '{goal}'") | |
| log_activity("[Eternity Mode] Querying Space 4 (The Library) for new feature research...") | |
| research_res = post_json(f"{SPACE4_URL}/api/research", {"query": f"novel scientific or chemical calculation features and algorithms for {goal}"}) | |
| brief = research_res.get("brief", "No new features found.") | |
| await update_db_brief(brief) | |
| log_activity(f"[Eternity Mode] Received research brief: {brief[:100]}...") | |
| await rate_limiter.wait_for_nim() | |
| prompt = f"Goal: '{goal}'. Research Brief: '{brief}'. Plan the next feature/optimization code. Respond ONLY with JSON: " + '{"prompt": "task description", "context_rules": "rules"}' | |
| try: | |
| res = await nim_client.chat.completions.create( | |
| model="nvidia/llama-3.1-nemotron-70b-instruct", | |
| messages=[{"role": "user", "content": prompt}], | |
| max_tokens=300 | |
| ) | |
| task = json.loads(res.choices[0].message.content.strip()) | |
| task_prompt = task.get("prompt") | |
| context_rules = task.get("context_rules", "") | |
| except Exception as e: | |
| log_activity(f"[Eternity Mode Error] NIM planning failed: {e}") | |
| return | |
| log_activity(f"[Eternity Mode] Dispatching task to Space 3: '{task_prompt}'") | |
| forge_res = post_json(f"{SPACE3_URL}/api/forge/execute", { | |
| "task_id": f"eternity_{int(time.time())}", | |
| "action": "execute_code", | |
| "prompt": task_prompt, | |
| "context_rules": context_rules | |
| }) | |
| if forge_res.get("status") == "success": | |
| log_activity(f"[Eternity Mode] Space 3 success: {forge_res.get('summary')}") | |
| log_activity("[Eternity Mode] Triggering Space 5 (The Vault) backup commit...") | |
| post_json(f"{SPACE5_URL}/api/vault/push", {}) | |
| def run_eternity_loop(): | |
| log_activity("[Eternity Loop] Daemon thread started.") | |
| while True: | |
| try: | |
| if not db_pool: | |
| time.sleep(10) | |
| continue | |
| loop = asyncio.new_event_loop() | |
| state = loop.run_until_complete(get_db_state()) | |
| loop.close() | |
| if not state: | |
| time.sleep(30) | |
| continue | |
| goal = state["goal"] | |
| deadline = state["deadline"] | |
| current_mode = state["current_mode"] | |
| now = datetime.now(timezone.utc) | |
| if current_mode == "build" and now >= deadline: | |
| loop = asyncio.new_event_loop() | |
| loop.run_until_complete(update_db_mode("eternity")) | |
| loop.close() | |
| current_mode = "eternity" | |
| log_activity("[Eternity Loop] Deadline reached. Transitioned to Eternity R&D Mode.") | |
| if current_mode == "build": | |
| loop = asyncio.new_event_loop() | |
| loop.run_until_complete(execute_build_cycle(goal)) | |
| loop.close() | |
| time.sleep(300) | |
| else: | |
| loop = asyncio.new_event_loop() | |
| loop.run_until_complete(execute_eternity_cycle(goal)) | |
| loop.close() | |
| interval = int(os.environ.get("ETERNITY_LOOP_INTERVAL", "3600")) | |
| time.sleep(interval) | |
| except Exception as e: | |
| log_activity(f"[Eternity Loop Error] Loop crash: {e}") | |
| time.sleep(60) | |
| async def dashboard(): | |
| return HTMLResponse(content=DASHBOARD_HTML) | |
| async def get_logs(): | |
| return list(activity_logs) | |
| async def get_models_status(): | |
| status_list = [] | |
| for model_id, display_name in ALL_MODELS.items(): | |
| status_info = MODEL_STATUSES.get(model_id, {"status": "ONLINE (Unchecked)", "latency": "N/A"}) | |
| is_rec = model_id == RECOMMENDED_MODEL | |
| is_agentic = model_id in TOOL_CAPABLE_MODELS | |
| status_list.append({ | |
| "id": model_id, | |
| "name": display_name, | |
| "status": status_info["status"], | |
| "latency": status_info["latency"], | |
| "is_recommended": is_rec, | |
| "type": "Agentic (Tools)" if is_agentic else "Chat Only", | |
| }) | |
| # Sort: Recommended first, then Agentic, then Chat | |
| status_list.sort(key=lambda m: (not m["is_recommended"], m["type"] != "Agentic (Tools)", m["name"])) | |
| return status_list | |
| import shutil | |
| import threading | |
| import signal | |
| from fastapi.responses import FileResponse | |
| async def download_backup(authorization: str = Header(None)): | |
| auth(authorization) | |
| snapshot_dir = "/tmp/workspace_snapshot" | |
| archive_base = "/tmp/workspace_backup_download" | |
| archive_zip = archive_base + ".zip" | |
| # Clean up old files/folders | |
| for path in [snapshot_dir, archive_zip]: | |
| if os.path.exists(path): | |
| try: | |
| if os.path.isdir(path): | |
| shutil.rmtree(path) | |
| else: | |
| os.unlink(path) | |
| except Exception: | |
| pass | |
| try: | |
| # 1. Atomic-like snapshot copy (ignoring temporary files) | |
| shutil.copytree(WORKSPACE_DIR, snapshot_dir, symlinks=True, ignore=shutil.ignore_patterns('.git', 'node_modules', '.next')) | |
| # 2. Archive the snapshot folder to disk to prevent OOM memory spike | |
| shutil.make_archive(archive_base, 'zip', snapshot_dir) | |
| # 3. Clean up the snapshot directory immediately | |
| shutil.rmtree(snapshot_dir) | |
| if not os.path.exists(archive_zip): | |
| raise HTTPException(status_code=500, detail="Failed to create zip archive") | |
| return FileResponse(archive_zip, media_type="application/zip", filename="workspace_backup.zip") | |
| except Exception as e: | |
| if os.path.exists(snapshot_dir): | |
| shutil.rmtree(snapshot_dir) | |
| raise HTTPException(status_code=500, detail=str(e)) | |
| async def health(): | |
| return { | |
| "status": "ok", | |
| "workspace": WORKSPACE_DIR, | |
| "workspace_exists": Path(WORKSPACE_DIR).exists(), | |
| "db_connected": db_pool is not None, | |
| "models_count": len(ALL_MODELS), | |
| "recommended_model": RECOMMENDED_MODEL, | |
| "active_sessions": len(ACTIVE_SESSIONS), | |
| } | |
| # --------------------------------------------------------------------------- | |
| # Watchdog Daemon for Claude Code Subprocesses (Orphan Reaper) | |
| # --------------------------------------------------------------------------- | |
| def run_watchdog(): | |
| log_activity("System Watchdog Daemon started (PPID-based Orphan detection)") | |
| while True: | |
| try: | |
| import psutil | |
| for proc in psutil.process_iter(['pid', 'ppid', 'name', 'cmdline', 'status']): | |
| try: | |
| cmd = " ".join(proc.info['cmdline'] or []) | |
| # Match the CLI binary (looks like claude-code or anthropic CLI wrapper) | |
| if "claude" in cmd.lower() or "anthropic" in cmd.lower(): | |
| ppid = proc.info['ppid'] | |
| pid = proc.info['pid'] | |
| # Is the parent still alive and not a zombie? | |
| parent_exists = False | |
| if ppid != 1: # Orphaned processes get reparented to PID 1 in Linux | |
| try: | |
| parent_proc = psutil.Process(ppid) | |
| if parent_proc.is_running() and parent_proc.status() != psutil.STATUS_ZOMBIE: | |
| parent_exists = True | |
| except psutil.NoSuchProcess: | |
| pass | |
| if not parent_exists: | |
| log_activity(f"[Watchdog SIGKILL] Reaping orphaned Claude Code process PID {pid} (PPID {ppid})") | |
| proc.terminate() | |
| time.sleep(2) | |
| if proc.is_running(): | |
| proc.kill() | |
| except Exception: | |
| continue | |
| except ImportError: | |
| # Fallback zero-dependency shell parser using /proc | |
| try: | |
| # Find all processes and examine their parent PID | |
| out = subprocess.check_output("ps -o pid,ppid,args | grep -E 'claude|anthropic' | grep -v grep", shell=True, text=True) | |
| for line in out.strip().split("\n"): | |
| parts = line.strip().split(None, 2) | |
| if len(parts) >= 2: | |
| pid = int(parts[0]) | |
| ppid = int(parts[1]) | |
| # Check if parent pid exists/is alive | |
| parent_exists = False | |
| if ppid != 1: | |
| # Check /proc/[ppid] directory | |
| if os.path.exists(f"/proc/{ppid}"): | |
| parent_exists = True | |
| if not parent_exists: | |
| log_activity(f"[Watchdog SIGKILL Fallback] Reaping orphaned process PID {pid} (PPID {ppid})") | |
| try: | |
| os.kill(pid, signal.SIGTERM) | |
| time.sleep(2) | |
| os.kill(pid, signal.SIGKILL) | |
| except Exception: | |
| pass | |
| except Exception: | |
| pass | |
| except Exception as e: | |
| log_activity(f"[Watchdog Error] {e}") | |
| time.sleep(60) | |
| async def db_heartbeat_loop(): | |
| log_activity("Database Heartbeat task started") | |
| while True: | |
| try: | |
| if db_pool: | |
| async with db_pool.acquire() as conn: | |
| await conn.execute("SELECT 1") | |
| log_activity("[Heartbeat] Pinged Aiven PostgreSQL successfully") | |
| except Exception as e: | |
| log_activity(f"[Heartbeat Warning] Failed to ping database: {e}") | |
| await asyncio.sleep(240) # Every 4 minutes | |
| def run_backup_loop(): | |
| log_activity("Local Git Backup Loop started") | |
| while True: | |
| # Wait 5 minutes between runs | |
| time.sleep(300) | |
| if not BACKUP_GIT_REPO: | |
| continue | |
| try: | |
| log_activity("[Backup] Starting local workspace backup...") | |
| # 1. Clean local backup directory | |
| backup_local_dir = "/tmp/git_backup_repo" | |
| if os.path.exists(backup_local_dir): | |
| shutil.rmtree(backup_local_dir) | |
| os.makedirs(backup_local_dir, exist_ok=True) | |
| # 2. Point-in-time snapshot copy | |
| snapshot_dir = "/tmp/workspace_snapshot" | |
| if os.path.exists(snapshot_dir): | |
| shutil.rmtree(snapshot_dir) | |
| shutil.copytree(WORKSPACE_DIR, snapshot_dir, symlinks=True, ignore=shutil.ignore_patterns('.git', 'node_modules', '.next')) | |
| # 3. Zip snapshot | |
| archive_base = "/tmp/workspace_backup_download" | |
| archive_zip = archive_base + ".zip" | |
| if os.path.exists(archive_zip): | |
| os.unlink(archive_zip) | |
| shutil.make_archive(archive_base, 'zip', snapshot_dir) | |
| shutil.rmtree(snapshot_dir) | |
| # 4. Handle Split if needed | |
| target_dest = os.path.join(backup_local_dir, "workspace_backup.zip") | |
| zip_size = os.path.getsize(archive_zip) | |
| max_part_size = 50 * 1024 * 1024 # 50MB | |
| if zip_size > max_part_size: | |
| subprocess.run(f"split -b 50M {archive_zip} {target_dest}.part", shell=True) | |
| else: | |
| shutil.copy(archive_zip, target_dest) | |
| os.unlink(archive_zip) | |
| # 5. Git commit and force push | |
| subprocess.run("git init", shell=True, cwd=backup_local_dir, stdout=subprocess.DEVNULL) | |
| subprocess.run("git config user.name 'Backup Agent'", shell=True, cwd=backup_local_dir, stdout=subprocess.DEVNULL) | |
| subprocess.run("git config user.email 'backup@agent.internal'", shell=True, cwd=backup_local_dir, stdout=subprocess.DEVNULL) | |
| subprocess.run(f"git remote add origin {BACKUP_GIT_REPO}", shell=True, cwd=backup_local_dir, stdout=subprocess.DEVNULL) | |
| subprocess.run("git checkout -b main", shell=True, cwd=backup_local_dir, stdout=subprocess.DEVNULL) | |
| subprocess.run("git add -A", shell=True, cwd=backup_local_dir, stdout=subprocess.DEVNULL) | |
| subprocess.run('git commit -m "Auto-backup: ' + time.strftime("%Y-%m-%d %H:%M:%S") + '"', shell=True, cwd=backup_local_dir, stdout=subprocess.DEVNULL) | |
| result = subprocess.run("git push origin main --force", shell=True, cwd=backup_local_dir, stdout=subprocess.DEVNULL) | |
| if result.returncode == 0: | |
| log_activity("[Backup] Sync completed successfully (git history purged)") | |
| else: | |
| log_activity("[Backup Error] Git push failed") | |
| except Exception as e: | |
| log_activity(f"[Backup Error] {e}") | |
| async def db_cleanup_loop(): | |
| log_activity("Database Retention Cleanup task started") | |
| while True: | |
| try: | |
| if db_pool: | |
| async with db_pool.acquire() as conn: | |
| result = await conn.execute( | |
| "DELETE FROM agent_session_entries WHERE created_at < NOW() - INTERVAL '30 days'" | |
| ) | |
| log_activity(f"[Cleanup] Nightly retention sweep complete. Status: {result}") | |
| except Exception as e: | |
| log_activity(f"[Cleanup Warning] Failed to run retention cleanup: {e}") | |
| await asyncio.sleep(86400) # Every 24 hours | |
| async def startup_event(): | |
| # Start the watchdog thread on startup | |
| threading.Thread(target=run_watchdog, daemon=True).start() | |
| # Start the local backup loop thread | |
| threading.Thread(target=run_backup_loop, daemon=True).start() | |
| # Start the eternity R&D loop thread | |
| threading.Thread(target=run_eternity_loop, daemon=True).start() | |
| # Start the db keep-alive loop on FastAPI event loop | |
| asyncio.create_task(db_heartbeat_loop()) | |
| # Start the db nightly retention cleanup loop | |
| asyncio.create_task(db_cleanup_loop()) | |
| # --------------------------------------------------------------------------- | |
| # Entrypoint | |
| # --------------------------------------------------------------------------- | |
| if __name__ == "__main__": | |
| import uvicorn | |
| uvicorn.run(app, host="0.0.0.0", port=7860) | |