diff --git a/docs/capability-domains.md b/docs/capability-domains.md index c886102..e35720a 100644 --- a/docs/capability-domains.md +++ b/docs/capability-domains.md @@ -157,7 +157,9 @@ Diagnostics live under `slopsmith.note_detection_capability.v1` and contain prov 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 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. +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. Backend provider code can call `register_provider`, `adopt`, `update_progress`, `complete`, `fail`, `cancelled`, `mark_provider_unavailable`, and `record_bridge_hit`; providers may also register private action callbacks for advertised operations such as `job.cancel` and `job.retry`. `/api/jobs/{id}/cancel` and `/api/jobs/{id}/retry` inspect the redaction-safe job owner, enforce the user-action gate for retry, and dispatch to the provider callback without exposing raw payloads to core state, diagnostics, or API responses. + +This slice supports dispatch before execution ownership. Core routes authorized backend job actions to provider callbacks (`cancel` for queued jobs, `retry` for failed jobs, and later provider-level actions where the contract is explicit) while providers keep their private queue stores and worker code. Moving scheduling or execution itself into core is a later, separate upgrade for shared worker patterns across multiple providers; it must not erase domain-specific semantics such as queue-level pause/resume, non-interruptible running jobs, remote service retries, or plugin-owned artifact indexing. 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`. diff --git a/docs/capability-roadmap.md b/docs/capability-roadmap.md index f604391..9cb468b 100644 --- a/docs/capability-roadmap.md +++ b/docs/capability-roadmap.md @@ -66,6 +66,8 @@ The jobs slice promotes `jobs` from a deferred domain to an active privileged pr 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. +Backend action dispatch is the middle ground before core-owned execution: provider route code declares supported actions such as queued-job cancel or failed-job retry; core checks the caller's authorization and job ownership; then it dispatches to provider-owned callbacks while the provider keeps raw payloads, queue semantics, and domain-specific execution. Future jobs upgrades should only promote shared backend scheduling/execution into core if multiple first-party plugins need the same worker behavior for provider-declared recoverable jobs. That larger step must define queue-vs-job actions, restart recovery, cancellation truthfulness, retry attempts, payload privacy, provider failure handling, diagnostics, and migration from plugin-owned queue stores before replacing plugin execution routes. + Jobs bridge removal gates are: bundled and first-party long-running workflows use native `jobs` provider registration/dispatch; normal conversion/import/update/cache smoke runs show no unexpected legacy bridge hits; diagnostics distinguish queued, denied, user-action-required, provider-selection-required, stale, cancelled, completed, failed, timeout, retry-started, orphaned, and provider-unavailable cases; repeated plugin hydration does not duplicate providers or jobs; and reload recovery restores only provider-declared safe references. ## Recommended Next Slices diff --git a/lib/jobs_backend.py b/lib/jobs_backend.py index a8a5a8d..fc123d3 100644 --- a/lib/jobs_backend.py +++ b/lib/jobs_backend.py @@ -9,6 +9,7 @@ artifact paths. from __future__ import annotations import datetime as _dt +import inspect as _inspect import re import threading import uuid @@ -53,6 +54,26 @@ _MAX_HISTORY_PER_JOB = 50 _MAX_TERMINAL_JOBS = 50 _MAX_OUTCOMES = 100 _MAX_BRIDGE_HITS = 100 +_ACTION_ALIASES = { + "cancel": "job.cancel", + "retry": "job.retry", + "pause": "job.pause", + "resume": "job.resume", + "recover": "job.recover", + "status": "job.status", +} +_SAFE_ACTION_PAYLOAD_KEYS = { + "action", + "backendJobId", + "bulkJobId", + "jobId", + "logicalJobKey", + "newJobId", + "providerId", + "safeRef", + "sourceJobId", + "state", +} def _now() -> str: @@ -118,6 +139,32 @@ def _result(outcome: str, payload: dict | None = None, reason: str = "") -> dict return out +def _canonical_action(action: Any) -> str: + raw = _safe_id(action, "") + return _ACTION_ALIASES.get(raw, raw) + + +def _action_names(action: Any) -> tuple[str, ...]: + canonical = _canonical_action(action) + names = [canonical] + for alias, target in _ACTION_ALIASES.items(): + if target == canonical: + names.append(alias) + return tuple(dict.fromkeys(n for n in names if n)) + + +def _safe_action_payload(payload: Any) -> dict: + if not isinstance(payload, dict): + return {} + out = {} + for key, value in payload.items(): + safe_key = _safe_id(key, "") + if safe_key not in _SAFE_ACTION_PAYLOAD_KEYS: + continue + out[safe_key] = _safe_id(value, "") + return {k: v for k, v in out.items() if v} + + class BackendJobs: """Thread-safe, redaction-safe backend jobs state.""" @@ -214,7 +261,16 @@ class BackendJobs: 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, "")] + actions = [_canonical_action(a) for a in provider.get("actions", []) if _canonical_action(a)] + callbacks = provider.get("callbacks") or provider.get("actionHandlers") or provider.get("operationHandlers") + action_handlers = {} + if isinstance(callbacks, dict): + for action, callback in callbacks.items(): + canonical = _canonical_action(action) + if canonical and callable(callback): + action_handlers[canonical] = callback + if canonical not in actions: + actions.append(canonical) now = _now() normalized = { "providerId": provider_id, @@ -222,6 +278,7 @@ class BackendJobs: "label": _safe_text(provider.get("label") or provider_id, provider_id, _MAX_LABEL), "jobTypes": job_types, "actions": actions, + "actionHandlers": action_handlers, "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 {}, @@ -233,6 +290,78 @@ class BackendJobs: self._emit("provider-registered", {"provider": self._provider_summary(normalized)}) return _result("handled", {"provider": self._provider_summary(normalized)}) + async def dispatch_action(self, job_id: str, action: str, request: dict | None = None) -> dict: + canonical = _canonical_action(action) + if not canonical: + return _result("validation-failed", reason="action is required") + request = request if isinstance(request, dict) else {} + authorization = _safe_id(request.get("authorization"), "") + requester_id = _safe_id(request.get("requesterId") or request.get("requester_id"), "api.jobs") + with self._lock: + job = self._jobs.get(_safe_id(job_id, "")) + if not job: + self._remember_outcome(canonical, "no-target", {"requesterId": requester_id}) + return _result("no-target", reason="job not found") + provider = self._providers.get(job.get("providerId")) + if not provider: + self._remember_outcome(canonical, "no-owner", {"jobId": job.get("jobId"), "providerId": job.get("providerId"), "requesterId": requester_id}) + return _result("no-owner", {"job": self._job_summary(job)}, "provider not found") + if not provider.get("available", True): + self._remember_outcome(canonical, "unavailable", {"jobId": job.get("jobId"), "providerId": provider.get("providerId"), "requesterId": requester_id}) + return _result("unavailable", {"job": self._job_summary(job), "provider": self._provider_summary(provider)}, "provider unavailable") + provider_actions = set(provider.get("actions") or []) + if canonical not in provider_actions: + self._remember_outcome(canonical, "unsupported-operation", {"jobId": job.get("jobId"), "providerId": provider.get("providerId"), "requesterId": requester_id}) + return _result("unsupported-operation", {"job": self._job_summary(job), "provider": self._provider_summary(provider)}, "provider does not advertise this action") + handlers = provider.get("actionHandlers") or {} + handler = None + for name in _action_names(canonical): + if callable(handlers.get(name)): + handler = handlers[name] + break + if handler is None: + self._remember_outcome(canonical, "no-handler", {"jobId": job.get("jobId"), "providerId": provider.get("providerId"), "requesterId": requester_id}) + return _result("no-handler", {"job": self._job_summary(job), "provider": self._provider_summary(provider)}, "provider action handler not registered") + state = job.get("state") + if state in _TERMINAL_STATES and canonical != "job.retry": + self._remember_outcome(canonical, "stale", {"jobId": job.get("jobId"), "providerId": provider.get("providerId"), "requesterId": requester_id, "safeReason": "job is terminal"}) + return _result("stale", {"job": self._job_summary(job)}, "job is terminal") + if canonical == "job.retry" and (state != "failed" or not job.get("retryable")): + self._remember_outcome(canonical, "stale", {"jobId": job.get("jobId"), "providerId": provider.get("providerId"), "requesterId": requester_id, "safeReason": "job is not retryable"}) + return _result("stale", {"job": self._job_summary(job)}, "job is not retryable") + if canonical in {"job.enqueue", "job.retry"} and authorization != "user-action": + self._remember_outcome(canonical, "user-action-required", {"jobId": job.get("jobId"), "providerId": provider.get("providerId"), "requesterId": requester_id}) + return _result("user-action-required", {"job": self._job_summary(job)}, "user action authorization required") + call_payload = { + "action": canonical, + "authorization": authorization, + "requesterId": requester_id, + "job": self._job_summary(job), + "provider": self._provider_summary(provider), + } + try: + provider_result = handler(call_payload) + if _inspect.isawaitable(provider_result): + provider_result = await provider_result + except Exception as exc: # noqa: BLE001 + safe_reason = _safe_text(str(exc), "provider action failed") + with self._lock: + self._remember_outcome(canonical, "failed", {"jobId": job_id, "providerId": provider.get("providerId"), "requesterId": requester_id, "safeReason": safe_reason}) + return _result("failed", reason=safe_reason) + + if not isinstance(provider_result, dict): + provider_result = {"outcome": "handled"} + outcome = _safe_id(provider_result.get("outcome"), "handled") + status = provider_result.get("status") + reason = _safe_text(provider_result.get("reason"), "") + payload = _safe_action_payload(provider_result.get("payload")) + with self._lock: + self._remember_outcome(canonical, outcome, {"jobId": job_id, "providerId": provider.get("providerId"), "requesterId": requester_id, "safeReason": reason}) + result = _result(outcome, payload or None, reason) + if status in {"applied", "rejected"}: + result["status"] = status + return result + def unregister_provider(self, provider_id: str) -> dict: provider_id = _safe_id(provider_id, "") with self._lock: diff --git a/server.py b/server.py index d6819eb..34d5b60 100644 --- a/server.py +++ b/server.py @@ -6328,29 +6328,25 @@ def api_jobs_inspect(job_id: str): @app.post("/api/jobs/{job_id}/cancel") -def api_jobs_cancel(job_id: str): +async def api_jobs_cancel(job_id: str, data: dict | None = Body(default=None)): 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", {}), - } + request = data if isinstance(data, dict) else {} + request.setdefault("authorization", "user-action") + request.setdefault("requesterId", "api.jobs") + return await backend_jobs.dispatch_action(job_id, "job.cancel", request) @app.post("/api/jobs/{job_id}/retry") -def api_jobs_retry(job_id: str): +async def api_jobs_retry(job_id: str, data: dict | None = Body(default=None)): 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", {}), - } + request = data if isinstance(data, dict) else {} + request.setdefault("authorization", "user-action") + request.setdefault("requesterId", "api.jobs") + return await backend_jobs.dispatch_action(job_id, "job.retry", request) @app.websocket("/ws/jobs") diff --git a/tests/test_jobs_backend.py b/tests/test_jobs_backend.py index bd9608c..ffa918e 100644 --- a/tests/test_jobs_backend.py +++ b/tests/test_jobs_backend.py @@ -1,4 +1,5 @@ from jobs_backend import BackendJobs +import asyncio def test_backend_jobs_adopt_progress_and_complete_redacts_raw_payloads(): @@ -57,3 +58,51 @@ def test_backend_jobs_provider_unavailable_settles_active_jobs(): assert result["outcome"] == "provider-unavailable" assert snapshot["jobs"]["active"] == [] assert snapshot["jobs"]["recentTerminal"][0]["state"] == "provider-unavailable" + + +def test_backend_jobs_dispatches_private_provider_action_and_redacts_payload(): + calls = [] + jobs = BackendJobs() + + def cancel_handler(request): + calls.append(request) + return { + "outcome": "cancelled", + "payload": { + "jobId": request["job"]["jobId"], + "rawPath": "/Users/example/Secret.psarc", + }, + } + + jobs.register_provider({ + "providerId": "provider.test", + "jobTypes": ["test.convert"], + "actions": ["job.status"], + "callbacks": {"job.cancel": cancel_handler}, + }) + jobs.adopt(provider_id="provider.test", job_type="test.convert", job_id="backend-3", state="queued") + + result = asyncio.run(jobs.dispatch_action("backend-3", "cancel", {"requesterId": "test"})) + + assert result["outcome"] == "cancelled" + assert result["payload"] == {"jobId": "backend-3"} + assert calls[0]["job"]["jobId"] == "backend-3" + assert "rawPath" not in str(result) + + +def test_backend_jobs_retry_requires_user_action_and_retryable_failed_job(): + jobs = BackendJobs() + jobs.register_provider({ + "providerId": "provider.test", + "jobTypes": ["test.convert"], + "callbacks": {"job.retry": lambda request: {"outcome": "retry-started", "payload": {"jobId": "backend-4b", "sourceJobId": request["job"]["jobId"]}}}, + }) + jobs.adopt(provider_id="provider.test", job_type="test.convert", job_id="backend-4", state="queued") + jobs.fail("provider.test", "backend-4", {"retryable": True, "safeReason": "failed"}) + + denied = asyncio.run(jobs.dispatch_action("backend-4", "retry", {"requesterId": "test"})) + allowed = asyncio.run(jobs.dispatch_action("backend-4", "retry", {"requesterId": "test", "authorization": "user-action"})) + + assert denied["outcome"] == "user-action-required" + assert allowed["outcome"] == "retry-started" + assert allowed["payload"] == {"jobId": "backend-4b", "sourceJobId": "backend-4"}