Fix remaining CodeRabbit jobs review findings

This commit is contained in:
barlind
2026-06-18 00:41:36 -07:00
committed by Bret Mogilefsky
parent 73dd472f52
commit 6c6934c29c
4 changed files with 91 additions and 3 deletions
+23 -3
View File
@@ -131,7 +131,7 @@ def _progress(source: Any = None) -> dict:
def _result(outcome: str, payload: dict | None = None, reason: str = "") -> 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: if payload:
out["payload"] = payload out["payload"] = payload
if reason: if reason:
@@ -364,13 +364,33 @@ class BackendJobs:
def unregister_provider(self, provider_id: str) -> dict: def unregister_provider(self, provider_id: str) -> dict:
provider_id = _safe_id(provider_id, "") provider_id = _safe_id(provider_id, "")
orphaned: list[dict] = []
with self._lock: with self._lock:
provider = self._providers.pop(provider_id, None) provider = self._providers.pop(provider_id, None)
if not provider: if not provider:
return _result("no-owner", reason="provider not found") 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}) 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}) 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: 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, "") provider_id = _safe_id(provider_id, "")
@@ -441,7 +461,7 @@ class BackendJobs:
job["queuedAt"] = job.get("queuedAt") or job["updatedAt"] job["queuedAt"] = job.get("queuedAt") or job["updatedAt"]
self._history(job, "event", f"Job adopted as {desired_state}") 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")}) 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" outcome = "handled"
summary = self._job_summary(job) summary = self._job_summary(job)
self._emit(event, {"job": self._job_summary(job, include_history=False)}) self._emit(event, {"job": self._job_summary(job, include_history=False)})
+5
View File
@@ -6352,6 +6352,11 @@ async def api_jobs_retry(job_id: str, data: dict | None = Body(default=None)):
@app.websocket("/ws/jobs") @app.websocket("/ws/jobs")
async def jobs_ws(websocket: WebSocket): async def jobs_ws(websocket: WebSocket):
"""Stream backend job snapshots and lifecycle updates.""" """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() await websocket.accept()
loop = asyncio.get_running_loop() loop = asyncio.get_running_loop()
queue: asyncio.Queue = asyncio.Queue(maxsize=200) queue: asyncio.Queue = asyncio.Queue(maxsize=200)
+11
View File
@@ -112,6 +112,17 @@ def test_demo_on_read_routes_not_blocked(tmp_path, monkeypatch):
_cleanup(server, client) _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 ─────────────── # ── Demo cookie: set on first GET /, not on subsequent requests ───────────────
def test_demo_cookie_set_on_first_get_root(tmp_path, monkeypatch): def test_demo_cookie_set_on_first_get_root(tmp_path, monkeypatch):
+52
View File
@@ -106,3 +106,55 @@ def test_backend_jobs_retry_requires_user_action_and_retryable_failed_job():
assert denied["outcome"] == "user-action-required" assert denied["outcome"] == "user-action-required"
assert allowed["outcome"] == "retry-started" assert allowed["outcome"] == "retry-started"
assert allowed["payload"] == {"jobId": "backend-4b", "sourceJobId": "backend-4"} 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