diff --git a/docs/capability-domains.md b/docs/capability-domains.md index 20bddb7..c886102 100644 --- a/docs/capability-domains.md +++ b/docs/capability-domains.md @@ -155,13 +155,13 @@ Diagnostics live under `slopsmith.note_detection_capability.v1` and contain prov ## Jobs Control Plane -The jobs slice promotes `jobs` as a privileged core provider-coordinator implemented by [static/capabilities/jobs.js](../static/capabilities/jobs.js). It coordinates long-running plugin work such as conversion, import, cache-building, update, preview, and studio-style background tasks while providers keep the actual file writes, subprocesses, downloads, and private payloads inside their own code. +The jobs slice promotes `jobs` as a privileged core provider-coordinator implemented by [static/capabilities/jobs.js](../static/capabilities/jobs.js), with a backend visibility companion exposed through `context["jobs"]`, `/api/jobs`, `/api/jobs/providers`, `/api/jobs/{id}`, and `/ws/jobs`. It coordinates long-running plugin work such as conversion, import, cache-building, update, preview, and studio-style background tasks while providers keep the actual file writes, subprocesses, downloads, and private payloads inside their own code. Backend plugin queues can register/adopt/progress/settle redaction-safe job summaries without requiring their browser screen to be loaded. -The public command surface is `register-provider`, `unregister-provider`, `list-providers`, `enqueue`, `list`, `inspect`, `cancel`, `pause`, `resume`, `retry`, and `record-bridge-hit`. Provider operations are `job.enqueue`, `job.status`, `job.cancel`, `job.pause`, `job.resume`, `job.retry`, and `job.recover`. `list` and `inspect` are prompt-free and side-effect-free: they never trigger provider work, writes, downloads, subprocesses, or external calls. Fresh privileged `enqueue` and `retry` requests require an explicit `authorization: "user-action"` or a matching approved-continuation scope before provider callbacks run. +The browser command surface is `register-provider`, `unregister-provider`, `list-providers`, `enqueue`, `adopt`, `list`, `inspect`, `cancel`, `pause`, `resume`, `retry`, and `record-bridge-hit`. Provider operations are `job.enqueue`, `job.status`, `job.cancel`, `job.pause`, `job.resume`, `job.retry`, `job.recover`, and `job.adopt`. `list` and `inspect` are prompt-free and side-effect-free: they never trigger provider work, writes, downloads, subprocesses, or external calls. Fresh privileged `enqueue` and `retry` requests require an explicit `authorization: "user-action"` or a matching approved-continuation scope before provider callbacks run. The backend companion is deliberately reporting-first in this slice: provider backends call `register_provider`, `adopt`, `update_progress`, `complete`, `fail`, `cancelled`, `mark_provider_unavailable`, and `record_bridge_hit`; `/api/jobs/{id}/cancel` and `/api/jobs/{id}/retry` return `unsupported-operation` until a backend dispatch/authorization bridge exists. Jobs are scheduled by provider capacity and priority (`user-approved-interactive` before `background-maintenance`). State transitions are explicit (`queued`, `running`, `paused`, `cancellation-requested`, terminal cancelled/completed/failed/provider-unavailable/orphaned), and outcomes use the shared canonical vocabulary: `handled`, `queued`, `denied`, `user-action-required`, `unavailable`, `no-owner`, `no-handler`, `no-target`, `unsupported-command`, `unsupported-operation`, `incompatible`, `incompatible-version`, `provider-selection-required`, `validation-failed`, `stale`, `cancelled`, `completed`, `failed`, `timeout`, and `retry-started`. -Diagnostics live under `slopsmith.jobs.diagnostics.v1`. Snapshots include provider summaries, selected/default providers, active/queued/paused/recent-terminal jobs, bounded per-job history, recent outcomes, bridge hits, limits, and notes. Browser persistence is limited to selected provider choices and provider-declared recoverable job references; it does not store raw payloads, non-recoverable active job state, file paths, command lines, URLs, tokens, media buffers, recordings, live handles, or provider-private data. On reload, only provider-declared recoverable jobs are restored; the rest become orphaned or provider-unavailable with a safe reason. +Diagnostics live under `slopsmith.jobs.diagnostics.v1`. Browser snapshots include provider summaries, selected/default providers, active/queued/paused/recent-terminal jobs, bounded per-job history, recent outcomes, bridge hits, limits, and notes. Backend snapshots are exported as `system/jobs.json` in diagnostics bundles and include server-side provider/job summaries even when no browser jobs runtime has hydrated. Browser persistence is limited to selected provider choices and provider-declared recoverable job references; backend-integrated queues should restore from their own authoritative server state and re-adopt safe summaries at plugin setup. Neither path stores raw payloads, non-recoverable private state, file paths, command lines, URLs, tokens, media buffers, recordings, live handles, or provider-private data in support surfaces. Compatibility bridge hits such as `jobs.legacy-plugin-queue`, `jobs.legacy-status-screen`, `jobs.legacy-backend-route`, `jobs.legacy-progress-poll`, and `jobs.legacy-update-flow` are diagnostics-only attribution. They never grant approval or execute queued work by themselves. diff --git a/docs/capability-roadmap.md b/docs/capability-roadmap.md index e42f257..f604391 100644 --- a/docs/capability-roadmap.md +++ b/docs/capability-roadmap.md @@ -62,7 +62,7 @@ Deferred follow-up slices: a `contributor` role so plugins ship their own challe ## Jobs Control Plane Slice -The jobs slice promotes `jobs` from a deferred domain to an active privileged provider-coordinator. It owns provider registration, selected-provider persistence, capacity-aware scheduling, explicit user-approved enqueue/retry gates, prompt-free list/inspect, cancellation/pause/resume/retry coordination, reload recovery for provider-declared recoverable references, compatibility bridge hit accounting, and redaction-safe diagnostics under `slopsmith.jobs.diagnostics.v1`. +The jobs slice promotes `jobs` from a deferred domain to an active privileged provider-coordinator. It owns browser provider registration, selected-provider persistence, capacity-aware scheduling, explicit user-approved enqueue/retry gates, prompt-free list/inspect, cancellation/pause/resume/retry coordination, reload recovery for provider-declared recoverable references, compatibility bridge hit accounting, and redaction-safe diagnostics under `slopsmith.jobs.diagnostics.v1`. A backend companion lets plugin route code register/adopt/progress/settle backend-owned jobs into `/api/jobs`, `/ws/jobs`, and diagnostics `system/jobs.json` without exposing provider-private execution payloads. Providers keep actual privileged work private. Core stores only safe job summaries, provider metadata, selected/default provider choices, bounded lifecycle history, terminal outcomes, and recovery references. It does not persist raw payloads, active non-recoverable work, DB schemas, paths, filenames, URLs, tokens, command lines, media/artifacts, recordings, live handles, or provider-private values. diff --git a/lib/diagnostics_bundle.py b/lib/diagnostics_bundle.py index 8d53a62..e4c1292 100644 --- a/lib/diagnostics_bundle.py +++ b/lib/diagnostics_bundle.py @@ -1019,6 +1019,7 @@ def _assemble_files_and_notes( client_contributions: dict | None, log: logging.Logger, plugins_root: "Path | list[Path] | None" = None, + backend_jobs_snapshot: dict | None = None, ) -> tuple[dict[str, bytes], list[str], str, "Redactor | None"]: """Assemble all diagnostic file bytes without packaging into a zip. @@ -1044,6 +1045,11 @@ def _assemble_files_and_notes( files["system/env.json"] = env_payload plugins_data = _system_plugins(loaded_plugins, plugins_root=plugins_root, redactor=redactor) files["system/plugins.json"] = _safe_json_dumps(plugins_data).encode("utf-8") + if backend_jobs_snapshot is not None: + payload = backend_jobs_snapshot + if redactor is not None: + payload = _redact_value(payload, redactor) + files["system/jobs.json"] = _safe_json_dumps(payload).encode("utf-8") # Plugin loading is async and takes a few seconds on cold # boot. If the bundle was captured during that window, every # plugin appears as an "orphan" — flag the likely race so the @@ -1195,6 +1201,7 @@ def build_bundle( client_contributions: dict | None = None, log: logging.Logger, plugins_root: "Path | list[Path] | None" = None, + backend_jobs_snapshot: dict | None = None, ) -> tuple[bytes, str, dict]: """Returns (zip_bytes, filename, manifest_dict).""" files, notes, runtime_kind, redactor = _assemble_files_and_notes( @@ -1212,6 +1219,7 @@ def build_bundle( client_contributions=client_contributions, log=log, plugins_root=plugins_root, + backend_jobs_snapshot=backend_jobs_snapshot, ) manifest = _make_manifest( @@ -1274,6 +1282,7 @@ def preview_bundle( redact: bool, log: logging.Logger, plugins_root: "Path | list[Path] | None" = None, + backend_jobs_snapshot: dict | None = None, ) -> dict: """Lightweight preview: returns manifest-shaped output without building the zip or executing plugin diagnostics callables. @@ -1317,6 +1326,7 @@ def preview_bundle( client_contributions=None, log=log, plugins_root=plugins_root, + backend_jobs_snapshot=backend_jobs_snapshot, ) # ── Callable-output placeholders ────────────────────────────────────── diff --git a/lib/jobs_backend.py b/lib/jobs_backend.py new file mode 100644 index 0000000..a8a5a8d --- /dev/null +++ b/lib/jobs_backend.py @@ -0,0 +1,526 @@ +"""Backend jobs registry for plugin-owned long-running work. + +This is a server-side companion to the browser jobs capability host. It keeps +only redaction-safe public job state; plugin backends remain responsible for +private execution payloads such as local filenames, subprocess handles, and +artifact paths. +""" + +from __future__ import annotations + +import datetime as _dt +import re +import threading +import uuid +from collections import deque +from typing import Any, Callable + +_SCHEMA = "slopsmith.jobs.diagnostics.v1" +_STATES = { + "queued", + "running", + "paused", + "cancellation-requested", + "cancelled", + "completed", + "failed", + "provider-unavailable", + "orphaned", +} +_TERMINAL_STATES = {"cancelled", "completed", "failed", "provider-unavailable", "orphaned"} +_ACTIVE_STATES = {"queued", "running", "paused", "cancellation-requested"} +_FAILURE_CATEGORIES = { + "invalid-input", + "permission-denied", + "provider-unavailable", + "unsupported-operation", + "timeout", + "cancellation", + "external-dependency", + "storage", + "provider-failure", + "unknown", +} +_PRIORITY_ORDER = { + "user-approved-interactive": 0, + "background-maintenance": 1, +} +_SAFE_ID_RE = re.compile(r"[^A-Za-z0-9_.:-]+") +_MAX_LABEL = 120 +_MAX_REASON = 240 +_MAX_STEP = 80 +_MAX_HISTORY_PER_JOB = 50 +_MAX_TERMINAL_JOBS = 50 +_MAX_OUTCOMES = 100 +_MAX_BRIDGE_HITS = 100 + + +def _now() -> str: + return _dt.datetime.now(_dt.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") + + +def _safe_id(value: Any, default: str = "") -> str: + if value is None: + return default + text = str(value).strip() + if not text: + return default + text = _SAFE_ID_RE.sub("-", text)[:160].strip("-_.:") + return text or default + + +def _safe_text(value: Any, default: str = "", limit: int = _MAX_REASON) -> str: + if value is None: + return default + text = str(value).replace("\x00", " ").strip() + if not text: + return default + # Conservative path/URL/secret scrubbing. Callers should pass safe text + # already; this protects support surfaces from accidents. + text = re.sub(r"https?://\S+", "", text) + text = re.sub(r"(?i)(token|secret|api[_-]?key|password|pwd)=([^\s&]+)", r"\1=", text) + text = re.sub(r"/Users/[^\s]+", "", text) + text = re.sub(r"[A-Za-z]:\\[^\s]+", "", text) + if len(text) > limit: + text = text[: max(0, limit - 3)].rstrip() + "..." + return text + + +def _number(value: Any, default: float = 0.0) -> float: + try: + return float(value) + except (TypeError, ValueError): + return default + + +def _progress(source: Any = None) -> dict: + data = source if isinstance(source, dict) else {} + raw_mode = data.get("mode") + mode = raw_mode if raw_mode in {"determinate", "indeterminate", "step-only"} else ("determinate" if data.get("percent") is not None else "indeterminate") + percent = None + if mode == "determinate": + percent = max(0.0, min(100.0, _number(data.get("percent"), 0.0))) + return { + "mode": mode, + "percent": percent, + "step": _safe_text(data.get("step") or data.get("currentStep"), "", _MAX_STEP), + "message": _safe_text(data.get("message") or data.get("safeMessage"), "", _MAX_REASON), + "updatedAt": _now(), + } + + +def _result(outcome: str, payload: dict | None = None, reason: str = "") -> dict: + out = {"outcome": outcome, "status": "applied" if outcome in {"handled", "queued", "completed", "cancelled", "failed", "timeout", "retry-started"} else "rejected"} + if payload: + out["payload"] = payload + if reason: + out["reason"] = _safe_text(reason) + return out + + +class BackendJobs: + """Thread-safe, redaction-safe backend jobs state.""" + + def __init__(self, *, log=None) -> None: + self._log = log + self._lock = threading.RLock() + self._providers: dict[str, dict] = {} + self._jobs: dict[str, dict] = {} + self._terminal_ids: deque[str] = deque(maxlen=_MAX_TERMINAL_JOBS) + self._outcomes: deque[dict] = deque(maxlen=_MAX_OUTCOMES) + self._bridge_hits: deque[dict] = deque(maxlen=_MAX_BRIDGE_HITS) + self._subscribers: set[Callable[[dict], None]] = set() + + def subscribe(self, callback: Callable[[dict], None]) -> Callable[[], None]: + with self._lock: + self._subscribers.add(callback) + + def unsubscribe() -> None: + with self._lock: + self._subscribers.discard(callback) + + return unsubscribe + + def _emit(self, event: str, payload: dict | None = None) -> None: + message = {"type": event, **(payload or {})} + with self._lock: + subscribers = list(self._subscribers) + for callback in subscribers: + try: + callback(message) + except Exception as exc: # noqa: BLE001 + if self._log: + self._log.warning("backend jobs subscriber failed: %s", exc) + + def _remember_outcome(self, operation: str, outcome: str, extra: dict | None = None) -> None: + entry = {"at": _now(), "operation": _safe_id(operation, "operation"), "outcome": _safe_id(outcome, "handled")} + if extra: + for key in ("jobId", "providerId", "requesterId", "category", "safeReason"): + if extra.get(key) is not None: + entry[key] = _safe_text(extra[key]) if key == "safeReason" else _safe_id(extra[key], "") + self._outcomes.append(entry) + + def _history(self, job: dict, kind: str, message: str, extra: dict | None = None) -> None: + entry = {"at": _now(), "kind": _safe_id(kind, "event"), "message": _safe_text(message)} + if extra: + entry.update(extra) + job.setdefault("history", []).append(entry) + if len(job["history"]) > _MAX_HISTORY_PER_JOB: + del job["history"][: len(job["history"]) - _MAX_HISTORY_PER_JOB] + + def _provider_summary(self, provider: dict) -> dict: + return { + "providerId": provider["providerId"], + "pluginId": provider.get("pluginId"), + "label": provider.get("label"), + "jobTypes": list(provider.get("jobTypes") or []), + "actions": list(provider.get("actions") or []), + "capacity": dict(provider.get("capacity") or {}), + "available": bool(provider.get("available", True)), + "recoverySupport": dict(provider.get("recoverySupport") or {}), + "updatedAt": provider.get("updatedAt"), + } + + def _job_summary(self, job: dict, *, include_history: bool = True) -> dict: + summary = { + "jobId": job["jobId"], + "providerId": job["providerId"], + "requesterId": job.get("requesterId"), + "jobType": job.get("jobType"), + "state": job.get("state"), + "priority": job.get("priority"), + "safeLabel": job.get("safeLabel"), + "target": dict(job.get("target") or {}), + "inputs": dict(job.get("inputs") or {}), + "progress": dict(job.get("progress") or {}), + "createdAt": job.get("createdAt"), + "queuedAt": job.get("queuedAt"), + "startedAt": job.get("startedAt"), + "updatedAt": job.get("updatedAt"), + "terminalAt": job.get("terminalAt"), + "terminalOutcome": dict(job.get("terminalOutcome") or {}) if job.get("terminalOutcome") else None, + "retryable": bool(job.get("retryable")), + "externallyManaged": bool(job.get("externallyManaged")), + "attempt": int(job.get("attempt", 1)), + } + if include_history: + summary["history"] = list(job.get("history") or []) + return summary + + def register_provider(self, provider: dict) -> dict: + provider_id = _safe_id(provider.get("providerId") or provider.get("id"), "") + if not provider_id: + return _result("validation-failed", reason="providerId is required") + job_types = [_safe_id(j, "") for j in provider.get("jobTypes", []) if _safe_id(j, "")] + if not job_types: + return _result("validation-failed", reason="jobTypes are required") + actions = [_safe_id(a, "") for a in provider.get("actions", []) if _safe_id(a, "")] + now = _now() + normalized = { + "providerId": provider_id, + "pluginId": _safe_id(provider.get("pluginId") or provider.get("plugin_id"), provider_id), + "label": _safe_text(provider.get("label") or provider_id, provider_id, _MAX_LABEL), + "jobTypes": job_types, + "actions": actions, + "capacity": provider.get("capacity") if isinstance(provider.get("capacity"), dict) else {}, + "available": bool(provider.get("available", True)), + "recoverySupport": provider.get("recoverySupport") if isinstance(provider.get("recoverySupport"), dict) else {}, + "updatedAt": now, + } + with self._lock: + self._providers[provider_id] = normalized + self._remember_outcome("register-provider", "handled", {"providerId": provider_id}) + self._emit("provider-registered", {"provider": self._provider_summary(normalized)}) + return _result("handled", {"provider": self._provider_summary(normalized)}) + + def unregister_provider(self, provider_id: str) -> dict: + provider_id = _safe_id(provider_id, "") + with self._lock: + provider = self._providers.pop(provider_id, None) + if not provider: + return _result("no-owner", reason="provider not found") + self._remember_outcome("unregister-provider", "handled", {"providerId": provider_id}) + self._emit("provider-unregistered", {"providerId": provider_id}) + return _result("handled") + + def adopt(self, *, provider_id: str, job_type: str, job_id: str | None = None, state: str = "running", requester_id: str | None = None, safe_label: str | None = None, target: dict | None = None, inputs: dict | None = None, priority: str = "user-approved-interactive", progress: dict | None = None, category: str | None = None, safe_reason: str | None = None, result_summary: str | None = None, retryable: bool = False, externally_managed: bool = True) -> dict: + provider_id = _safe_id(provider_id, "") + job_type = _safe_id(job_type, "") + if not provider_id or not job_type: + return _result("validation-failed", reason="provider_id and job_type are required") + job_id = _safe_id(job_id, "") or f"job-{uuid.uuid4().hex[:12]}" + desired_state = _safe_id(state, "running") + if desired_state not in _STATES and desired_state not in {"done", "error"}: + desired_state = "running" + if desired_state == "done": + desired_state = "completed" + if desired_state == "error": + desired_state = "failed" + now = _now() + with self._lock: + provider = self._providers.get(provider_id) + if provider and job_type not in provider.get("jobTypes", []): + return _result("no-handler", reason="provider does not handle job type") + job = self._jobs.get(job_id) + if job and job.get("state") in _TERMINAL_STATES: + return _result("stale", {"job": self._job_summary(job)}, "job is terminal") + if not job: + job = { + "jobId": job_id, + "providerId": provider_id, + "requesterId": _safe_id(requester_id, provider_id), + "jobType": job_type, + "state": desired_state, + "priority": _safe_id(priority, "user-approved-interactive"), + "safeLabel": _safe_text(safe_label or job_type, job_type, _MAX_LABEL), + "target": self._safe_ref_dict(target, prefix="target"), + "inputs": self._safe_ref_dict(inputs, prefix="input"), + "progress": _progress(progress), + "createdAt": now, + "queuedAt": now if desired_state == "queued" else None, + "startedAt": now if desired_state == "running" else None, + "updatedAt": now, + "terminalAt": None, + "terminalOutcome": None, + "retryable": False, + "externallyManaged": bool(externally_managed), + "attempt": 1, + "history": [], + } + self._jobs[job_id] = job + self._history(job, "event", "Job adopted from backend provider") + else: + job.update({ + "safeLabel": _safe_text(safe_label or job.get("safeLabel"), job.get("safeLabel", job_type), _MAX_LABEL), + "target": self._safe_ref_dict(target, prefix="target") or job.get("target") or {}, + "inputs": self._safe_ref_dict(inputs, prefix="input") or job.get("inputs") or {}, + "updatedAt": now, + "externallyManaged": bool(externally_managed), + }) + if progress is not None: + job["progress"] = _progress(progress) + if desired_state in _TERMINAL_STATES: + self._settle_locked(job, desired_state, category=category, safe_reason=safe_reason, result_summary=result_summary, retryable=retryable) + event = desired_state + outcome = desired_state + else: + job["state"] = desired_state + job["updatedAt"] = _now() + if desired_state == "running": + job["startedAt"] = job.get("startedAt") or job["updatedAt"] + if desired_state == "queued": + job["queuedAt"] = job.get("queuedAt") or job["updatedAt"] + self._history(job, "event", f"Job adopted as {desired_state}") + self._remember_outcome("adopt", "handled", {"jobId": job_id, "providerId": provider_id, "requesterId": job.get("requesterId")}) + event = "started" if desired_state in {"running", "cancellation-requested"} else desired_state + outcome = "handled" + summary = self._job_summary(job) + self._emit(event, {"job": self._job_summary(job, include_history=False)}) + return _result(outcome, {"job": summary}) + + def _safe_ref_dict(self, data: dict | None, *, prefix: str) -> dict: + if not isinstance(data, dict): + return {} + out: dict[str, str] = {} + for key in ("safeRef", "safeFingerprint", "logicalJobKey", "backendJobId", "bulkJobId"): + if key in data: + out[key] = _safe_id(data.get(key), "") + if not out: + for key in ("id", "ref", "fingerprint"): + if data.get(key): + out[f"{prefix}Ref"] = _safe_id(data.get(key), "") + break + return {k: v for k, v in out.items() if v} + + def update_progress(self, provider_id: str, job_id: str, progress: dict | None = None) -> dict: + provider_id = _safe_id(provider_id, "") + job_id = _safe_id(job_id, "") + with self._lock: + job = self._jobs.get(job_id) + if not job or job.get("providerId") != provider_id: + return _result("no-target", reason="job not found") + if job.get("state") in _TERMINAL_STATES: + self._remember_outcome("progress", "stale", {"jobId": job_id, "providerId": provider_id, "requesterId": job.get("requesterId"), "safeReason": "progress after terminal state"}) + return _result("stale", {"job": self._job_summary(job)}, "job is terminal") + job["progress"] = _progress(progress) + job["updatedAt"] = job["progress"]["updatedAt"] + self._history(job, "progress", job["progress"].get("message") or job["progress"].get("step") or job["progress"].get("mode"), {"progress": dict(job["progress"])}) + self._remember_outcome("progress", "handled", {"jobId": job_id, "providerId": provider_id, "requesterId": job.get("requesterId")}) + summary = self._job_summary(job) + self._emit("progress", {"job": self._job_summary(job, include_history=False)}) + return _result("handled", {"job": summary}) + + def complete(self, provider_id: str, job_id: str, result: dict | None = None) -> dict: + return self._terminal(provider_id, job_id, "completed", result or {}) + + def fail(self, provider_id: str, job_id: str, result: dict | None = None) -> dict: + result = dict(result or {}) + status = "provider-unavailable" if result.get("status") == "provider-unavailable" else "failed" + return self._terminal(provider_id, job_id, status, result) + + def cancelled(self, provider_id: str, job_id: str, result: dict | None = None) -> dict: + return self._terminal(provider_id, job_id, "cancelled", result or {}) + + def mark_provider_unavailable(self, provider_id: str, safe_reason: str = "Provider unavailable") -> dict: + provider_id = _safe_id(provider_id, "") + changed: list[dict] = [] + with self._lock: + provider = self._providers.get(provider_id) + if provider: + provider["available"] = False + provider["updatedAt"] = _now() + for job in self._jobs.values(): + if job.get("providerId") == provider_id and job.get("state") in _ACTIVE_STATES: + self._settle_locked(job, "provider-unavailable", category="provider-unavailable", safe_reason=safe_reason, result_summary="Provider unavailable") + changed.append(self._job_summary(job, include_history=False)) + self._remember_outcome("provider-unavailable", "provider-unavailable", {"providerId": provider_id, "safeReason": safe_reason}) + for job in changed: + self._emit("provider-unavailable", {"job": job}) + return _result("provider-unavailable", {"jobs": changed}) + + def _terminal(self, provider_id: str, job_id: str, state: str, result: dict) -> dict: + provider_id = _safe_id(provider_id, "") + job_id = _safe_id(job_id, "") + with self._lock: + job = self._jobs.get(job_id) + if not job or job.get("providerId") != provider_id: + return _result("no-target", reason="job not found") + if job.get("state") in _TERMINAL_STATES: + return _result("stale", {"job": self._job_summary(job)}, "job is terminal") + self._settle_locked( + job, + state, + category=result.get("category"), + safe_reason=result.get("safeReason") or result.get("safe_reason") or result.get("reason"), + result_summary=result.get("resultSummary") or result.get("summary"), + retryable=bool(result.get("retryable")), + ) + summary = self._job_summary(job) + event = "completed" if state == "completed" else ("cancelled" if state == "cancelled" else ("provider-unavailable" if state == "provider-unavailable" else "failed")) + self._emit(event, {"job": self._job_summary(job, include_history=False)}) + return _result("completed" if state == "completed" else ("cancelled" if state == "cancelled" else "failed"), {"job": summary}) + + def _settle_locked(self, job: dict, state: str, *, category: str | None = None, safe_reason: str | None = None, result_summary: str | None = None, retryable: bool = False) -> None: + if state not in _TERMINAL_STATES: + state = "failed" + terminal_at = _now() + category = _safe_id(category, "") + if state == "completed": + category = "" + elif category not in _FAILURE_CATEGORIES: + category = "unknown" + job["state"] = state + job["retryable"] = bool(retryable) + job["terminalAt"] = terminal_at + job["updatedAt"] = terminal_at + job["terminalOutcome"] = { + "status": state, + "category": category or None, + "retryable": bool(retryable), + "safeReason": _safe_text(safe_reason or ("Completed" if state == "completed" else state.replace("-", " "))), + "resultSummary": _safe_text(result_summary or ("Completed" if state == "completed" else state.replace("-", " "))), + } + if job["jobId"] not in self._terminal_ids: + self._terminal_ids.append(job["jobId"]) + self._history(job, "event", f"Job {state}") + self._remember_outcome(state, state, {"jobId": job["jobId"], "providerId": job["providerId"], "requesterId": job.get("requesterId"), "category": category, "safeReason": job["terminalOutcome"].get("safeReason")}) + + def record_bridge_hit(self, bridge_id: str, detail: dict | None = None) -> dict: + bridge_id = _safe_id(bridge_id, "") + if not bridge_id: + return _result("validation-failed", reason="bridge id required") + data = detail if isinstance(detail, dict) else {} + hit = { + "at": _now(), + "bridgeId": bridge_id, + "providerId": _safe_id(data.get("providerId"), ""), + "jobType": _safe_id(data.get("jobType"), ""), + "operation": _safe_id(data.get("operation"), ""), + "safeReason": _safe_text(data.get("safeReason") or data.get("reason"), ""), + } + hit = {k: v for k, v in hit.items() if v} + with self._lock: + self._bridge_hits.append(hit) + self._remember_outcome("record-bridge-hit", "handled", {"providerId": hit.get("providerId"), "safeReason": hit.get("safeReason")}) + self._emit("bridge-hit", {"bridgeHit": hit}) + return _result("handled", {"bridgeHit": hit}) + + def list_jobs(self, filters: dict | None = None) -> dict: + filters = filters if isinstance(filters, dict) else {} + include_terminal = filters.get("includeTerminal", True) is not False + with self._lock: + jobs = [] + for job in self._jobs.values(): + if filters.get("providerId") and job.get("providerId") != filters.get("providerId"): + continue + if filters.get("jobType") and job.get("jobType") != filters.get("jobType"): + continue + if filters.get("state") and job.get("state") != filters.get("state"): + continue + if not include_terminal and job.get("state") in _TERMINAL_STATES: + continue + jobs.append(self._job_summary(job, include_history=False)) + jobs.sort(key=lambda j: (_PRIORITY_ORDER.get(j.get("priority"), 99), j.get("createdAt") or "")) + providers = [self._provider_summary(p) for p in self._providers.values()] + return {"schema": _SCHEMA, "jobs": jobs, "providers": providers} + + def inspect(self, job_id: str) -> dict: + job_id = _safe_id(job_id, "") + with self._lock: + job = self._jobs.get(job_id) + if not job: + return _result("no-target", reason="job not found") + provider = self._providers.get(job.get("providerId")) + return _result("handled", {"job": self._job_summary(job), "provider": self._provider_summary(provider) if provider else None}) + + def snapshot(self) -> dict: + with self._lock: + providers = [self._provider_summary(p) for p in self._providers.values()] + active = [] + queued = [] + paused = [] + recent_terminal = [] + for job in self._jobs.values(): + summary = self._job_summary(job) + state = job.get("state") + if state == "queued": + queued.append(summary) + elif state == "paused": + paused.append(summary) + elif state in _TERMINAL_STATES: + if job["jobId"] in self._terminal_ids: + recent_terminal.append(summary) + else: + active.append(summary) + recent_terminal.sort(key=lambda j: j.get("terminalAt") or "") + return { + "schema": _SCHEMA, + "backendIntegrated": True, + "providers": providers, + "jobs": { + "active": active, + "queued": queued, + "paused": paused, + "recentTerminal": recent_terminal[-_MAX_TERMINAL_JOBS:], + }, + "bridgeHits": list(self._bridge_hits), + "recentOutcomes": list(self._outcomes), + "limits": { + "maxHistoryPerJob": _MAX_HISTORY_PER_JOB, + "maxTerminalJobs": _MAX_TERMINAL_JOBS, + "maxOutcomes": _MAX_OUTCOMES, + "maxBridgeHits": _MAX_BRIDGE_HITS, + }, + "notes": ["Backend jobs contain redaction-safe summaries only; plugin-private execution payloads remain with provider backends."], + } + + def reset_for_tests(self) -> None: + with self._lock: + self._providers.clear() + self._jobs.clear() + self._terminal_ids.clear() + self._outcomes.clear() + self._bridge_hits.clear() + + +backend_jobs = BackendJobs() diff --git a/server.py b/server.py index 5498359..d6819eb 100644 --- a/server.py +++ b/server.py @@ -45,6 +45,7 @@ import loosefolder as loosefolder_mod # scan workers can import + unpickle _scan_one without re-running this # module's import-time side effects (see lib/scan_worker.py). from scan_worker import _extract_meta_for_file, _relpath, _scan_one +from jobs_backend import backend_jobs import concurrent.futures import contextvars @@ -198,6 +199,9 @@ _DEMO_BLOCKED: list[tuple[str, re.Pattern]] = [ ("POST", re.compile(r"^/api/diagnostics/export$")), ("GET", re.compile(r"^/api/diagnostics/preview$")), ("GET", re.compile(r"^/api/diagnostics/hardware$")), + ("GET", re.compile(r"^/api/jobs(?:/[^/]+)?$")), + ("GET", re.compile(r"^/api/jobs/providers$")), + ("POST", re.compile(r"^/api/jobs/[^/]+/(cancel|retry)$")), # Bundled core plugin — video background upload/delete ("POST", re.compile(r"^/api/plugins/highway_3d/files$")), ("DELETE", re.compile(r"^/api/plugins/highway_3d/files$")), @@ -3357,6 +3361,7 @@ async def startup_events(): "unregister_tuning_provider": unregister_tuning_provider, "get_sloppak_cache_dir": lambda: SLOPPAK_CACHE_DIR, "register_demo_janitor_hook": register_demo_janitor_hook, + "jobs": backend_jobs, # Unified XP service (fee[dB]ack v0.3.0). Plugins that award XP # (minigames, tutorials, …) should feed the single core store via these # instead of keeping a private XP curve. `award_xp` returns the new @@ -6242,6 +6247,7 @@ def export_diagnostics(payload: dict = Body(default_factory=dict)): client_contributions=client_contributions, log=log, plugins_root=_diag_plugins_roots(), + backend_jobs_snapshot=backend_jobs.snapshot(), ) return Response( content=zip_bytes, @@ -6283,6 +6289,7 @@ def preview_diagnostics( redact=redact, log=log, plugins_root=_diag_plugins_roots(), + backend_jobs_snapshot=backend_jobs.snapshot(), ) @@ -6293,6 +6300,88 @@ def diagnostics_hardware(): return _diag_hardware() +@app.get("/api/jobs/providers") +def api_jobs_providers(): + """Return backend-registered job providers with redaction-safe metadata.""" + return {"schema": "slopsmith.jobs.providers.v1", "providers": backend_jobs.list_jobs().get("providers", [])} + + +@app.get("/api/jobs") +def api_jobs_list(providerId: str | None = None, jobType: str | None = None, + state: str | None = None, includeTerminal: bool = True): + """Return backend job summaries without triggering provider work.""" + filters = { + "providerId": providerId, + "jobType": jobType, + "state": state, + "includeTerminal": includeTerminal, + } + return backend_jobs.list_jobs(filters) + + +@app.get("/api/jobs/{job_id}") +def api_jobs_inspect(job_id: str): + result = backend_jobs.inspect(job_id) + if result.get("outcome") == "no-target": + raise HTTPException(status_code=404, detail=result.get("reason") or "job not found") + return result + + +@app.post("/api/jobs/{job_id}/cancel") +def api_jobs_cancel(job_id: str): + result = backend_jobs.inspect(job_id) + if result.get("outcome") == "no-target": + raise HTTPException(status_code=404, detail=result.get("reason") or "job not found") + return { + "outcome": "unsupported-operation", + "status": "rejected", + "reason": "Backend job cancellation must be handled by the owning provider route in this slice", + "payload": result.get("payload", {}), + } + + +@app.post("/api/jobs/{job_id}/retry") +def api_jobs_retry(job_id: str): + result = backend_jobs.inspect(job_id) + if result.get("outcome") == "no-target": + raise HTTPException(status_code=404, detail=result.get("reason") or "job not found") + return { + "outcome": "unsupported-operation", + "status": "rejected", + "reason": "Backend job retry must be handled by the owning provider route in this slice", + "payload": result.get("payload", {}), + } + + +@app.websocket("/ws/jobs") +async def jobs_ws(websocket: WebSocket): + """Stream backend job snapshots and lifecycle updates.""" + await websocket.accept() + loop = asyncio.get_running_loop() + queue: asyncio.Queue = asyncio.Queue(maxsize=200) + + def _on_event(event: dict) -> None: + def _put() -> None: + if queue.full(): + try: + queue.get_nowait() + except asyncio.QueueEmpty: + pass + queue.put_nowait(event) + loop.call_soon_threadsafe(_put) + + unsubscribe = backend_jobs.subscribe(_on_event) + try: + await websocket.send_json({"type": "snapshot", "snapshot": backend_jobs.snapshot()}) + while True: + event = await queue.get() + await websocket.send_json(event) + except WebSocketDisconnect: + pass + finally: + unsubscribe() + + # ── Plugin-provided routes are registered at startup via plugins/__init__.py ─ # (CustomsForge, Ultimate Guitar, etc. are loaded from plugins/ directory) diff --git a/tests/test_jobs_backend.py b/tests/test_jobs_backend.py new file mode 100644 index 0000000..bd9608c --- /dev/null +++ b/tests/test_jobs_backend.py @@ -0,0 +1,59 @@ +from jobs_backend import BackendJobs + + +def test_backend_jobs_adopt_progress_and_complete_redacts_raw_payloads(): + jobs = BackendJobs() + jobs.register_provider({ + "providerId": "provider.test", + "pluginId": "plugin_test", + "label": "Test Provider", + "jobTypes": ["test.convert"], + "actions": ["status"], + }) + + adopted = jobs.adopt( + provider_id="provider.test", + job_type="test.convert", + job_id="backend-1", + state="queued", + requester_id="plugin_test", + safe_label="Convert Example Song", + target={"safeRef": "target-abc", "backendJobId": "backend-1"}, + inputs={"safeFingerprint": "input-abc"}, + ) + assert adopted["outcome"] == "handled" + + progressed = jobs.update_progress( + "provider.test", + "backend-1", + {"mode": "determinate", "percent": 42, "step": "convert", "message": "/Users/example/Secret.psarc token=abc"}, + ) + assert progressed["outcome"] == "handled" + + completed = jobs.complete("provider.test", "backend-1", {"resultSummary": "Wrote /Users/example/out.sloppak"}) + assert completed["outcome"] == "completed" + + snapshot = jobs.snapshot() + text = str(snapshot) + assert "provider.test" in text + assert "target-abc" in text + assert "Secret.psarc" not in text + assert "token=abc" not in text + assert "/Users/example" not in text + assert snapshot["jobs"]["recentTerminal"][0]["state"] == "completed" + + +def test_backend_jobs_provider_unavailable_settles_active_jobs(): + jobs = BackendJobs() + jobs.register_provider({ + "providerId": "provider.test", + "jobTypes": ["test.convert"], + }) + jobs.adopt(provider_id="provider.test", job_type="test.convert", job_id="backend-2", state="running") + + result = jobs.mark_provider_unavailable("provider.test", "backend stopped") + snapshot = jobs.snapshot() + + assert result["outcome"] == "provider-unavailable" + assert snapshot["jobs"]["active"] == [] + assert snapshot["jobs"]["recentTerminal"][0]["state"] == "provider-unavailable"