feat(services): backend service registry foundation (encryption, definitions, CRUD)

PR 1 of 4 for the runtime service registry change.

- Add Fernet encryption helper (services/secrets.py) with a required
  MANAGE_ENCRYPTION_KEY; validate it on startup.
- Add closed integrations/ registry with Pydantic config + widget-config
  definitions for grafana, prometheus, jellyfin, nextcloud, and ssh_tasks.
- Add services + service_task_runs tables and SettingsStore CRUD with
  cascade-delete (defensive until widgets carry service_id).
- Add /api/services/types and /api/services/instances CRUD (encrypted secrets,
  secrets_set flags only; never plaintext).
- Declare cryptography as a direct dependency.
- Require MANAGE_ENCRYPTION_KEY in compose + .env.example + README.
- Add 25 backend tests (registry, encryption, CRUD, cascade, task-run history).

Verification: ruff clean; pytest 225 passed; frontend lint/build green.
This commit is contained in:
Developer
2026-06-22 12:56:03 +00:00
parent d1819c0186
commit 8cdeadd6dd
20 changed files with 1470 additions and 0 deletions
@@ -0,0 +1,97 @@
"""Encryption-at-rest for service secrets.
Service API keys / tokens are stored encrypted in the ``services.secrets_json``
column. Encryption uses Fernet (symmetric authenticated encryption) with a single
master key provided via the ``MANAGE_ENCRYPTION_KEY`` environment variable.
* The key **must** be a urlsafe base64-encoded 32-byte value (Fernet format).
* The key is **always required** — there is no development fallback, so secrets
are never accidentally stored in plaintext.
* Secrets are encrypted field-by-field; the ``"which secrets are set"`` metadata
can be derived from the ciphertext blob without decrypting.
"""
from __future__ import annotations
import os
from functools import lru_cache
from cryptography.fernet import Fernet, InvalidToken
ENCRYPTION_KEY_ENV = "MANAGE_ENCRYPTION_KEY"
class EncryptionKeyError(RuntimeError):
"""Raised when the encryption key is missing or invalid."""
@lru_cache(maxsize=1)
def get_encryption_key() -> bytes:
"""Return the raw Fernet key, or raise if missing/invalid.
The result is cached for the process lifetime. Tests should call
:func:`reset_encryption_key_cache` after changing the environment.
"""
raw = os.environ.get(ENCRYPTION_KEY_ENV)
if not raw:
raise EncryptionKeyError(
f"{ENCRYPTION_KEY_ENV} is required to store service secrets"
)
key = raw.strip().encode()
try:
Fernet(key)
except (ValueError, TypeError) as exc: # pragma: no cover - validated by tests
raise EncryptionKeyError(
f"{ENCRYPTION_KEY_ENV} must be a valid Fernet key"
) from exc
return key
def reset_encryption_key_cache() -> None:
"""Drop the cached encryption key (used by tests that swap keys)."""
get_encryption_key.cache_clear()
def _fernet() -> Fernet:
return Fernet(get_encryption_key())
def encrypt_value(plaintext: str) -> str:
"""Encrypt a single secret value and return the ciphertext string."""
return _fernet().encrypt(plaintext.encode()).decode()
def decrypt_value(ciphertext: str) -> str:
"""Decrypt a single ciphertext value."""
try:
return _fernet().decrypt(ciphertext.encode()).decode()
except InvalidToken as exc:
raise EncryptionKeyError("Service secret could not be decrypted") from exc
def encrypt_secrets(values: dict[str, str]) -> dict[str, str]:
"""Encrypt every provided secret value."""
fernet = _fernet()
return {key: fernet.encrypt(value.encode()).decode() for key, value in values.items()}
def decrypt_secrets(blob: dict[str, str]) -> dict[str, str]:
"""Decrypt every secret value in a blob."""
fernet = _fernet()
result: dict[str, str] = {}
for key, ciphertext in blob.items():
try:
result[key] = fernet.decrypt(ciphertext.encode()).decode()
except InvalidToken as exc:
raise EncryptionKeyError(f"Service secret '{key}' could not be decrypted") from exc
return result
def generate_development_key() -> str:
"""Return a freshly generated Fernet key (helper for operators/docs)."""
return Fernet.generate_key().decode()
def validate_encryption_key() -> None:
"""Eagerly validate that the encryption key is present and well-formed."""
get_encryption_key() # raises EncryptionKeyError on failure
@@ -226,6 +226,45 @@ class SettingsStore:
""")
conn.execute("CREATE INDEX IF NOT EXISTS idx_backup_alerts_job_id ON backup_alerts(job_id)")
conn.execute("CREATE INDEX IF NOT EXISTS idx_backup_alerts_acknowledged ON backup_alerts(acknowledged)")
conn.execute(
"""
CREATE TABLE IF NOT EXISTS services (
id TEXT PRIMARY KEY,
service_type TEXT NOT NULL,
name TEXT NOT NULL,
config_json TEXT NOT NULL DEFAULT '{}',
secrets_json TEXT NOT NULL DEFAULT '{}',
enabled INTEGER NOT NULL DEFAULT 1,
created_at INTEGER NOT NULL,
updated_at INTEGER NOT NULL
)
"""
)
conn.execute("CREATE INDEX IF NOT EXISTS idx_services_type ON services(service_type)")
conn.execute(
"""
CREATE TABLE IF NOT EXISTS service_task_runs (
id TEXT PRIMARY KEY,
task_id TEXT NOT NULL,
service_id TEXT NOT NULL,
status TEXT NOT NULL,
exit_status INTEGER,
duration_ms INTEGER,
stdout_tail TEXT NOT NULL DEFAULT '',
stderr_tail TEXT NOT NULL DEFAULT '',
error TEXT NOT NULL DEFAULT '',
created_at INTEGER NOT NULL
)
"""
)
conn.execute(
"CREATE INDEX IF NOT EXISTS idx_service_task_runs_service "
"ON service_task_runs(service_id, created_at DESC)"
)
conn.execute(
"CREATE INDEX IF NOT EXISTS idx_service_task_runs_task "
"ON service_task_runs(task_id, created_at DESC)"
)
@staticmethod
def _normalize_services(value: Any, fallback: list[str] | None = None) -> list[str]:
@@ -1476,6 +1515,211 @@ class SettingsStore:
with self.connect() as conn:
conn.execute("DELETE FROM dashboard_widgets WHERE id = ?", (widget_id,))
# ------------------------------------------------------------------
# Service registry
# ------------------------------------------------------------------
def _row_to_service(self, row: sqlite3.Row) -> dict[str, Any]:
secrets_blob = json.loads(row["secrets_json"] or "{}")
return {
"id": row["id"],
"service_type": row["service_type"],
"name": row["name"],
"config": json.loads(row["config_json"] or "{}"),
"secrets": secrets_blob,
"enabled": bool(row["enabled"]),
"created_at": row["created_at"],
"updated_at": row["updated_at"],
}
def list_services(self, service_type: str | None = None) -> list[dict[str, Any]]:
self.init_schema()
with self.connect() as conn:
if service_type:
rows = conn.execute(
"SELECT * FROM services WHERE service_type = ? ORDER BY name ASC",
(service_type,),
).fetchall()
else:
rows = conn.execute("SELECT * FROM services ORDER BY name ASC").fetchall()
return [self._row_to_service(row) for row in rows]
def get_service(self, service_id: str) -> dict[str, Any] | None:
self.init_schema()
with self.connect() as conn:
row = conn.execute(
"SELECT * FROM services WHERE id = ?", (service_id,)
).fetchone()
return self._row_to_service(row) if row else None
def _normalize_service_payload(
self,
payload: dict[str, Any],
service_id: str | None = None,
) -> dict[str, Any]:
current = self.get_service(service_id) if service_id else None
service_id = (
str(payload.get("id") or service_id or uuid.uuid4().hex[:12]).strip()
or uuid.uuid4().hex[:12]
)
service_type = str(
payload.get("service_type") or (current or {}).get("service_type", "")
).strip()
name = str(payload.get("name") or (current or {}).get("name", "") or "").strip()
config = payload.get("config", (current or {}).get("config", {}))
if not isinstance(config, dict):
config = {}
enabled = bool(payload.get("enabled", (current or {}).get("enabled", True)))
return {
"id": service_id,
"service_type": service_type,
"name": name,
"config": config,
"enabled": enabled,
}
def upsert_service(
self,
payload: dict[str, Any],
secret_values: dict[str, str] | None = None,
service_id: str | None = None,
) -> dict[str, Any]:
"""Insert or update a service instance.
``secret_values`` carries plaintext secrets to encrypt and store. A key
absent from ``secret_values`` preserves the existing ciphertext; a key
mapped to an empty string clears it.
"""
self.init_schema()
service = self._normalize_service_payload(payload, service_id)
now = int(time.time())
existing = self.get_service(service["id"])
secrets_blob: dict[str, str]
if existing is not None:
secrets_blob = dict(existing["secrets"])
else:
secrets_blob = {}
if secret_values:
from media_library_viewer_api.services.secrets import encrypt_value
for key, value in secret_values.items():
if value == "":
secrets_blob.pop(key, None)
else:
secrets_blob[key] = encrypt_value(value)
with self.connect() as conn:
created_at = int(existing["created_at"]) if existing else now
conn.execute(
"""
INSERT INTO services (
id, service_type, name, config_json, secrets_json,
enabled, created_at, updated_at
)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(id) DO UPDATE SET
service_type = excluded.service_type,
name = excluded.name,
config_json = excluded.config_json,
secrets_json = excluded.secrets_json,
enabled = excluded.enabled,
updated_at = excluded.updated_at
""",
(
service["id"],
service["service_type"],
service["name"],
json.dumps(service["config"]),
json.dumps(secrets_blob),
1 if service["enabled"] else 0,
created_at,
now,
),
)
return self.get_service(service["id"]) or service
def delete_service(self, service_id: str) -> None:
"""Delete a service and cascade-delete widgets referencing it."""
self.init_schema()
with self.connect() as conn:
# The service_id column on dashboard_widgets is added in a later
# slice; only cascade when it is present.
widget_cols = {row[1] for row in conn.execute("PRAGMA table_info(dashboard_widgets)").fetchall()}
if "service_id" in widget_cols:
conn.execute(
"DELETE FROM dashboard_widgets WHERE service_id = ?",
(service_id,),
)
conn.execute("DELETE FROM services WHERE id = ?", (service_id,))
def record_service_task_run(self, payload: dict[str, Any]) -> dict[str, Any]:
"""Append a service task run history row."""
self.init_schema()
run_id = str(payload.get("id") or uuid.uuid4().hex[:12])
now = int(time.time())
with self.connect() as conn:
conn.execute(
"""
INSERT INTO service_task_runs (
id, task_id, service_id, status, exit_status, duration_ms,
stdout_tail, stderr_tail, error, created_at
)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
run_id,
str(payload.get("task_id") or ""),
str(payload.get("service_id") or ""),
str(payload.get("status") or "error"),
payload.get("exit_status"),
payload.get("duration_ms"),
str(payload.get("stdout_tail") or "")[:8000],
str(payload.get("stderr_tail") or "")[:8000],
str(payload.get("error") or "")[:1000],
int(payload.get("created_at") or now),
),
)
return {"id": run_id}
def list_service_task_runs(
self,
service_id: str | None = None,
task_id: str | None = None,
limit: int = 50,
) -> list[dict[str, Any]]:
self.init_schema()
clauses: list[str] = []
params: list[Any] = []
if service_id:
clauses.append("service_id = ?")
params.append(service_id)
if task_id:
clauses.append("task_id = ?")
params.append(task_id)
where = ("WHERE " + " AND ".join(clauses)) if clauses else ""
params.append(int(limit))
with self.connect() as conn:
rows = conn.execute(
f"SELECT * FROM service_task_runs {where} ORDER BY created_at DESC LIMIT ?",
params,
).fetchall()
return [
{
"id": row["id"],
"task_id": row["task_id"],
"service_id": row["service_id"],
"status": row["status"],
"exit_status": row["exit_status"],
"duration_ms": row["duration_ms"],
"stdout_tail": row["stdout_tail"],
"stderr_tail": row["stderr_tail"],
"error": row["error"],
"created_at": row["created_at"],
}
for row in rows
]
_store: SettingsStore | None = None