From 6c6934c29cd4dc3db2fc742bccbbba486312c422 Mon Sep 17 00:00:00 2001 From: barlind Date: Fri, 12 Jun 2026 13:01:14 +0200 Subject: [PATCH] Fix remaining CodeRabbit jobs review findings --- lib/jobs_backend.py | 26 ++++++++++++++++--- server.py | 5 ++++ tests/test_demo_mode.py | 11 ++++++++ tests/test_jobs_backend.py | 52 ++++++++++++++++++++++++++++++++++++++ 4 files changed, 91 insertions(+), 3 deletions(-) diff --git a/lib/jobs_backend.py b/lib/jobs_backend.py index fc123d3..a433a32 100644 --- a/lib/jobs_backend.py +++ b/lib/jobs_backend.py @@ -131,7 +131,7 @@ def _progress(source: Any = None) -> dict: 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"} + out = {"outcome": outcome, "status": "applied" if outcome in {"handled", "queued", "completed", "cancelled", "failed", "timeout", "retry-started", "provider-unavailable", "orphaned"} else "rejected"} if payload: out["payload"] = payload if reason: @@ -364,13 +364,33 @@ class BackendJobs: def unregister_provider(self, provider_id: str) -> dict: provider_id = _safe_id(provider_id, "") + orphaned: list[dict] = [] with self._lock: provider = self._providers.pop(provider_id, None) if not provider: return _result("no-owner", reason="provider not found") + for job in self._jobs.values(): + if job.get("providerId") == provider_id and job.get("state") in _ACTIVE_STATES: + self._settle_locked( + job, + "orphaned", + category="provider-failure", + safe_reason="Provider unregistered", + result_summary="Provider unregistered", + ) + orphaned.append(self._job_summary(job, include_history=False)) self._remember_outcome("unregister-provider", "handled", {"providerId": provider_id}) + if orphaned: + self._remember_outcome( + "orphaned", + "orphaned", + {"providerId": provider_id, "safeReason": "Provider unregistered"}, + ) + for job in orphaned: + self._emit("orphaned", {"job": job}) self._emit("provider-unregistered", {"providerId": provider_id}) - return _result("handled") + payload = {"orphanedJobs": orphaned} if orphaned else None + return _result("handled", payload) 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, "") @@ -441,7 +461,7 @@ class BackendJobs: 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 + event = "started" if desired_state == "running" else desired_state outcome = "handled" summary = self._job_summary(job) self._emit(event, {"job": self._job_summary(job, include_history=False)}) diff --git a/server.py b/server.py index 34d5b60..855cd07 100644 --- a/server.py +++ b/server.py @@ -6352,6 +6352,11 @@ async def api_jobs_retry(job_id: str, data: dict | None = Body(default=None)): @app.websocket("/ws/jobs") async def jobs_ws(websocket: WebSocket): """Stream backend job snapshots and lifecycle updates.""" + if os.environ.get("SLOPSMITH_DEMO_MODE") == "1": + await websocket.accept() + await websocket.send_json({"error": "demo mode: read-only"}) + await websocket.close(code=1008) + return await websocket.accept() loop = asyncio.get_running_loop() queue: asyncio.Queue = asyncio.Queue(maxsize=200) diff --git a/tests/test_demo_mode.py b/tests/test_demo_mode.py index 686019e..d29a111 100644 --- a/tests/test_demo_mode.py +++ b/tests/test_demo_mode.py @@ -112,6 +112,17 @@ def test_demo_on_read_routes_not_blocked(tmp_path, monkeypatch): _cleanup(server, client) +def test_demo_on_ws_jobs_blocked(tmp_path, monkeypatch): + """Jobs WebSocket should not expose control-plane state in demo mode.""" + server, client = _make_client(tmp_path, monkeypatch, demo=True) + try: + with client.websocket_connect("/ws/jobs") as ws: + msg = ws.receive_json() + assert msg == {"error": "demo mode: read-only"} + finally: + _cleanup(server, client) + + # ── Demo cookie: set on first GET /, not on subsequent requests ─────────────── def test_demo_cookie_set_on_first_get_root(tmp_path, monkeypatch): diff --git a/tests/test_jobs_backend.py b/tests/test_jobs_backend.py index ffa918e..873224e 100644 --- a/tests/test_jobs_backend.py +++ b/tests/test_jobs_backend.py @@ -106,3 +106,55 @@ def test_backend_jobs_retry_requires_user_action_and_retryable_failed_job(): assert denied["outcome"] == "user-action-required" assert allowed["outcome"] == "retry-started" assert allowed["payload"] == {"jobId": "backend-4b", "sourceJobId": "backend-4"} + + +def test_backend_jobs_provider_unavailable_and_orphaned_are_applied_outcomes(): + jobs = BackendJobs() + jobs.register_provider({ + "providerId": "provider.test", + "jobTypes": ["test.convert"], + }) + jobs.adopt(provider_id="provider.test", job_type="test.convert", job_id="backend-5", state="running") + unavailable = jobs.mark_provider_unavailable("provider.test", "offline") + + assert unavailable["outcome"] == "provider-unavailable" + assert unavailable["status"] == "applied" + + +def test_backend_jobs_unregister_provider_orphans_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-6", state="queued") + + removed = jobs.unregister_provider("provider.test") + inspect = jobs.inspect("backend-6") + + assert removed["outcome"] == "handled" + assert removed["status"] == "applied" + assert removed["payload"]["orphanedJobs"][0]["jobId"] == "backend-6" + assert inspect["payload"]["job"]["state"] == "orphaned" + + +def test_backend_jobs_adopt_emits_cancellation_requested_event(): + jobs = BackendJobs() + jobs.register_provider({ + "providerId": "provider.test", + "jobTypes": ["test.convert"], + }) + seen = [] + unsubscribe = jobs.subscribe(lambda event: seen.append(event.get("type"))) + try: + jobs.adopt( + provider_id="provider.test", + job_type="test.convert", + job_id="backend-7", + state="cancellation-requested", + ) + finally: + unsubscribe() + + assert "cancellation-requested" in seen + assert "started" not in seen