Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 19 additions & 0 deletions docs/persistence.md
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,25 @@ Backup and restore are deliberately excluded from repositories. The specialized
SQLite's backup API and are shared by autoupdate and the explicit migration
command.

## Dashboard read model

`github_agent_bridge.dashboard_data.DashboardQueries` is the read-only boundary
used by HTTP handlers, CLI status/list commands, and monitoring. It owns
dashboard SQL and row-to-JSON mapping separately from write repositories so
denormalized UI queries can evolve without leaking `sqlite3.Row` or SQL into
application orchestration.

Job-list callers pass a typed `JobListFilters` value. The query maps each field
to a fixed column and binds every value; callers cannot supply column names.
The list remains one SQL statement and keeps the indexed dashboard ordering.

The read model assumes the current migrated schema. Health/status first
validates migration history and reports `schema_ok=false` for an old or partial
database; normal queries do not carry indefinite `table_exists()` or
`column_exists()` branches. Apply migrations explicitly before starting a newer
dashboard. The module-level query functions remain compatibility shims for
Python callers, while services use `DashboardQueries` directly.

Executor heartbeat writes, acknowledgement claims and streamed session-event
writes treat `SQLITE_BUSY` and `SQLITE_LOCKED` as transient contention after the
connection timeout: they wait and retry instead of terminating the background
Expand Down
77 changes: 36 additions & 41 deletions src/github_agent_bridge/backend.py
Original file line number Diff line number Diff line change
Expand Up @@ -42,16 +42,8 @@
update_rule_scope,
)
from .dashboard_data import (
get_job_detail,
inspect_db_read_only,
job_logs,
job_session,
job_session_events,
job_session_transcript,
list_all_job_actor_logins,
list_job_actors,
list_jobs,
metrics_summary,
DashboardQueries,
JobListFilters,
transcript_entry_from_session_event,
)
from .monitor import monitor
Expand Down Expand Up @@ -517,11 +509,12 @@ async def _sleep_or_shutdown(shutdown_event: asyncio.Event | None, sleep_seconds


async def _session_stream_events(db: str | Path, job_id: int, *, after_id: int | None = None, sleep_seconds: float = 2.0, shutdown_event: asyncio.Event | None = None):
queries = DashboardQueries(db)
last_id = after_id or 0
sent_transcript_keys: set[str] = set()
while shutdown_event is None or not shutdown_event.is_set():
emitted = False
events = job_session_events(db, job_id, after_id=last_id, limit=100)
events = queries.job_session_events(job_id, after_id=last_id, limit=100)
for event in events:
if shutdown_event is not None and shutdown_event.is_set():
return
Expand All @@ -534,7 +527,7 @@ async def _session_stream_events(db: str | Path, job_id: int, *, after_id: int |
if key not in sent_transcript_keys:
sent_transcript_keys.add(key)
yield _sse_event("transcript_entry", {"job_id": job_id, "entry": entry})
transcript = job_session_transcript(db, job_id, limit=500)
transcript = queries.job_session_transcript(job_id, limit=500)
for entry in transcript:
if shutdown_event is not None and shutdown_event.is_set():
return
Expand Down Expand Up @@ -656,7 +649,7 @@ def _known_mcp_user_profiles(config: DashboardConfig, *, current_login: str = ""
logins = {login.lower() for login in config.allowed_users | config.admin_users if login}
if current_login:
logins.add(current_login.lower())
for actor_login in list_all_job_actor_logins(config.db):
for actor_login in DashboardQueries(config.db).list_all_job_actor_logins():
login = str(actor_login).strip().lower()
if login:
logins.add(login)
Expand Down Expand Up @@ -879,6 +872,7 @@ async def lifespan(app: FastAPI):
app.state.dashboard_static_snapshot = static_snapshot
app.state.dashboard_shutdown_event = shutdown_event
ensure_webhook_schema = _webhook_schema_initializer(config)
queries = DashboardQueries(config.db)
webhook_repository = WebhookRepository(Database(config.db))

assets_dir = runtime_static_dir / "assets"
Expand Down Expand Up @@ -935,7 +929,7 @@ async def database_unavailable(_: Request, exc: sqlite3.OperationalError) -> JSO

@app.get("/api/health")
def health() -> dict[str, Any]:
metrics = inspect_db_read_only(config.db)
metrics = queries.status()
return {
"ok": bool(metrics.get("db_exists") and metrics.get("schema_ok", True)),
"service": "github-agent-bridge-dashboard",
Expand Down Expand Up @@ -1248,7 +1242,7 @@ async def dashboard_webhooks(request: Request, webhook_path: str = "") -> Respon

@app.get("/api/status")
def api_status(request: Request, profile: dict[str, Any] = Depends(current_profile)) -> dict[str, Any]:
metrics = inspect_db_read_only(config.db)
metrics = queries.status()
dashboard_url, dashboard_url_source = _dashboard_public_url_with_source(request)
admin_actions = [
"retry_job",
Expand Down Expand Up @@ -1396,56 +1390,57 @@ def api_jobs(
limit: int = 50,
) -> dict[str, Any]:
return {
"jobs": list_jobs(
config.db,
status_filter=status_filter,
repo=repo,
thread=thread,
action=action,
intent=intent,
actor=actor,
since=since,
until=until,
"jobs": queries.list_jobs(
JobListFilters(
status=status_filter,
repo=repo,
thread=thread,
action=action,
intent=intent,
actor=actor,
since=since,
until=until,
),
limit=limit,
)
}

@app.get("/api/jobs/actors")
def api_job_actors(_: str = Depends(current_user), limit: int = 100) -> dict[str, Any]:
return {"actors": list_job_actors(config.db, limit=limit)}
return {"actors": queries.list_job_actors(limit=limit)}

@app.get("/api/jobs/{job_id}")
def api_job(job_id: int, _: str = Depends(current_user)) -> dict[str, Any]:
job = get_job_detail(config.db, job_id)
job = queries.get_job_detail(job_id)
if job is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="job_not_found")
return {"job": job}

@app.get("/api/jobs/{job_id}/logs")
def api_job_logs(job_id: int, limit: int = 100, _: str = Depends(current_user)) -> dict[str, Any]:
return {"logs": job_logs(config.db, job_id, limit=limit)}
return {"logs": queries.job_logs(job_id, limit=limit)}

@app.post("/api/jobs/{job_id}/retry")
def api_job_retry(job_id: int, profile: dict[str, Any] = Depends(current_admin_profile)) -> dict[str, Any]:
if get_job_detail(config.db, job_id) is None:
if queries.get_job_detail(job_id) is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="job_not_found")
if not JobQueue(config.db).retry(job_id, actor=str(profile["login"])):
raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="job_not_retryable")
job = get_job_detail(config.db, job_id)
job = queries.get_job_detail(job_id)
return {"job": job, "detail": "job_requeued"}

@app.post("/api/jobs/{job_id}/dismiss")
def api_job_dismiss(job_id: int, profile: dict[str, Any] = Depends(current_admin_profile)) -> dict[str, Any]:
if get_job_detail(config.db, job_id) is None:
if queries.get_job_detail(job_id) is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="job_not_found")
if not JobQueue(config.db).dismiss(job_id, f"dismissed by @{profile['login']}"):
raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="job_not_dismissable")
job = get_job_detail(config.db, job_id)
job = queries.get_job_detail(job_id)
return {"job": job, "detail": "job_dismissed"}

@app.post("/api/jobs/{job_id}/cancel")
async def api_job_cancel(job_id: int, request: Request, profile: dict[str, Any] = Depends(current_profile)) -> dict[str, Any]:
if get_job_detail(config.db, job_id) is None:
if queries.get_job_detail(job_id) is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="job_not_found")
if not can_cancel_job(job_id, profile):
raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="job_cancel_not_allowed")
Expand All @@ -1459,7 +1454,7 @@ async def api_job_cancel(job_id: int, request: Request, profile: dict[str, Any]
result = cancel_running_job(JobQueue(config.db), job_id, actor=str(profile["login"]), reason=reason)
if not result.cancelled:
raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="job_not_running")
job = get_job_detail(config.db, job_id)
job = queries.get_job_detail(job_id)
return {
"job": job,
"detail": "job_cancelled",
Expand All @@ -1470,26 +1465,26 @@ async def api_job_cancel(job_id: int, request: Request, profile: dict[str, Any]

@app.get("/api/jobs/{job_id}/session")
def api_job_session(job_id: int, _: str = Depends(current_user)) -> dict[str, Any]:
session = job_session(config.db, job_id)
session = queries.job_session(job_id)
if session is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="job_not_found")
return {"session": session}

@app.get("/api/jobs/{job_id}/session/events")
def api_job_session_events(job_id: int, after_id: int | None = None, limit: int = 100, _: str = Depends(current_user)) -> dict[str, Any]:
if job_session(config.db, job_id) is None:
if queries.job_session(job_id) is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="job_not_found")
return {"events": job_session_events(config.db, job_id, after_id=after_id, limit=limit)}
return {"events": queries.job_session_events(job_id, after_id=after_id, limit=limit)}

@app.get("/api/jobs/{job_id}/session/transcript")
def api_job_session_transcript(job_id: int, limit: int = 500, _: str = Depends(current_user)) -> dict[str, Any]:
if job_session(config.db, job_id) is None:
if queries.job_session(job_id) is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="job_not_found")
return {"entries": job_session_transcript(config.db, job_id, limit=limit)}
return {"entries": queries.job_session_transcript(job_id, limit=limit)}

@app.get("/api/jobs/{job_id}/session/stream")
def api_job_session_stream(job_id: int, after_id: int | None = None, _: str = Depends(current_user)) -> StreamingResponse:
if job_session(config.db, job_id) is None:
if queries.job_session(job_id) is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="job_not_found")

return StreamingResponse(
Expand All @@ -1500,7 +1495,7 @@ def api_job_session_stream(job_id: int, after_id: int | None = None, _: str = De

@app.get("/api/metrics/summary")
def api_metrics(timezone: str = "UTC", _: str = Depends(current_user)) -> dict[str, Any]:
return {"metrics": metrics_summary(config.db, timezone_name=timezone)}
return {"metrics": queries.metrics_summary(timezone_name=timezone)}

@app.get("/api/processes")
def api_processes(_: str = Depends(current_user)) -> dict[str, Any]:
Expand Down
9 changes: 6 additions & 3 deletions src/github_agent_bridge/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@
record_update_plan,
)
from .cancellation import cancel_running_job
from .dashboard_data import inspect_db_read_only, list_jobs
from .dashboard_data import DashboardQueries, JobListFilters
from .dispatch import FEEDBACK_LEARNING_RULES, GitHubClient, OpenClawDispatcher, RunMode, prompt_rule
from .executor import ExecutorConfig, ExecutorPool
from .models import Notification, utc_now
Expand Down Expand Up @@ -280,7 +280,7 @@ def job_dict(job):


def cmd_status(args: argparse.Namespace) -> int:
metrics = inspect_db_read_only(args.db)
metrics = DashboardQueries(args.db).status()
print(json.dumps({"stats": metrics.get("counts", {}), "oldest_pending_age_seconds": metrics.get("oldest_pending_age_seconds"), "executor_pause": metrics.get("executor_pause", {"paused": False})}, ensure_ascii=False, indent=2))
return 0

Expand All @@ -300,7 +300,10 @@ def cmd_resume_executor(args: argparse.Namespace) -> int:


def cmd_jobs(args: argparse.Namespace) -> int:
rows = list_jobs(args.db, status_filter=args.status, limit=args.limit)
rows = DashboardQueries(args.db).list_jobs(
JobListFilters(status=args.status),
limit=args.limit,
)
print(json.dumps([job_dict(j) for j in rows], ensure_ascii=False, indent=2))
return 0

Expand Down
Loading
Loading