[Haik]: ckpt, copied over health and dashboard sub apps. Also made a generic PydanticStore and swapped the agents subapp to use this instead. Now gonna abstract and clean up this agents sub app a bit, then gonna make the dashboards sub app functional via the new codebase not the legacy shit

This commit is contained in:
haikdc
2026-04-05 08:21:13 -07:00
parent ef688f0a2e
commit eca8cccf1a
6 changed files with 477 additions and 168 deletions
+67 -54
View File
@@ -7,40 +7,41 @@ The ws module is used ONLY in this file — the Agent class and its internals
communicate via the on_event callback, never importing ws directly.
"""
import asyncio
import os
from contextlib import asynccontextmanager
from datetime import datetime
from typeguard import typechecked
from typing import Optional, List, Dict
from uuid import uuid4
import asyncio
from typeguard import typechecked
from fastapi import HTTPException
from pydantic import BaseModel
from typing import Optional, List, Dict
from backend.config.Apps import SubApp
from backend.config.paths import DB_ROOT
from backend.core.Agent.Agent import Agent
from backend.core.db.PydanticStore import PydanticStore
from backend.core.shared_structs.agent.Message.Message import UserMessage
from backend.core.events.events import (
AnyEvent, AgentStatusEvent, AgentClosedEvent, BranchSwitchedEvent,
ApprovalRequestEvent, EventCallback,
)
from backend.apps.agents.session_store import (
load_all,
save,
delete,
build_search_text,
get_history as session_store_get_history,
reconcile_on_startup,
load,
)
from backend.apps.agents import ws
from backend.apps.agents.compose_system_prompt import compose_system_prompt
from backend.core.tools.make_builtin_toolkit.make_builtin_toolkit import make_builtin_toolkit
from backend.apps.agents.create_sdk_hooks import create_sdk_hooks
from claude_agent_sdk import ClaudeAgentOptions
from claude_agent_sdk.types import HookMatcher
from claude_agent_sdk.types import HookMatcher, McpServerConfig
from backend.core.tools.shared_structs.Toolkit import Toolkit
from claude_agent_sdk.types import McpServerConfig
AGENT_STORE: PydanticStore[Agent] = PydanticStore[Agent](
model_cls=Agent,
data_dir=os.path.join(DB_ROOT, "sessions"),
id_field="session_id",
dump_mode="json",
not_found_detail="Session not found in history",
)
SESSIONS: dict[str, Agent] = {}
@@ -109,28 +110,19 @@ def get_agent(session_id: str) -> Agent:
@asynccontextmanager
async def agents_lifespan():
await reconcile_on_startup()
for sid, data in load_all():
if data.get("closed_at") is not None:
continue
for stored in AGENT_STORE.load_all():
try:
data.pop("task", None)
data.pop("lock", None)
agent: Agent = Agent(**data)
agent.status = "stopped"
agent.on_event = p_make_session_emitter(agent.session_id)
toolkit: Toolkit = make_builtin_toolkit(agent, SESSIONS, p_send_browser_command)
agent.toolkit = toolkit
SESSIONS[agent.session_id] = agent
delete(sid)
stored.status = "stopped"
stored.on_event = p_make_session_emitter(stored.session_id)
toolkit: Toolkit = make_builtin_toolkit(stored, SESSIONS, p_send_browser_command)
stored.toolkit = toolkit
SESSIONS[stored.session_id] = stored
except Exception as e:
print(f"[agents lifespan] Skipping corrupt session {sid}: {e}")
print(f"[agents lifespan] Skipping corrupt session {stored.session_id}: {e}")
yield
for agent in list[Agent](SESSIONS.values()):
await agent.stop_agent()
data: dict = agent.model_dump(mode="json")
data["search_text"] = build_search_text(data)
save(agent.session_id, data)
AGENT_STORE.save(agent)
SESSIONS.clear()
@@ -222,7 +214,7 @@ async def delete_session(session_id: str) -> dict:
agent: Optional[Agent] = SESSIONS.pop(session_id, None)
if agent is not None:
await agent.stop_agent()
delete(session_id)
AGENT_STORE.delete(session_id)
return {"ok": True}
@@ -329,11 +321,9 @@ async def close_session(session_id: str) -> dict:
if not agent:
raise HTTPException(status_code=404, detail="Session not found")
await agent.stop_agent()
closed_at: str = datetime.now().isoformat()
data: dict = agent.model_dump(mode="json")
data["search_text"] = build_search_text(data)
data["closed_at"] = closed_at
save(session_id, data)
AGENT_STORE.save(agent)
msgs = agent.messages.messages
closed_at: str = msgs[-1].timestamp.isoformat() if msgs else datetime.now().isoformat()
await agent.emit(AgentClosedEvent(
session_id=session_id, status=agent.status,
closed_at=closed_at,
@@ -345,20 +335,15 @@ async def close_session(session_id: str) -> dict:
async def resume_session(session_id: str) -> dict:
if session_id in SESSIONS:
return {"session": SESSIONS[session_id].model_dump(mode="json")}
data: Optional[dict] = load(session_id)
if not data:
agent: Optional[Agent] = AGENT_STORE.load_or_none(session_id)
if not agent:
raise HTTPException(status_code=404, detail="Session not found in history")
data.pop("task", None)
data.pop("lock", None)
data.pop("search_text", None)
data.pop("closed_at", None)
agent: Agent = Agent(**data)
agent.status = "stopped"
agent.on_event = p_make_session_emitter(agent.session_id)
toolkit: Toolkit = make_builtin_toolkit(agent, SESSIONS, p_send_browser_command)
agent.toolkit = toolkit
SESSIONS[agent.session_id] = agent
delete(session_id)
AGENT_STORE.delete(session_id)
await agent.emit(AgentStatusEvent(
session_id=session_id, status=agent.status,
session=agent.snapshot(),
@@ -370,12 +355,9 @@ async def resume_session(session_id: str) -> dict:
async def duplicate_session(session_id: str, body: dict = {}) -> dict:
source: Optional[Agent] = SESSIONS.get(session_id)
if source is None:
data: Optional[dict] = load(session_id)
if not data:
source = AGENT_STORE.load_or_none(session_id)
if not source:
raise HTTPException(status_code=404, detail="Session not found")
data.pop("task", None)
data.pop("lock", None)
source = Agent(**data)
clone: Agent = source.model_copy(deep=True)
clone.session_id = uuid4().hex
@@ -395,9 +377,40 @@ async def duplicate_session(session_id: str, body: dict = {}) -> dict:
return {"session": clone.snapshot().model_dump(mode="json")}
@typechecked
def p_build_search_text(agent: Agent, max_len: int = 5000) -> str:
parts: List[str] = []
for msg in agent.messages.messages:
if msg.role in ("user", "assistant") and isinstance(msg.content, str):
parts.append(msg.content)
return " ".join(parts)[:max_len]
@agents.router.get("/history")
async def get_history(q: str = "", limit: int = 20, offset: int = 0, dashboard_id: str = "") -> dict:
return session_store_get_history(
q=q, limit=limit, offset=offset,
dashboard_id=dashboard_id or None,
)
all_agents: List[Agent] = AGENT_STORE.load_all()
all_agents.sort(
key=lambda a: a.messages.messages[-1].timestamp if a.messages.messages else datetime.min,
reverse=True,
)
q_lower: str = q.strip().lower()
history: List[dict] = []
for agent in all_agents:
if q_lower:
search_text: str = p_build_search_text(agent).lower()
if q_lower not in search_text:
continue
msgs = agent.messages.messages
closed_at: str = msgs[-1].timestamp.isoformat() if msgs else ""
history.append({
"id": agent.session_id,
"status": agent.status,
"model": agent.model,
"mode": agent.mode,
"closed_at": closed_at,
})
total: int = len(history)
page: List[dict] = history[offset : offset + limit]
return {"sessions": page, "total": total, "has_more": offset + limit < total}
-114
View File
@@ -1,114 +0,0 @@
"""On-disk JSON persistence for agent sessions.
Each session is stored as {session_id}.json in SESSIONS_DIR.
The agents subapp calls these functions during close/resume/startup/shutdown.
"""
import json
import os
from typing import List, Tuple, Optional
from backend.config.paths import DB_ROOT
from typeguard import typechecked
SESSIONS_DIR = os.path.join(DB_ROOT, "sessions")
@typechecked
def p_path(session_id: str) -> str:
return os.path.join(SESSIONS_DIR, f"{session_id}.json")
@typechecked
def save(session_id: str, data: dict) -> None:
os.makedirs(SESSIONS_DIR, exist_ok=True)
with open(p_path(session_id), "w") as f:
json.dump(data, f, indent=2)
@typechecked
def load(session_id: str) -> Optional[dict]:
path: str = p_path(session_id)
if not os.path.exists(path):
return None
with open(path) as f:
return json.load(f)
@typechecked
def delete(session_id: str) -> None:
path: str = p_path(session_id)
if os.path.exists(path):
os.remove(path)
# TODO: better type spec for the dict type
@typechecked
def load_all() -> List[Tuple[str, dict]]:
results: List[Tuple[str, dict]] = []
if not os.path.exists(SESSIONS_DIR):
return results
for fname in os.listdir(SESSIONS_DIR):
if fname.endswith(".json"):
try:
with open(os.path.join(SESSIONS_DIR, fname)) as f:
results.append((fname[:-5], json.load(f)))
except (json.JSONDecodeError, OSError) as e:
print(f"[session_store.load_all] Skipping corrupt session file {fname}: {e}")
return results
# TODO: better type spec for the whole damn thing
@typechecked
def build_search_text(agent_data: dict, max_len: int = 5000) -> str:
parts = [agent_data.get("name", "")]
for msg in agent_data.get("messages", {}).get("messages", []):
role = msg.get("role")
content = msg.get("content")
if role in ("user", "assistant") and isinstance(content, str):
parts.append(content)
return " ".join(parts)[:max_len]
# TODO: all the dict get guesswork shld be swapped to smthn less ambiguous/guessy
@typechecked
def get_history(
q: str = "",
limit: int = 20,
offset: int = 0,
dashboard_id: Optional[str] = None,
) -> dict:
all_data: List[Tuple[str, dict]] = load_all()
all_data.sort(key=lambda pair: pair[1].get("closed_at") or "", reverse=True)
q_lower: str = q.strip().lower()
history: List[dict] = []
for sid, data in all_data:
if dashboard_id and data.get("dashboard_id") != dashboard_id:
continue
if q_lower:
name: str = (data.get("name") or "").lower()
search_text: str = (data.get("search_text") or "").lower()
if q_lower not in name and q_lower not in search_text:
continue
history.append({
"id": data.get("session_id", sid),
"name": data.get("name", "Untitled"),
"status": data.get("status", "stopped"),
"model": data.get("model", "sonnet"),
"mode": data.get("mode", "agent"),
"created_at": data.get("created_at"),
"closed_at": data.get("closed_at"),
"cost_usd": data.get("cost_usd", 0),
"dashboard_id": data.get("dashboard_id"),
})
total: int = len(history)
page: List[dict] = history[offset : offset + limit]
return {"sessions": page, "total": total, "has_more": offset + limit < total}
@typechecked
async def reconcile_on_startup() -> None:
for sid, data in load_all():
if data.get("status") in ("running", "waiting_approval"):
data["status"] = "stopped"
save(sid, data)
print(f"[session_store.reconcile_on_startup] Marked stale session {sid} as stopped")
+55
View File
@@ -0,0 +1,55 @@
from pydantic import BaseModel, Field
from typing import Optional, List, Dict
from datetime import datetime
from uuid import uuid4
class CardPosition(BaseModel):
session_id: str
x: float = 0
y: float = 0
width: float = 420
height: float = 280
class ViewCardPosition(BaseModel):
output_id: str
x: float = 0
y: float = 0
width: float = 480
height: float = 360
class BrowserTab(BaseModel):
id: str
url: str = ""
title: str = ""
favicon: Optional[str] = None
class BrowserCardPosition(BaseModel):
browser_id: str
url: str = ""
tabs: List[BrowserTab] = Field(default_factory=list)
activeTabId: str = ""
x: float = 0
y: float = 0
width: float = 1280
height: float = 800
class DashboardLayout(BaseModel):
cards: Dict[str, CardPosition] = Field(default_factory=dict)
view_cards: Dict[str, ViewCardPosition] = Field(default_factory=dict)
browser_cards: Dict[str, BrowserCardPosition] = Field(default_factory=dict)
expanded_session_ids: list[str] = Field(default_factory=list)
class Dashboard(BaseModel):
id: str = Field(default_factory=lambda: uuid4().hex)
name: str = "Untitled Dashboard"
auto_named: bool = False
created_at: datetime = Field(default_factory=datetime.now)
updated_at: datetime = Field(default_factory=datetime.now)
layout: DashboardLayout = Field(default_factory=DashboardLayout)
thumbnail: Optional[str] = None
+242
View File
@@ -0,0 +1,242 @@
import json
import os
import logging
from contextlib import asynccontextmanager
from datetime import datetime
from uuid import uuid4
from backend.config.Apps import SubApp
# from backend.apps.common.json_store import JsonStore
from backend.core.db.PydanticStore import PydanticStore
from backend.apps.dashboards.models import (
Dashboard,
DashboardCreate,
DashboardUpdate,
DashboardLayout,
)
from backend.apps.common.llm_helpers import _resolve_model as _rm
from backend.apps.settings.settings import load_settings
from backend.apps.settings.credentials import get_anthropic_client
logger = logging.getLogger(__name__)
from backend.config.paths import DASHBOARDS_DIR as DATA_DIR, SESSIONS_DIR, DASHBOARD_LAYOUT_DIR as OLD_LAYOUT_DIR
OLD_LAYOUT_FILE = os.path.join(OLD_LAYOUT_DIR, "layout.json")
_store = JsonStore(
Dashboard, DATA_DIR, dump_mode="json", not_found_detail="Dashboard not found",
)
_load_all = _store.load_all
_save = _store.save
_load = _store.load
_delete = _store.delete
def _migrate_if_needed():
"""One-time migration: if no dashboards exist, create 'Dashboard 1' from old layout."""
existing = _load_all()
if existing:
return
logger.info("No dashboards found — running one-time migration")
layout = DashboardLayout()
if os.path.exists(OLD_LAYOUT_FILE):
try:
with open(OLD_LAYOUT_FILE) as f:
data = json.load(f)
if "cards" in data:
layout = DashboardLayout(**data)
logger.info("Migrated layout from old layout.json")
except Exception:
logger.exception("Failed to read old layout.json, using empty layout")
dashboard = Dashboard(name="Dashboard 1", layout=layout)
_save(dashboard)
logger.info(f"Created default dashboard: {dashboard.id}")
if os.path.exists(SESSIONS_DIR):
count = 0
for fname in os.listdir(SESSIONS_DIR):
if not fname.endswith(".json"):
continue
fpath = os.path.join(SESSIONS_DIR, fname)
with open(fpath) as f:
session_data = json.load(f)
session_data["dashboard_id"] = dashboard.id
with open(fpath, "w") as f:
json.dump(session_data, f, indent=2)
count += 1
if count:
logger.info(f"Tagged {count} existing chat sessions with dashboard_id={dashboard.id}")
@asynccontextmanager
async def dashboards_lifespan():
os.makedirs(DATA_DIR, exist_ok=True)
_migrate_if_needed()
yield
dashboards = SubApp("dashboards", dashboards_lifespan)
@dashboards.router.get("/list")
async def list_dashboards():
all_dashboards = _load_all()
all_dashboards.sort(key=lambda d: d.updated_at or d.created_at, reverse=True)
items = []
for d in all_dashboards:
dumped = d.model_dump(mode="json")
items.append({
"id": dumped["id"],
"name": dumped.get("name", "Untitled"),
"auto_named": dumped.get("auto_named", False),
"created_at": dumped.get("created_at"),
"updated_at": dumped.get("updated_at"),
"thumbnail": dumped.get("thumbnail"),
})
return {"dashboards": items}
@dashboards.router.post("/create")
async def create_dashboard(body: DashboardCreate):
dashboard = Dashboard(name=body.name)
_save(dashboard)
_analytics("dashboard.created", {}, dashboard_id=dashboard.id)
return dashboard.model_dump(mode="json")
@dashboards.router.post("/{dashboard_id}/generate-name")
async def generate_name(dashboard_id: str):
from backend.apps.agents.manager.agent_manager import agent_manager
dashboard = _load(dashboard_id)
if not dashboard.auto_named and dashboard.name != "Untitled Dashboard":
return {"name": dashboard.name, "auto_named": dashboard.auto_named}
prompts = []
for session in agent_manager.sessions.values():
if getattr(session, "dashboard_id", None) != dashboard_id:
continue
for msg in session.messages:
if msg.role == "user" and isinstance(msg.content, str) and msg.content.strip():
prompts.append(msg.content.strip()[:200])
break
if not prompts:
return {"name": dashboard.name, "auto_named": dashboard.auto_named}
fallback = prompts[0][:40]
try:
global_settings = load_settings()
client = get_anthropic_client(global_settings)
if len(prompts) == 1:
system = "Generate a concise 2-5 word workspace name for a project based on this task. Return only the name, nothing else."
user_content = prompts[0]
else:
system = "Generate a concise 2-5 word workspace name that captures the overall theme of these tasks. Return only the name, nothing else."
user_content = "\n".join(f"- {p}" for p in prompts)
resp = await client.messages.create(
model=_rm("claude-haiku-4-5-20251001", global_settings),
max_tokens=30,
system=system,
messages=[{"role": "user", "content": user_content}],
)
generated = resp.content[0].text.strip().strip('"\'')
if generated:
fallback = generated
except Exception as e:
logger.warning(f"Dashboard name generation failed, using fallback: {e}")
dashboard.name = fallback
dashboard.auto_named = True
dashboard.updated_at = datetime.now()
_save(dashboard)
return {"name": dashboard.name, "auto_named": True}
@dashboards.router.get("/{dashboard_id}")
async def get_dashboard(dashboard_id: str):
dashboard = _load(dashboard_id)
return dashboard.model_dump(mode="json")
@dashboards.router.put("/{dashboard_id}")
async def update_dashboard(dashboard_id: str, body: DashboardUpdate):
dashboard = _load(dashboard_id)
if body.name is not None:
dashboard.name = body.name
dashboard.auto_named = False
if body.layout is not None:
dashboard.layout = body.layout
if body.thumbnail is not None:
dashboard.thumbnail = body.thumbnail
dashboard.updated_at = datetime.now()
_save(dashboard)
return dashboard.model_dump(mode="json")
@dashboards.router.delete("/{dashboard_id}")
async def delete_dashboard(dashboard_id: str):
from backend.apps.agents.manager.agent_manager import agent_manager
_load(dashboard_id)
if os.path.exists(SESSIONS_DIR):
for fname in os.listdir(SESSIONS_DIR):
if not fname.endswith(".json"):
continue
fpath = os.path.join(SESSIONS_DIR, fname)
try:
with open(fpath) as f:
data = json.load(f)
if data.get("dashboard_id") == dashboard_id:
os.remove(fpath)
except Exception:
logger.warning(f"Failed to read/delete session file {fname}")
to_remove = [
sid for sid, sess in agent_manager.sessions.items()
if getattr(sess, "dashboard_id", None) == dashboard_id
]
for sid in to_remove:
try:
await agent_manager.delete_session(sid)
except Exception:
logger.warning(f"Failed to delete active session {sid} during dashboard deletion")
_delete(dashboard_id)
return {"ok": True}
@dashboards.router.post("/{dashboard_id}/duplicate")
async def duplicate_dashboard(dashboard_id: str):
source = _load(dashboard_id)
source_data = source.model_dump(mode="json")
new_id = uuid4().hex
now = datetime.now().isoformat()
new_dashboard = {
**source_data,
"id": new_id,
"name": f"{source_data.get('name', 'Untitled')} (copy)",
"created_at": now,
"updated_at": now,
"layout": {
"cards": {},
"view_cards": source_data.get("layout", {}).get("view_cards", {}),
"browser_cards": source_data.get("layout", {}).get("browser_cards", {}),
},
}
with open(os.path.join(DATA_DIR, f"{new_id}.json"), "w") as f:
json.dump(new_dashboard, f, indent=2)
return new_dashboard
+33
View File
@@ -0,0 +1,33 @@
from backend.config.Apps import SubApp
from contextlib import asynccontextmanager
from fastapi.responses import PlainTextResponse
from typeguard import typechecked
# import debug
from fastapi import status, HTTPException
@asynccontextmanager
async def health_lifespan():
# debug("START")
yield
# debug("END")
health = SubApp("health", health_lifespan)
######################################
# Health Check Endpoints #
######################################
@health.router.get("/check")
@typechecked
async def check() -> PlainTextResponse:
# debug("Health check successful")
# Use PlainTextResponse instead of JSONResponse for AWS ALB compatibility
# ALB health checks can be sensitive to JSON responses and Content-Length headers
return PlainTextResponse(
content="OK",
status_code=status.HTTP_200_OK,
headers={
"Content-Type": "text/plain",
"Content-Length": "2"
}
)
+80
View File
@@ -0,0 +1,80 @@
"""Generic JSON-file CRUD store for Pydantic models in a flat directory.
Every sub-app in the backend persists entities as one JSON file per record
inside a data directory. This module eliminates that copy-paste.
"""
import json
import os
from typing import Generic, List, Optional, TypeVar
from fastapi import HTTPException
from pydantic import BaseModel
from typeguard import typechecked
T = TypeVar("T", bound=BaseModel)
class PydanticStore(BaseModel, Generic[T]):
model_cls: type[T]
data_dir: str
id_field: str = "id"
dump_mode: str | None = None
not_found_detail: str = "Not found"
# -- private methods ------------------------------------------------------
@typechecked
def p_path(self, item_id: str) -> str:
return os.path.join(self.data_dir, f"{item_id}.json")
@typechecked
def p_dump(self, item: T) -> dict:
if self.dump_mode:
return item.model_dump(mode=self.dump_mode)
return item.model_dump()
# -- public methods ------------------------------------------------------
@typechecked
def load_all(self) -> list[T]:
result: List[T] = []
if not os.path.exists(self.data_dir):
return result
for fname in os.listdir(self.data_dir):
if fname.endswith(".json"):
with open(os.path.join(self.data_dir, fname)) as f:
result.append(self.model_cls(**json.load(f)))
return result
@typechecked
def save(self, item: T) -> None:
os.makedirs(self.data_dir, exist_ok=True)
item_id = getattr(item, self.id_field)
with open(self.p_path(item_id), "w") as f:
json.dump(self.p_dump(item), f, indent=2)
@typechecked
def load(self, item_id: str) -> T:
path = self.p_path(item_id)
if not os.path.exists(path):
raise HTTPException(status_code=404, detail=self.not_found_detail)
with open(path) as f:
return self.model_cls(**json.load(f))
@typechecked
def load_or_none(self, item_id: str) -> Optional[T]:
path = self.p_path(item_id)
if not os.path.exists(path):
return None
with open(path) as f:
return self.model_cls(**json.load(f))
@typechecked
def delete(self, item_id: str) -> None:
path = self.p_path(item_id)
if os.path.exists(path):
os.remove(path)
@typechecked
def exists(self, item_id: str) -> bool:
return os.path.exists(self.p_path(item_id))