diff --git a/docs/persistence.md b/docs/persistence.md index 5145ad6..4486928 100644 --- a/docs/persistence.md +++ b/docs/persistence.md @@ -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 diff --git a/src/github_agent_bridge/backend.py b/src/github_agent_bridge/backend.py index 761ec4b..c0aebf0 100644 --- a/src/github_agent_bridge/backend.py +++ b/src/github_agent_bridge/backend.py @@ -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 @@ -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 @@ -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 @@ -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) @@ -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" @@ -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", @@ -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", @@ -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") @@ -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", @@ -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( @@ -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]: diff --git a/src/github_agent_bridge/cli.py b/src/github_agent_bridge/cli.py index a40a685..6e87af2 100644 --- a/src/github_agent_bridge/cli.py +++ b/src/github_agent_bridge/cli.py @@ -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 @@ -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 @@ -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 diff --git a/src/github_agent_bridge/dashboard_data.py b/src/github_agent_bridge/dashboard_data.py index 002a8d8..ec7c30d 100644 --- a/src/github_agent_bridge/dashboard_data.py +++ b/src/github_agent_bridge/dashboard_data.py @@ -4,6 +4,7 @@ import os import sqlite3 from collections import Counter +from dataclasses import dataclass from datetime import UTC, datetime, time, timedelta from pathlib import Path from typing import Any @@ -12,6 +13,7 @@ from .persistence import Database from .session_events import redact_event_detail from .session_correlation import job_session_metadata +from .sql.migrations import MigrationError, validate_migrations JOB_LIST_ORDER_SQL = """ CASE status @@ -33,15 +35,6 @@ def readonly_connect(db: str | Path) -> sqlite3.Connection: return Database(db).read_only() -def table_exists(con: sqlite3.Connection, name: str) -> bool: - row = con.execute("SELECT 1 FROM sqlite_master WHERE type='table' AND name=?", (name,)).fetchone() - return row is not None - - -def column_exists(con: sqlite3.Connection, table: str, column: str) -> bool: - return any(row["name"] == column for row in con.execute(f"PRAGMA table_info({table})")) - - def parse_utc(value: str | None) -> datetime | None: if not value: return None @@ -134,13 +127,7 @@ def intent_classifier_summary(metadata: dict[str, Any]) -> dict[str, Any] | None } -def jobs_select_sql(con: sqlite3.Connection) -> str: - if not table_exists(con, "job_session_events"): - return """SELECT jobs.*, ( - SELECT running.id FROM jobs AS running - WHERE jobs.status='pending' AND running.work_key=jobs.work_key AND running.status='running' - ORDER BY running.id LIMIT 1 - ) AS blocked_by_job_id FROM jobs""" +def jobs_select_sql(_con: sqlite3.Connection | None = None) -> str: return """ SELECT jobs.*, ( @@ -200,23 +187,66 @@ def job_summary(row: sqlite3.Row) -> dict[str, Any]: } -def where_clause(filters: dict[str, Any]) -> tuple[str, list[Any]]: - clauses: list[str] = [] - args: list[Any] = [] - for column, value in filters.items(): - if value is None: - continue - clauses.append(f"{column}=?") - args.append(value) - return (" WHERE " + " AND ".join(clauses), args) if clauses else ("", args) - - -def actor_where_clause(actor: str | None, *, has_trigger_actor: bool) -> tuple[str, list[Any]]: - if not actor or not actor.strip(): - return "", [] - if not has_trigger_actor: - return " WHERE 1=0", [] - return " WHERE lower(trigger_actor)=lower(?)", [actor.strip().lstrip("@")] +@dataclass(frozen=True) +class JobListFilters: + """Allowlisted filters supported by the dashboard job read model.""" + + status: str | None = None + repo: str | None = None + thread: int | None = None + action: str | None = None + intent: str | None = None + actor: str | None = None + since: str | None = None + until: str | None = None + + def where_clause(self) -> tuple[str, list[Any]]: + filters = ( + ("jobs.status", self.status), + ("jobs.repo", self.repo), + ("jobs.thread", self.thread), + ("jobs.action", self.action), + ("jobs.work_intent", self.intent), + ) + clauses: list[str] = [] + args: list[Any] = [] + for column, value in filters: + if value is None: + continue + clauses.append(f"{column}=?") + args.append(value) + if self.actor and self.actor.strip(): + clauses.append("lower(jobs.trigger_actor)=lower(?)") + args.append(self.actor.strip().lstrip("@")) + if self.since: + clauses.append("jobs.created_at>=?") + args.append(self.since) + if self.until: + clauses.append("jobs.created_at<=?") + args.append(self.until) + return (" WHERE " + " AND ".join(clauses), args) if clauses else ("", args) + + +def where_clause(filters: JobListFilters) -> tuple[str, list[Any]]: + """Build SQL only from the fixed dashboard filter contract.""" + return filters.where_clause() + + +def _prefixed_row( + row: sqlite3.Row, + prefix: str, + *, + include_kind: bool = False, +) -> dict[str, Any] | None: + if row[f"{prefix}_id"] is None: + return None + keys = ("ts", "phase", "summary") + if include_kind: + keys = ("id", "ts", "kind", "phase", "summary", "detail") + item = {key: row[f"{prefix}_{key}"] for key in keys} + if include_kind: + item["age_seconds"] = duration_seconds(item.get("ts")) + return item def inspect_db_read_only(db: str | Path) -> dict[str, Any]: @@ -225,69 +255,90 @@ def inspect_db_read_only(db: str | Path) -> dict[str, Any]: if not path.exists(): return out with readonly_connect(path) as con: - if not table_exists(con, "jobs"): - return out | {"schema_ok": False} + try: + validate_migrations(con) + except (MigrationError, sqlite3.OperationalError) as exc: + return out | {"schema_ok": False, "schema_error": str(exc)} out["schema_ok"] = True counts = {r["status"]: int(r["count"]) for r in con.execute("SELECT status, count(*) count FROM jobs GROUP BY status")} out["counts"] = counts out.update(counts) pending_age = con.execute("SELECT CAST((julianday('now') - julianday(min(created_at))) * 86400 AS INTEGER) age FROM jobs WHERE status='pending'").fetchone()["age"] out["oldest_pending_age_seconds"] = None if pending_age is None else int(pending_age) - if table_exists(con, "state"): - state = {r["key"]: r["value"] for r in con.execute("SELECT key,value FROM state")} - out["last_uid"] = state.get("last_uid") - out["executor_process_tracking_id"] = state.get("executor_process_tracking_id") - out["executor_worker_count"] = int(state.get("executor_worker_count") or 0) - out["executor_pause"] = executor_pause_state_from_raw(state.get(EXECUTOR_PAUSE_STATE_KEY)) - if table_exists(con, "worker_heartbeats"): - heartbeats = [] - for row in con.execute( - """SELECT worker_id, executor_id, pid, last_seen, active_job_id, loop_state, - recent_error_count, - CAST((julianday('now') - julianday(last_seen)) * 86400 AS INTEGER) age_seconds - FROM worker_heartbeats ORDER BY worker_id""" - ): - item = dict(row) - item["age_seconds"] = max(0, int(item["age_seconds"] or 0)) - heartbeats.append(item) - out["worker_heartbeats"] = heartbeats - if table_exists(con, "feedback_rule_proposals"): - knowledge_counts = { - r["status"]: int(r["count"]) - for r in con.execute("SELECT status, count(*) count FROM feedback_rule_proposals GROUP BY status") - } - out["knowledge"] = { - "proposed": knowledge_counts.get("proposed", 0), - "approved": knowledge_counts.get("approved", 0), - "rejected": knowledge_counts.get("rejected", 0), - "errors": knowledge_counts.get("error", 0), - } - if table_exists(con, "quarantined_notifications"): - unresolved_quarantines = con.execute( - "SELECT count(*) count FROM quarantined_notifications WHERE resolved_at IS NULL" - ).fetchone()["count"] - out["quarantined_notifications"] = int(unresolved_quarantines) - latest_quarantine = con.execute( - """ - SELECT id, uid, message_id, subject, reason, error, created_at - FROM quarantined_notifications - WHERE resolved_at IS NULL - ORDER BY created_at DESC, id DESC - LIMIT 1 - """ - ).fetchone() - if latest_quarantine: - out["latest_quarantined_notification"] = dict(latest_quarantine) - if table_exists(con, "worklog"): - last_log = con.execute("SELECT ts, phase, summary FROM worklog ORDER BY id DESC LIMIT 1").fetchone() - if last_log: - out["last_worklog"] = dict(last_log) + state = {r["key"]: r["value"] for r in con.execute("SELECT key,value FROM state")} + out["last_uid"] = state.get("last_uid") + out["executor_process_tracking_id"] = state.get("executor_process_tracking_id") + out["executor_worker_count"] = int(state.get("executor_worker_count") or 0) + out["executor_pause"] = executor_pause_state_from_raw(state.get(EXECUTOR_PAUSE_STATE_KEY)) + heartbeats = [] + for row in con.execute( + """SELECT worker_id, executor_id, pid, last_seen, active_job_id, loop_state, + recent_error_count, + CAST((julianday('now') - julianday(last_seen)) * 86400 AS INTEGER) age_seconds + FROM worker_heartbeats ORDER BY worker_id""" + ): + item = dict(row) + item["age_seconds"] = max(0, int(item["age_seconds"] or 0)) + heartbeats.append(item) + out["worker_heartbeats"] = heartbeats + knowledge_counts = { + r["status"]: int(r["count"]) + for r in con.execute("SELECT status, count(*) count FROM feedback_rule_proposals GROUP BY status") + } + out["knowledge"] = { + "proposed": knowledge_counts.get("proposed", 0), + "approved": knowledge_counts.get("approved", 0), + "rejected": knowledge_counts.get("rejected", 0), + "errors": knowledge_counts.get("error", 0), + } + unresolved_quarantines = con.execute( + "SELECT count(*) count FROM quarantined_notifications WHERE resolved_at IS NULL" + ).fetchone()["count"] + out["quarantined_notifications"] = int(unresolved_quarantines) + latest_quarantine = con.execute( + """ + SELECT id, uid, message_id, subject, reason, error, created_at + FROM quarantined_notifications + WHERE resolved_at IS NULL + ORDER BY created_at DESC, id DESC + LIMIT 1 + """ + ).fetchone() + if latest_quarantine: + out["latest_quarantined_notification"] = dict(latest_quarantine) + last_log = con.execute("SELECT ts, phase, summary FROM worklog ORDER BY id DESC LIMIT 1").fetchone() + if last_log: + out["last_worklog"] = dict(last_log) running_rows = con.execute( """ - SELECT id, work_key, work_intent, locked_by, attempts, started_at, updated_at, metadata_json + WITH ranked_worklog AS ( + SELECT *, row_number() OVER (PARTITION BY job_id ORDER BY id DESC) AS row_number + FROM worklog + ), ranked_progress AS ( + SELECT *, row_number() OVER (PARTITION BY job_id, kind ORDER BY id DESC) AS row_number + FROM job_progress + WHERE kind IN ('semantic', 'visible') + ) + SELECT jobs.id, jobs.work_key, jobs.work_intent, jobs.locked_by, jobs.attempts, + jobs.started_at, jobs.updated_at, jobs.metadata_json, + worklog.id AS worklog_id, worklog.ts AS worklog_ts, + worklog.phase AS worklog_phase, worklog.summary AS worklog_summary, + worklog.detail AS worklog_detail, + semantic.id AS semantic_id, semantic.ts AS semantic_ts, + semantic.kind AS semantic_kind, semantic.phase AS semantic_phase, + semantic.summary AS semantic_summary, semantic.detail AS semantic_detail, + visible.id AS visible_id, visible.ts AS visible_ts, + visible.kind AS visible_kind, visible.phase AS visible_phase, + visible.summary AS visible_summary, visible.detail AS visible_detail FROM jobs - WHERE status='running' - ORDER BY id + LEFT JOIN ranked_worklog AS worklog + ON worklog.job_id=jobs.id AND worklog.row_number=1 + LEFT JOIN ranked_progress AS semantic + ON semantic.job_id=jobs.id AND semantic.kind='semantic' AND semantic.row_number=1 + LEFT JOIN ranked_progress AS visible + ON visible.job_id=jobs.id AND visible.kind='visible' AND visible.row_number=1 + WHERE jobs.status='running' + ORDER BY jobs.id """ ).fetchall() out["running_jobs"] = [ @@ -301,9 +352,9 @@ def inspect_db_read_only(db: str | Path) -> dict[str, Any]: "runtime_process": json.loads(row["metadata_json"] or "{}").get("runtime_process"), "age_seconds": duration_seconds(row["started_at"]), "idle_seconds": duration_seconds(row["updated_at"]), - "last_worklog": _last_worklog(con, int(row["id"])), - "semantic_progress": _latest_job_progress(con, int(row["id"]), "semantic"), - "visible_progress": _latest_job_progress(con, int(row["id"]), "visible"), + "last_worklog": _prefixed_row(row, "worklog"), + "semantic_progress": _prefixed_row(row, "semantic", include_kind=True), + "visible_progress": _prefixed_row(row, "visible", include_kind=True), } for row in running_rows ] @@ -328,20 +379,18 @@ def list_jobs( if not path.exists(): return [] with readonly_connect(path) as con: - if not table_exists(con, "jobs"): - return [] - filters = {"status": status_filter, "repo": repo, "thread": thread, "action": action, "work_intent": intent} - where, args = where_clause(filters) - actor_where, actor_args = actor_where_clause(actor, has_trigger_actor=column_exists(con, "jobs", "trigger_actor")) - if actor_where: - where += actor_where.replace(" WHERE ", " AND ", 1) if where else actor_where - args.extend(actor_args) - if since: - where += " AND created_at >= ?" if where else " WHERE created_at >= ?" - args.append(since) - if until: - where += " AND created_at <= ?" if where else " WHERE created_at <= ?" - args.append(until) + where, args = where_clause( + JobListFilters( + status=status_filter, + repo=repo, + thread=thread, + action=action, + intent=intent, + actor=actor, + since=since, + until=until, + ) + ) args.append(coerce_limit(limit)) rows = con.execute(f"{jobs_select_sql(con)}{where} ORDER BY {JOB_LIST_ORDER_SQL} LIMIT ?", args).fetchall() return [job_summary(row) for row in rows] @@ -352,8 +401,6 @@ def list_job_actors(db: str | Path, *, limit: int = 100) -> list[dict[str, Any]] if not path.exists(): return [] with readonly_connect(path) as con: - if not table_exists(con, "jobs") or not column_exists(con, "jobs", "trigger_actor"): - return [] rows = con.execute( """ SELECT @@ -385,8 +432,6 @@ def list_all_job_actor_logins(db: str | Path) -> list[str]: if not path.exists(): return [] with readonly_connect(path) as con: - if not table_exists(con, "jobs") or not column_exists(con, "jobs", "trigger_actor"): - return [] rows = con.execute( """ SELECT lower(trigger_actor) AS login @@ -404,8 +449,6 @@ def get_job_detail(db: str | Path, job_id: int) -> dict[str, Any] | None: if not path.exists(): return None with readonly_connect(path) as con: - if not table_exists(con, "jobs"): - return None row = con.execute(f"{jobs_select_sql(con)} WHERE jobs.id=?", (job_id,)).fetchone() if row is None: return None @@ -418,14 +461,14 @@ def get_job_detail(db: str | Path, job_id: int) -> dict[str, Any] | None: "SELECT id, ts, phase, summary, detail FROM worklog WHERE job_id=? ORDER BY id", (job_id,), ).fetchall() - ] if table_exists(con, "worklog") else [] + ] job["progress"] = [ dict(progress) for progress in con.execute( "SELECT id, ts, kind, phase, summary, detail FROM job_progress WHERE job_id=? ORDER BY id", (job_id,), ).fetchall() - ] if table_exists(con, "job_progress") else [] + ] job["runs"] = [ dict(run) for run in con.execute( @@ -433,7 +476,7 @@ def get_job_detail(db: str | Path, job_id: int) -> dict[str, Any] | None: FROM job_runs WHERE job_id=? ORDER BY attempt""", (job_id,), ).fetchall() - ] if table_exists(con, "job_runs") else [] + ] job["coalesced_notifications"] = [ { "id": row["id"], @@ -449,7 +492,7 @@ def get_job_detail(db: str | Path, job_id: int) -> dict[str, Any] | None: "SELECT * FROM coalesced_notifications WHERE job_id=? ORDER BY id", (job_id,), ).fetchall() - ] if table_exists(con, "coalesced_notifications") else [] + ] return job @@ -458,8 +501,6 @@ def job_logs(db: str | Path, job_id: int, limit: int = 100) -> list[dict[str, An if not path.exists(): return [] with readonly_connect(path) as con: - if not table_exists(con, "worklog"): - return [] rows = con.execute( "SELECT id, ts, phase, summary, detail FROM worklog WHERE job_id=? ORDER BY id DESC LIMIT ?", (job_id, coerce_limit(limit, maximum=500)), @@ -472,8 +513,6 @@ def latest_job_progress(db: str | Path, job_id: int, kind: str | None = None) -> if not path.exists(): return None with readonly_connect(path) as con: - if not table_exists(con, "job_progress"): - return None return _latest_job_progress(con, job_id, kind) @@ -501,8 +540,6 @@ def job_session_events(db: str | Path, job_id: int, *, after_id: int | None = No if not path.exists(): return [] with readonly_connect(path) as con: - if not table_exists(con, "job_session_events"): - return [] where = "job_id=?" args: list[Any] = [job_id] if after_id is not None: @@ -739,15 +776,13 @@ def metrics_summary(db: str | Path, *, timezone_name: str = "UTC") -> dict[str, return {"db_exists": False, "status_counts": {}, "runtime_seconds": {}} timezone = dashboard_timezone(timezone_name) with readonly_connect(path) as con: - if not table_exists(con, "jobs"): - return {"db_exists": True, "schema_ok": False, "status_counts": {}, "runtime_seconds": {}} rows = con.execute("SELECT status, repo, action, work_intent, created_at, started_at, finished_at FROM jobs").fetchall() run_rows = con.execute( """SELECT job_runs.started_at, job_runs.finished_at, jobs.work_intent FROM job_runs JOIN jobs ON jobs.id=job_runs.job_id WHERE job_runs.finished_at IS NOT NULL""" - ).fetchall() if table_exists(con, "job_runs") else rows + ).fetchall() status_counts = Counter(row["status"] for row in rows) by_repo = Counter(row["repo"] or "unknown" for row in rows) by_action = Counter(row["action"] for row in rows) @@ -882,19 +917,7 @@ def nearest_rank(values: list[int], percentile: float) -> int: return values[index] -def _last_worklog(con: sqlite3.Connection, job_id: int) -> dict[str, Any] | None: - if not table_exists(con, "worklog"): - return None - row = con.execute( - "SELECT ts, phase, summary FROM worklog WHERE job_id=? ORDER BY id DESC LIMIT 1", - (job_id,), - ).fetchone() - return dict(row) if row else None - - def _latest_job_progress(con: sqlite3.Connection, job_id: int, kind: str | None = None) -> dict[str, Any] | None: - if not table_exists(con, "job_progress"): - return None where = "job_id=?" args: list[Any] = [job_id] if kind: @@ -909,3 +932,85 @@ def _latest_job_progress(con: sqlite3.Connection, job_id: int, kind: str | None progress = dict(row) progress["age_seconds"] = duration_seconds(progress.get("ts")) return progress + + +class DashboardQueries: + """Read-only dashboard persistence boundary. + + The compatibility functions in this module remain as a public API shim, but + HTTP handlers and services depend on this object so connection policy, + filters, SQL, and row mapping stay outside application orchestration. + """ + + def __init__(self, db: str | Path) -> None: + self.path = Path(db).expanduser() + + def status(self) -> dict[str, Any]: + return inspect_db_read_only(self.path) + + def list_jobs( + self, + filters: JobListFilters | None = None, + *, + limit: int = 50, + ) -> list[dict[str, Any]]: + selected = filters or JobListFilters() + return list_jobs( + self.path, + status_filter=selected.status, + repo=selected.repo, + thread=selected.thread, + action=selected.action, + intent=selected.intent, + actor=selected.actor, + since=selected.since, + until=selected.until, + limit=limit, + ) + + def list_job_actors(self, *, limit: int = 100) -> list[dict[str, Any]]: + return list_job_actors(self.path, limit=limit) + + def list_all_job_actor_logins(self) -> list[str]: + return list_all_job_actor_logins(self.path) + + def get_job_detail(self, job_id: int) -> dict[str, Any] | None: + return get_job_detail(self.path, job_id) + + def job_logs(self, job_id: int, *, limit: int = 100) -> list[dict[str, Any]]: + return job_logs(self.path, job_id, limit=limit) + + def latest_job_progress( + self, + job_id: int, + kind: str | None = None, + ) -> dict[str, Any] | None: + return latest_job_progress(self.path, job_id, kind) + + def job_session(self, job_id: int) -> dict[str, Any] | None: + return job_session(self.path, job_id) + + def job_session_events( + self, + job_id: int, + *, + after_id: int | None = None, + limit: int = 100, + ) -> list[dict[str, Any]]: + return job_session_events( + self.path, + job_id, + after_id=after_id, + limit=limit, + ) + + def job_session_transcript( + self, + job_id: int, + *, + limit: int = 500, + ) -> list[dict[str, Any]]: + return job_session_transcript(self.path, job_id, limit=limit) + + def metrics_summary(self, *, timezone_name: str = "UTC") -> dict[str, Any]: + return metrics_summary(self.path, timezone_name=timezone_name) diff --git a/src/github_agent_bridge/monitor.py b/src/github_agent_bridge/monitor.py index 074b558..8f50bb7 100644 --- a/src/github_agent_bridge/monitor.py +++ b/src/github_agent_bridge/monitor.py @@ -11,7 +11,7 @@ from pathlib import Path from typing import Any -from .dashboard_data import inspect_db_read_only +from .dashboard_data import DashboardQueries from .observability import DEFAULT_PROCESS_SAMPLE_RETENTION_SECONDS, recent_process_samples, record_monitor_observation from .process_inspection import direct_children, process_identity_matches @@ -155,7 +155,7 @@ def _direct_children(pid: int) -> list[dict[str, Any]]: def inspect_db(path: str | Path) -> dict[str, Any]: - return inspect_db_read_only(path) + return DashboardQueries(path).status() def _package_version() -> str: diff --git a/tests/test_dashboard_queries.py b/tests/test_dashboard_queries.py new file mode 100644 index 0000000..576381d --- /dev/null +++ b/tests/test_dashboard_queries.py @@ -0,0 +1,87 @@ +import inspect + +from github_agent_bridge import backend, cli, monitor +from github_agent_bridge.dashboard_data import DashboardQueries, JobListFilters +from github_agent_bridge.models import Notification +from github_agent_bridge.persistence import Database +from github_agent_bridge.policy import Policy +from github_agent_bridge.queue import JobQueue + + +def _notification() -> Notification: + return Notification( + uid=1, + message_id="", + subject="Re: [gisce/github-agent-bridge] dashboard", + from_addr="Edu ", + body=( + "@giscebot inspect " + "https://github.com/gisce/github-agent-bridge/issues/260#issuecomment-1" + ), + auth={"spf": True, "dkim": True, "dmarc": True}, + ) + + +def test_job_list_filters_are_allowlisted_and_values_are_bound(): + where, args = JobListFilters( + status="done' OR 1=1 --", + repo="gisce/github-agent-bridge", + thread=260, + action="reply_comment", + intent="work_allowed", + actor="@ecarreras", + since="2026-10-01T00:00:00Z", + until="2026-10-31T23:59:59Z", + ).where_clause() + + assert where == ( + " WHERE jobs.status=? AND jobs.repo=? AND jobs.thread=? " + "AND jobs.action=? AND jobs.work_intent=? " + "AND lower(jobs.trigger_actor)=lower(?) " + "AND jobs.created_at>=? AND jobs.created_at<=?" + ) + assert args[0] == "done' OR 1=1 --" + assert "OR 1=1" not in where + assert args[5] == "ecarreras" + + +def test_dashboard_job_list_executes_one_read_query(tmp_path, monkeypatch): + db = tmp_path / "bridge.sqlite3" + queue = JobQueue(db) + queue.enqueue( + _notification(), + Policy(trusted_orgs={"gisce"}, bot_logins={"giscebot"}), + ) + statements: list[str] = [] + original_read_only = Database.read_only + + def traced_read_only(database): + connection = original_read_only(database) + connection.set_trace_callback(statements.append) + return connection + + monkeypatch.setattr(Database, "read_only", traced_read_only) + + rows = DashboardQueries(db).list_jobs( + JobListFilters(repo="gisce/github-agent-bridge", thread=260), + limit=20, + ) + + read_queries = [ + statement + for statement in statements + if statement.lstrip().upper().startswith(("SELECT", "WITH")) + ] + assert len(rows) == 1 + assert len(read_queries) == 1 + assert "sqlite_master" not in read_queries[0] + assert "PRAGMA table_info" not in read_queries[0] + + +def test_http_cli_and_monitor_use_dashboard_query_boundary(): + for module in (backend, cli, monitor): + source = inspect.getsource(module) + assert "sqlite3.connect(" not in source + assert ".execute(" not in source + + assert "DashboardQueries(config.db)" in inspect.getsource(backend.create_app)