feat(widgets): add backend source adapters and per-widget data endpoint

PR 2 of 4 for configurable dashboard widgets.

- Add grafana_url and prometheus_url settings (config.py + compose/env).
- Create WidgetSource protocol and adapters for jellyfin, backups, grafana,
  prometheus, ssh_task, and static sources.
- Add GET /api/widgets/instances/{id}/data endpoint.
- Extract shared dashboard helpers into domain/dashboard.py so widgets and
  the dashboard router reuse the same logic.
- Add adapter and data-endpoint tests.
- Update apply-progress.md.

Verification: ruff clean; backend pytest 200 passed; frontend lint/build green.
This commit is contained in:
Developer
2026-06-21 10:09:45 +00:00
parent 1a52dfb087
commit 1cd8e926de
10 changed files with 607 additions and 85 deletions
@@ -57,6 +57,8 @@ class Settings(BaseSettings):
prometheus_file_sd_dir: str = "/app/backend/.cache/prometheus-file-sd"
alertmanager_url: str = "http://alertmanager:9093"
alertmanager_webhook_url: str = "" # Optional receiver for alertmanager webhook notifications
grafana_url: str = "http://grafana:3000"
prometheus_url: str = "http://prometheus:9090"
# Remote paths
remote_media_root: str = ""
@@ -0,0 +1,91 @@
"""Dashboard domain helpers shared between routers and widget adapters."""
from __future__ import annotations
import time
from typing import Any
from media_library_viewer_api.models.backups import BackupDashboardSummary
from media_library_viewer_api.services.settings_store import SettingsStore
def _map_sessions_to_activity_rows(sessions: list[dict[str, Any]]) -> list[dict[str, Any]]:
"""Normalize Jellyfin sessions into dashboard activity rows."""
results: list[dict[str, Any]] = []
for session in sessions:
item = session.get("NowPlayingItem") or {}
play_state = session.get("PlayState") or {}
transcoding = session.get("TranscodingInfo") or {}
has_item = bool(item)
series = item.get("SeriesName") or ""
title = (
(f"{series} - {item.get('Name', '')}" if series else item.get("Name", "Unknown"))
if has_item
else "(idle)"
)
if not has_item:
state_label = "idle"
else:
state_label = "paused" if play_state.get("IsPaused") else "playing"
is_transcoding = bool(transcoding)
transcode_type: list[str] = []
if is_transcoding:
if transcoding.get("IsVideoDirect") is False:
transcode_type.append("video")
if transcoding.get("IsAudioDirect") is False:
transcode_type.append("audio")
if not transcode_type:
transcode_type.append("active")
results.append(
{
"user": session.get("UserName") or "Unknown",
"title": title,
"type": item.get("Type", "") if has_item else "",
"state": state_label,
"transcoding": "yes" if is_transcoding else "no",
"transcoding_type": ", ".join(transcode_type),
"device": session.get("DeviceName") or session.get("Client") or "",
"session_id": session.get("Id") or "",
}
)
return results
def build_backup_dashboard_summary(store: SettingsStore) -> BackupDashboardSummary:
"""Compute the backup summary shown on the dashboard."""
jobs = store.list_backup_jobs()
total_jobs = len(jobs)
cutoff = int(time.time()) - (24 * 60 * 60)
recent_runs = []
for job in jobs:
runs = store.list_backup_runs(job_id=job["id"], limit=1)
if runs and runs[0]["started_at"] >= cutoff:
recent_runs.append(runs[0])
successful = sum(1 for r in recent_runs if r["status"] == "success")
success_rate = (successful / len(recent_runs) * 100) if recent_runs else 100.0
alerts = store.list_backup_alerts(acknowledged=False)
active_alerts = len(alerts)
failed_runs = []
for job in jobs:
runs = store.list_backup_runs(job_id=job["id"], status="failure", limit=1)
if runs:
failed_runs.append(runs[0])
last_failed_at = None
if failed_runs:
last_failed_at = max(r["started_at"] for r in failed_runs)
return BackupDashboardSummary(
total_jobs=total_jobs,
success_rate_24h=round(success_rate, 1),
active_alerts=active_alerts,
last_failed_at=last_failed_at,
)
@@ -3,7 +3,6 @@
from __future__ import annotations
import logging
import time
from typing import Any
from fastapi import APIRouter, Depends
@@ -14,6 +13,10 @@ from media_library_viewer_api.dependencies import (
get_settings_store,
get_user_id,
)
from media_library_viewer_api.domain.dashboard import (
_map_sessions_to_activity_rows,
build_backup_dashboard_summary,
)
from media_library_viewer_api.models.backups import BackupDashboardSummary
from media_library_viewer_api.services.settings_store import SettingsStore
@@ -85,50 +88,6 @@ def delete_shortcut(
return {"status": "deleted"}
def _map_sessions_to_activity_rows(sessions: list[dict[str, Any]]) -> list[dict[str, Any]]:
"""Normalize Jellyfin sessions into dashboard activity rows."""
results: list[dict[str, Any]] = []
for session in sessions:
item = session.get("NowPlayingItem") or {}
play_state = session.get("PlayState") or {}
transcoding = session.get("TranscodingInfo") or {}
has_item = bool(item)
series = item.get("SeriesName") or ""
title = (
(f"{series} - {item.get('Name', '')}" if series else item.get("Name", "Unknown")) if has_item else "(idle)"
)
if not has_item:
state_label = "idle"
else:
state_label = "paused" if play_state.get("IsPaused") else "playing"
is_transcoding = bool(transcoding)
transcode_type: list[str] = []
if is_transcoding:
if transcoding.get("IsVideoDirect") is False:
transcode_type.append("video")
if transcoding.get("IsAudioDirect") is False:
transcode_type.append("audio")
if not transcode_type:
transcode_type.append("active")
results.append(
{
"user": session.get("UserName") or "Unknown",
"title": title,
"type": item.get("Type", "") if has_item else "",
"state": state_label,
"transcoding": "yes" if is_transcoding else "no",
"transcoding_type": ", ".join(transcode_type),
"device": session.get("DeviceName") or session.get("Client") or "",
"session_id": session.get("Id") or "",
}
)
return results
@router.get("/activity")
def get_activity(
client: JellyfinClient = Depends(get_jellyfin_client),
@@ -154,38 +113,4 @@ def get_now_playing(
def get_backup_dashboard(
store: SettingsStore = Depends(get_settings_store),
) -> BackupDashboardSummary:
jobs = store.list_backup_jobs()
total_jobs = len(jobs)
# Calculate 24h success rate
cutoff = int(time.time()) - (24 * 60 * 60)
recent_runs = []
for job in jobs:
runs = store.list_backup_runs(job_id=job["id"], limit=1)
if runs and runs[0]["started_at"] >= cutoff:
recent_runs.append(runs[0])
successful = sum(1 for r in recent_runs if r["status"] == "success")
success_rate = (successful / len(recent_runs) * 100) if recent_runs else 100.0
# Active alerts
alerts = store.list_backup_alerts(acknowledged=False)
active_alerts = len(alerts)
# Last failed
failed_runs = []
for job in jobs:
runs = store.list_backup_runs(job_id=job["id"], status="failure", limit=1)
if runs:
failed_runs.append(runs[0])
last_failed_at = None
if failed_runs:
last_failed_at = max(r["started_at"] for r in failed_runs)
return BackupDashboardSummary(
total_jobs=total_jobs,
success_rate_24h=round(success_rate, 1),
active_alerts=active_alerts,
last_failed_at=last_failed_at,
)
return build_backup_dashboard_summary(store)
@@ -1,20 +1,30 @@
"""REST API for dashboard widget instances and registry metadata."""
import logging
import time
from typing import Any
from fastapi import APIRouter, Depends, HTTPException, status
from media_library_viewer_api.dependencies import get_settings_store
from media_library_viewer_api.models.widgets import WidgetInstance, WidgetInstanceInput, WidgetTypeInfo
from media_library_viewer_api.models.widgets import (
WidgetDataResponse,
WidgetInstance,
WidgetInstanceInput,
)
from media_library_viewer_api.services.settings_store import SettingsStore
from media_library_viewer_api.widgets.registry import (
get_widget_info,
list_source_types,
list_widget_types,
validate_config,
)
from media_library_viewer_api.widgets.sources import get_source_adapter
router = APIRouter(prefix="/api/widgets", tags=["widgets"])
logger = logging.getLogger(__name__)
def _registry_for_type(widget_type: str) -> dict[str, Any]:
from media_library_viewer_api.widgets.registry import WIDGET_REGISTRY
@@ -56,7 +66,7 @@ def list_sources() -> list[str]:
@router.get("/types")
def list_types() -> list[WidgetTypeInfo]:
def list_types() -> list[dict[str, Any]]:
"""Return metadata for all registered widget types."""
return [info.model_dump() for info in list_widget_types()]
@@ -111,3 +121,53 @@ def delete_instance(
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Widget not found")
store.delete_widget(widget_id)
return {"status": "deleted"}
@router.get("/instances/{widget_id}/data")
async def fetch_data(
widget_id: str,
store: SettingsStore = Depends(get_settings_store),
) -> dict[str, Any]:
"""Fetch widget data through the registered source adapter."""
widget = store.get_widget(widget_id)
if not widget:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Widget not found")
widget_type = widget["widget_type"]
info = get_widget_info(widget_type)
if info is None:
return WidgetDataResponse(
widget_id=widget_id,
widget_type=widget_type,
data=None,
error=f"Unknown widget type: {widget_type}",
fetched_at=int(time.time()),
).model_dump()
adapter = get_source_adapter(info.source_type)
if adapter is None:
# Defensive: registry should prevent this, but return a safe error.
return WidgetDataResponse(
widget_id=widget_id,
widget_type=widget_type,
data=None,
error=f"No adapter registered for source type: {info.source_type}",
fetched_at=int(time.time()),
).model_dump()
try:
data = await adapter.fetch(widget["config"])
except Exception as exc:
logger.exception("Unhandled adapter exception widget_id=%s", widget_id)
raise HTTPException(
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
detail="Widget data fetch failed",
) from exc
return WidgetDataResponse(
widget_id=widget_id,
widget_type=widget_type,
data=data if "error" not in data else None,
error=data.get("error"),
fetched_at=int(time.time()),
).model_dump()
@@ -0,0 +1,213 @@
"""Widget source adapters.
Each adapter implements a uniform async interface and translates widget
configuration into data for the dashboard. Adapters reuse existing clients,
machine registries, and environment settings; they never accept arbitrary
commands or store credentials.
"""
from __future__ import annotations
import asyncio
import logging
import shlex
from typing import Any, Protocol
import requests
from starlette.requests import Request
from media_library_viewer_api.config import get_settings
from media_library_viewer_api.dependencies import get_jellyfin_client
from media_library_viewer_api.domain.dashboard import (
_map_sessions_to_activity_rows,
build_backup_dashboard_summary,
)
from media_library_viewer_api.routers.tasks import _client_for_machine, _resolve_machine_for_task
from media_library_viewer_api.services.settings_store import get_settings_store
logger = logging.getLogger(__name__)
def _request_with_machine_id(machine_id: str | None = None) -> Request:
"""Build a minimal Starlette Request carrying a machine_id query param."""
query = f"machine_id={machine_id}".encode() if machine_id else b""
return Request({"type": "http", "query_string": query})
class WidgetSource(Protocol):
"""Protocol for widget source adapters."""
source_type: str
async def fetch(self, config: dict[str, Any]) -> dict[str, Any]: ...
class JellyfinWidgetSource:
"""Fetch Jellyfin sessions and map them to activity rows."""
source_type = "jellyfin"
timeout = 10
async def fetch(self, config: dict[str, Any]) -> dict[str, Any]:
try:
request = _request_with_machine_id(config.get("machine_id") or None)
client = await asyncio.wait_for(
asyncio.to_thread(get_jellyfin_client, request),
timeout=self.timeout,
)
sessions = await asyncio.wait_for(
asyncio.to_thread(client.sessions),
timeout=self.timeout,
)
rows = _map_sessions_to_activity_rows(sessions)
return {"sessions": rows}
except asyncio.TimeoutError:
return {"error": "Widget data fetch timed out"}
except Exception as exc:
logger.exception("jellyfin adapter failed")
return {"error": f"Jellyfin data fetch failed: {exc}"}
class BackupsWidgetSource:
"""Compute the backup dashboard summary."""
source_type = "backups"
timeout = 10
async def fetch(self, config: dict[str, Any]) -> dict[str, Any]:
try:
store = get_settings_store()
summary = build_backup_dashboard_summary(store)
return summary.model_dump()
except asyncio.TimeoutError:
return {"error": "Widget data fetch timed out"}
except Exception as exc:
logger.exception("backups adapter failed")
return {"error": f"Backup summary failed: {exc}"}
class GrafanaWidgetSource:
"""Build a Grafana deep-link (no embedding)."""
source_type = "grafana"
timeout = 5
async def fetch(self, config: dict[str, Any]) -> dict[str, Any]:
try:
settings = get_settings()
dashboard_uid = config.get("dashboard_uid")
if not dashboard_uid:
return {"error": "dashboard_uid is required"}
url = f"{settings.grafana_url.rstrip('/')}/d/{dashboard_uid}"
panel_id = config.get("panel_id")
if panel_id is not None:
url = f"{url}?viewPanel={panel_id}"
return {"url": url}
except Exception as exc:
logger.exception("grafana adapter failed")
return {"error": f"Grafana link failed: {exc}"}
class PrometheusWidgetSource:
"""Run a PromQL instant query against Prometheus."""
source_type = "prometheus"
timeout = 10
async def fetch(self, config: dict[str, Any]) -> dict[str, Any]:
try:
settings = get_settings()
promql = config.get("promql")
if not promql:
return {"error": "promql is required"}
url = f"{settings.prometheus_url.rstrip('/')}/api/v1/query"
response = await asyncio.wait_for(
asyncio.to_thread(
requests.get,
url,
params={"query": promql},
timeout=self.timeout,
),
timeout=self.timeout,
)
response.raise_for_status()
payload = response.json()
return {"result": payload.get("data", {})}
except asyncio.TimeoutError:
return {"error": "Widget data fetch timed out"}
except requests.RequestException as exc:
logger.exception("prometheus adapter failed")
return {"error": f"Prometheus query failed: {exc}"}
except Exception as exc:
logger.exception("prometheus adapter failed")
return {"error": f"Prometheus query failed: {exc}"}
class SshTaskWidgetSource:
"""Run a saved task from the registry and return its output."""
source_type = "ssh_task"
timeout = 30
async def fetch(self, config: dict[str, Any]) -> dict[str, Any]:
try:
store = get_settings_store()
task_id = config.get("task_id")
if not task_id:
return {"error": "task_id is required"}
task = store.get_task(task_id)
if not task:
return {"error": f"Task {task_id} not found"}
if not task.get("enabled", True):
return {"error": "Task is disabled"}
machine = _resolve_machine_for_task(store, task, None)
if not machine:
return {"error": "No machine available for this task"}
client = _client_for_machine(store, machine)
task_type = str(task.get("task_type") or "shell").lower()
command = str(task.get("content") or "")
if task_type == "python":
command = f"python3 -c {shlex.quote(command)}"
elif task_type != "shell":
return {"error": f"Unknown task type: {task_type}"}
result = await asyncio.wait_for(
asyncio.to_thread(client.run, command, timeout=self.timeout),
timeout=self.timeout,
)
return {
"exit_status": result.exit_status,
"stdout": result.stdout or "",
"stderr": result.stderr or "",
}
except asyncio.TimeoutError:
return {"error": "Widget data fetch timed out"}
except Exception as exc:
logger.exception("ssh_task adapter failed")
return {"error": f"SSH task failed: {exc}"}
class StaticWidgetSource:
"""Return static text/markdown unchanged."""
source_type = "static"
async def fetch(self, config: dict[str, Any]) -> dict[str, Any]:
return {"text": config.get("text", "")}
SOURCE_REGISTRY: dict[str, WidgetSource] = {
"jellyfin": JellyfinWidgetSource(),
"backups": BackupsWidgetSource(),
"grafana": GrafanaWidgetSource(),
"prometheus": PrometheusWidgetSource(),
"ssh_task": SshTaskWidgetSource(),
"static": StaticWidgetSource(),
}
def get_source_adapter(source_type: str) -> WidgetSource | None:
"""Return the adapter for a source type, or None if unknown."""
return SOURCE_REGISTRY.get(source_type)