Commit 71aa8270 authored by Admin's avatar Admin

feat: add agent log viewer, SSE streaming, and L1 sandbox

- services/agent_log_writer.py: Writes ExecutionRunLog rows via AIAgent
  callbacks (tool_start, tool_complete, step, status). Errors swallowed
  to never block agent execution.

- services/workspace_manager.py: Creates isolated workspace directories
  per agent run (workspaces/{company_id}/{run_id}/) with src/,
  .agent_output/, README.md structure.

- services/path_security.py: Path validation to prevent directory
  traversal, blocks dangerous paths (C:\Windows, /etc, .env, .git).

- api/routes/run_logs.py: Added company-scoped runs_router with:
  GET /companies/{id}/runs/{runId} — run detail with all logs
  GET /companies/{id}/runs/{runId}/stream — SSE real-time log stream

- api/routes/agents.py: run_agent_in_background now creates workspace,
  instantiates AgentLogWriter, passes callbacks to AIAgent, and cleans
  up workspace in finally block.

- tests/test_agent_flow.py: Updated to mock WorkspaceManager and
  AgentLogWriter, relaxed AIAgent assertion to check core kwargs only.

All 3 tests pass. Dashboard verified 0 JS errors via Playwright.
parent 5bed3e6c
......@@ -51,7 +51,7 @@ from .routes.authz import router as authz_router
from .routes.issue_recovery import router as issue_recovery_router
from .routes.auth_sessions import router as auth_sessions_router
from .routes.feedback import router as feedback_router
from .routes.run_logs import router as run_logs_router
from .routes.run_logs import router as run_logs_router, runs_router as runs_router
from .routes.environment_selection import router as environment_selection_router
from .routes.workspace_command_authz import router as workspace_command_authz_router
from .routes.workspace_runtime_service_authz import router as workspace_runtime_service_authz_router
......@@ -127,6 +127,7 @@ api_sub_router.include_router(issue_recovery_router)
api_sub_router.include_router(auth_sessions_router)
api_sub_router.include_router(feedback_router)
api_sub_router.include_router(run_logs_router)
api_sub_router.include_router(runs_router)
api_sub_router.include_router(environment_selection_router)
api_sub_router.include_router(workspace_command_authz_router)
api_sub_router.include_router(workspace_runtime_service_authz_router)
......
......@@ -6,6 +6,7 @@ from typing import Optional, List, Dict, Any
from fastapi import APIRouter, Depends, HTTPException, Query, status, BackgroundTasks
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select, func
import os
import uuid
from datetime import datetime
......@@ -16,6 +17,8 @@ from schemas.common import PaginationParams, PaginatedResponse, CompanyScope
from schemas.auth import AuthUserResponse
from models import Agent, CompanyMembership, HeartbeatRun
from services.activity_logger import log_activity
from services.workspace_manager import WorkspaceManager
from services.agent_log_writer import AgentLogWriter
router = APIRouter(prefix="/agents", tags=["agents"])
......@@ -186,10 +189,21 @@ async def run_agent_in_background(
system_message: Optional[str],
conversation_history: Optional[List[Dict[str, Any]]]
):
workspace = None
try:
# Create isolated workspace for this run
workspace = await WorkspaceManager.create(company_id, run_id)
os.environ["AGENT_WORKSPACE"] = workspace.root_path
log_writer = AgentLogWriter(run_id=run_id, company_id=company_id)
ai_agent = AIAgent(
model_name=model_name,
company_id=company_id,
tool_start_callback=log_writer.on_tool_start,
tool_complete_callback=log_writer.on_tool_complete,
step_callback=log_writer.on_step,
status_callback=log_writer.on_status,
)
loop_result = await ai_agent.run_conversation(
......@@ -227,6 +241,9 @@ async def run_agent_in_background(
db_run.status = "failed"
db_run.completed_at = datetime.utcnow()
db_run.error_message = str(e)
finally:
if workspace:
await WorkspaceManager.cleanup(company_id, run_id)
@router.post("/{agent_id}/wakeup")
async def wakeup_agent(
......
......@@ -4,12 +4,15 @@ Execution run logs routes.
from typing import Optional
from fastapi import APIRouter, Depends, HTTPException, Query, status, Body
from fastapi.responses import StreamingResponse
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select, func, desc
import asyncio
import json
import uuid
from schemas.auth import AuthUserResponse
from database import get_db_session
from database import get_db_session, get_db_context
from middleware.auth import require_company_scope, get_current_user
from models import ExecutionRunLog, HeartbeatRun
from schemas.execution_run_log import ExecutionRunLogCreate, ExecutionRunLogUpdate, ExecutionRunLogResponse
......@@ -231,3 +234,131 @@ async def delete_run_log(
await db.commit()
return None
# ------------------------------------------------------------------
# Company-scoped run endpoints (mounted separately in main_router)
# ------------------------------------------------------------------
runs_router = APIRouter(prefix="/companies/{company_id}/runs", tags=["execution"])
@runs_router.get("/{run_id}")
async def get_run_detail(
company_id: str,
run_id: str,
scope: CompanyScope = Depends(require_company_scope),
db: AsyncSession = Depends(get_db_session),
):
"""Return a HeartbeatRun with its execution logs."""
result = await db.execute(
select(HeartbeatRun).where(
HeartbeatRun.id == run_id,
HeartbeatRun.company_id == company_id,
)
)
run = result.scalar_one_or_none()
if not run:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Run not found")
logs_result = await db.execute(
select(ExecutionRunLog)
.where(ExecutionRunLog.run_id == run_id)
.order_by(ExecutionRunLog.created_at.asc())
)
logs = logs_result.scalars().all()
return {
"id": run.id,
"company_id": run.company_id,
"agent_id": run.agent_id,
"issue_id": run.issue_id,
"status": run.status,
"started_at": run.started_at,
"completed_at": run.completed_at,
"error_message": run.error_message,
"created_at": run.created_at,
"updated_at": run.updated_at,
"logs": [
ExecutionRunLogResponse.model_validate(log) for log in logs
],
}
@runs_router.get("/{run_id}/stream")
async def stream_run_logs(
company_id: str,
run_id: str,
scope: CompanyScope = Depends(require_company_scope),
db: AsyncSession = Depends(get_db_session),
):
"""SSE stream of execution logs for a running agent.
Polls the ExecutionRunLog table every second and yields new entries.
Automatically closes when the HeartbeatRun status leaves "running".
"""
# Verify run exists and belongs to company
result = await db.execute(
select(HeartbeatRun).where(
HeartbeatRun.id == run_id,
HeartbeatRun.company_id == company_id,
)
)
run = result.scalar_one_or_none()
if not run:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Run not found")
async def event_generator():
last_seen_id: Optional[str] = None
while True:
async with get_db_context() as poll_db:
# Fetch new log entries since last seen
stmt = (
select(ExecutionRunLog)
.where(ExecutionRunLog.run_id == run_id)
.order_by(ExecutionRunLog.created_at.asc(), ExecutionRunLog.id.asc())
)
if last_seen_id:
# Get the created_at of the last seen log for cursor-based paging
prev_result = await poll_db.execute(
select(ExecutionRunLog.created_at).where(ExecutionRunLog.id == last_seen_id)
)
prev_row = prev_result.scalar_one_or_none()
if prev_row:
stmt = stmt.where(
(ExecutionRunLog.created_at > prev_row)
| (
(ExecutionRunLog.created_at == prev_row)
& (ExecutionRunLog.id > last_seen_id)
)
)
logs_result = await poll_db.execute(stmt.limit(100))
logs = logs_result.scalars().all()
for log in logs:
data = {
"id": log.id,
"log_type": log.log_type,
"message": log.message,
"metadata": log.log_metadata,
"created_at": log.created_at.isoformat(),
}
yield f"data: {json.dumps(data, default=str)}\n\n"
last_seen_id = log.id
# Check if the run is still active
run_result = await poll_db.execute(
select(HeartbeatRun.status).where(HeartbeatRun.id == run_id)
)
current_status = run_result.scalar_one_or_none()
if current_status and current_status != "running":
# Emit a final status event
yield f"data: {json.dumps({'event': 'done', 'status': current_status})}\n\n"
return
await asyncio.sleep(1)
return StreamingResponse(event_generator(), media_type="text/event-stream")
"""
Agent log writer — writes ExecutionRunLog rows from agent callbacks.
"""
import uuid
import json
import logging
from datetime import datetime
from typing import Any, Dict, Optional
from database import get_db_context
from models import ExecutionRunLog
logger = logging.getLogger(__name__)
class AgentLogWriter:
"""Writes agent lifecycle events to the execution_run_logs table.
Provides callback methods compatible with AIAgent's callback interface.
Each callback writes one ExecutionRunLog row via an independent DB session.
"""
def __init__(self, run_id: str, company_id: str):
self.run_id = run_id
self.company_id = company_id
async def _write_log(self, log_type: str, message: str, metadata: Optional[Dict[str, Any]] = None):
"""Persist a single log row. Errors are swallowed so agent execution is never blocked."""
try:
async with get_db_context() as db:
log = ExecutionRunLog(
id=uuid.uuid4().hex,
run_id=self.run_id,
log_type=log_type,
message=message,
log_metadata=metadata,
created_at=datetime.utcnow(),
updated_at=datetime.utcnow(),
)
db.add(log)
except Exception:
logger.warning("Failed to write agent log for run %s", self.run_id, exc_info=True)
# -- Callbacks for AIAgent --
async def on_tool_start(self, tool_name: str, args: Any = None):
args_str = json.dumps(args, default=str, ensure_ascii=False) if args else ""
await self._write_log(
log_type="info",
message=f"Tool started: {tool_name}",
metadata={"event": "tool_start", "tool_name": tool_name, "args": args_str},
)
async def on_tool_complete(self, tool_name: str, result: Any = None):
result_str = str(result)[:2000] if result else ""
await self._write_log(
log_type="info",
message=f"Tool completed: {tool_name}",
metadata={"event": "tool_complete", "tool_name": tool_name, "result_preview": result_str},
)
async def on_thinking(self, text: str):
await self._write_log(
log_type="info",
message=text[:4000] if text else "",
metadata={"event": "thinking"},
)
async def on_step(self, step_info: Any = None):
message = str(step_info)[:4000] if step_info else "Step"
await self._write_log(
log_type="info",
message=message,
metadata={"event": "step"},
)
async def on_status(self, status: str, detail: str = ""):
await self._write_log(
log_type="info",
message=detail or status,
metadata={"event": "status", "status": status},
)
"""
Path security utilities for workspace directory jail.
Validates and sanitizes file paths to prevent directory traversal
and access to sensitive system locations.
"""
from pathlib import Path
from typing import List
# Dangerous path patterns that should never be accessed by agents
BLOCKED_PATTERNS: List[str] = [
# Unix sensitive paths
"/etc/",
"/root/",
"/var/",
"/proc/",
"/sys/",
"/dev/",
"/boot/",
"/usr/",
"/sbin/",
"/tmp/",
"/.ssh/",
"/.gnupg/",
# Windows sensitive paths
"C:\\Windows\\",
"C:\\Program Files\\",
"C:\\Program Files (x86)\\",
"C:\\Users\\",
"C:\\ProgramData\\",
"C:\\System Volume Information\\",
# Common env/secret files
".env",
".git/",
"__pycache__/",
]
def is_path_safe(requested_path: str, workspace_root: str) -> bool:
"""
Check if a requested path is safely within the workspace root.
Resolves symlinks and blocks:
- Directory traversal via '..'
- Absolute paths outside workspace
- Paths matching blocked patterns
"""
try:
root = Path(workspace_root).resolve()
# Handle relative and absolute paths
candidate = Path(requested_path)
if not candidate.is_absolute():
candidate = root / candidate
resolved = candidate.resolve()
# Must be within workspace root
if not str(resolved).startswith(str(root)):
return False
# Check against blocked patterns
resolved_str = str(resolved)
for pattern in BLOCKED_PATTERNS:
if pattern.lower() in resolved_str.lower():
return False
return True
except (OSError, ValueError):
return False
def sanitize_path(requested_path: str, workspace_root: str) -> str:
"""
Return a safe, resolved path within the workspace.
Raises ValueError if the path escapes the workspace or matches
a blocked pattern.
"""
root = Path(workspace_root).resolve()
candidate = Path(requested_path)
if not candidate.is_absolute():
candidate = root / candidate
resolved = candidate.resolve()
# Must be within workspace root
if not str(resolved).startswith(str(root)):
raise ValueError(
f"Path '{requested_path}' resolves outside workspace root"
)
# Check against blocked patterns
resolved_str = str(resolved)
for pattern in BLOCKED_PATTERNS:
if pattern.lower() in resolved_str.lower():
raise ValueError(
f"Path '{requested_path}' matches blocked pattern '{pattern}'"
)
return str(resolved)
"""
Workspace manager for agent directory isolation (L1 Sandbox).
Creates isolated workspace directories per agent run and provides
path validation to prevent directory traversal.
"""
import shutil
import logging
from pathlib import Path
from dataclasses import dataclass
from services.path_security import is_path_safe, sanitize_path
logger = logging.getLogger(__name__)
# Base directory for all workspaces, relative to project root
WORKSPACES_BASE = Path(__file__).resolve().parent.parent.parent / "workspaces"
@dataclass
class Workspace:
"""Represents an isolated agent workspace."""
root_path: str
company_id: str
run_id: str
def validate_path(self, path: str) -> bool:
"""Check if a path is safely within this workspace."""
return is_path_safe(path, self.root_path)
def safe_path(self, path: str) -> str:
"""Return a resolved safe path within this workspace. Raises ValueError if unsafe."""
return sanitize_path(path, self.root_path)
class WorkspaceManager:
"""Manages isolated workspace directories for agent runs."""
@staticmethod
async def create(company_id: str, run_id: str) -> Workspace:
"""
Create an isolated workspace directory for an agent run.
Structure:
workspaces/{company_id}/{run_id}/
src/
.agent_output/
README.md
"""
workspace_dir = WORKSPACES_BASE / company_id / run_id
workspace_dir.mkdir(parents=True, exist_ok=True)
# Create initial structure
(workspace_dir / "src").mkdir(exist_ok=True)
(workspace_dir / ".agent_output").mkdir(exist_ok=True)
readme_path = workspace_dir / "README.md"
if not readme_path.exists():
readme_path.write_text(
f"# Agent Workspace\n\n"
f"- **Company**: {company_id}\n"
f"- **Run**: {run_id}\n",
encoding="utf-8",
)
workspace = Workspace(
root_path=str(workspace_dir.resolve()),
company_id=company_id,
run_id=run_id,
)
logger.info("Created workspace: %s", workspace.root_path)
return workspace
@staticmethod
async def cleanup(company_id: str, run_id: str) -> None:
"""Remove a workspace directory after a run completes."""
workspace_dir = WORKSPACES_BASE / company_id / run_id
if workspace_dir.exists():
shutil.rmtree(workspace_dir, ignore_errors=True)
logger.info("Cleaned up workspace: %s", workspace_dir)
@staticmethod
async def get(company_id: str, run_id: str) -> Workspace | None:
"""Get an existing workspace, or None if it doesn't exist."""
workspace_dir = WORKSPACES_BASE / company_id / run_id
if not workspace_dir.exists():
return None
return Workspace(
root_path=str(workspace_dir.resolve()),
company_id=company_id,
run_id=run_id,
)
......@@ -75,7 +75,15 @@ async def test_agent_wake_route_integration():
try:
async with httpx.AsyncClient(transport=transport, base_url="http://testserver") as client:
# Mock AIAgent loop
with patch("api.routes.agents.AIAgent") as mock_ai_agent_cls, patch("api.routes.agents.get_db_context", mock_get_db_context):
with patch("api.routes.agents.AIAgent") as mock_ai_agent_cls, \
patch("api.routes.agents.get_db_context", mock_get_db_context), \
patch("api.routes.agents.WorkspaceManager") as mock_ws_mgr, \
patch("api.routes.agents.AgentLogWriter"):
# Mock workspace creation
mock_workspace = MagicMock()
mock_workspace.root_path = "C:\\temp\\workspace"
mock_ws_mgr.create = AsyncMock(return_value=mock_workspace)
mock_ws_mgr.cleanup = AsyncMock()
# Mock CompanyMembership check query result
mock_membership = MagicMock()
......@@ -137,11 +145,11 @@ async def test_agent_wake_route_integration():
assert json_data["agentId"] == "agent_123"
assert json_data["companyId"] == "company_123"
# Verify AIAgent was instantiated and run correctly
mock_ai_agent_cls.assert_called_once_with(
model_name="gemini-3.1-flash",
company_id="company_123"
)
# Verify AIAgent was instantiated with correct core params
mock_ai_agent_cls.assert_called_once()
call_kwargs = mock_ai_agent_cls.call_args[1]
assert call_kwargs["model_name"] == "gemini-3.1-flash"
assert call_kwargs["company_id"] == "company_123"
mock_agent_instance.run_conversation.assert_called_once_with(
user_message="Run testing task",
system_message="Override prompt",
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment