diff --git a/panel/src/app/(dashboard)/a2a/__tests__/page.test.tsx b/panel/src/app/(dashboard)/a2a/__tests__/page.test.tsx
index be1f8732..473b09aa 100644
--- a/panel/src/app/(dashboard)/a2a/__tests__/page.test.tsx
+++ b/panel/src/app/(dashboard)/a2a/__tests__/page.test.tsx
@@ -147,7 +147,7 @@ describe("A2APage", () => {
expect(screen.getByPlaceholderText(/chime in/i)).toBeInTheDocument();
expect(
screen.getByText(
- /direct A2A message from you to the selected participant/i,
+ /posts into this conversation — visible to both participants/i,
),
).toBeInTheDocument();
expect(screen.getByText("Live")).toBeInTheDocument();
diff --git a/panel/src/app/(dashboard)/agents/[agentId]/__tests__/page.test.tsx b/panel/src/app/(dashboard)/agents/[agentId]/__tests__/page.test.tsx
new file mode 100644
index 00000000..090226e0
--- /dev/null
+++ b/panel/src/app/(dashboard)/agents/[agentId]/__tests__/page.test.tsx
@@ -0,0 +1,90 @@
+import { describe, it, expect, vi } from "vitest";
+import { render, screen } from "@testing-library/react";
+
+// Detail-page parity fix: the grid card's spawn affordance (SpawnAgentDialog,
+// which collects task id + message) was already correct, but this page called
+// spawnAgent.mutateAsync({ agentId }) directly from a bare button — no task,
+// no message, no double-fire guard. Both bare buttons must now render the
+// shared SpawnAgentDialog instead.
+
+vi.mock("next/navigation", () => ({
+ useParams: () => ({ agentId: "fe-dev-2" }),
+ useRouter: () => ({ back: vi.fn() }),
+}));
+
+vi.mock("@/hooks/use-agents", () => ({
+ useAgentStatus: vi.fn(),
+ useAgentDefinition: vi.fn(() => ({ data: undefined })),
+ useStopAgent: vi.fn(() => ({ mutateAsync: vi.fn() })),
+}));
+
+vi.mock("@/components/agents", () => ({
+ AgentStatusCards: () => null,
+ ResolveWaitDialog: () => null,
+ AgentStreamViewer: () => null,
+ SpawnAgentDialog: ({
+ agentId,
+ agentName,
+ trigger,
+ }: {
+ agentId: string;
+ agentName: string;
+ trigger: React.ReactNode;
+ }) => (
+
+ {trigger}
+
+ ),
+}));
+
+import { useAgentStatus } from "@/hooks/use-agents";
+import AgentDetailPage from "../page";
+
+describe("AgentDetailPage — spawn dialog parity", () => {
+ it("renders SpawnAgentDialog (not a bare button) in the error state", () => {
+ vi.mocked(useAgentStatus).mockReturnValue({
+ data: undefined,
+ isLoading: false,
+ error: new Error("not found"),
+ refetch: vi.fn(),
+ } as unknown as ReturnType);
+
+ render();
+
+ const dialog = screen.getByTestId("spawn-agent-dialog");
+ expect(dialog).toHaveAttribute("data-agent-id", "fe-dev-2");
+ expect(
+ screen.getByRole("button", { name: /Spawn Agent/i }),
+ ).toBeInTheDocument();
+ });
+
+ it("renders SpawnAgentDialog in the header when the agent is not active", () => {
+ vi.mocked(useAgentStatus).mockReturnValue({
+ data: {
+ agent_id: "fe-dev-2",
+ state: "stopped",
+ task_id: null,
+ error_count: 0,
+ started_at: null,
+ waiting_for: null,
+ },
+ isLoading: false,
+ error: undefined,
+ refetch: vi.fn(),
+ } as unknown as ReturnType);
+
+ render();
+
+ const dialog = screen.getByTestId("spawn-agent-dialog");
+ expect(dialog).toHaveAttribute("data-agent-id", "fe-dev-2");
+ expect(screen.getByRole("button", { name: "Spawn" })).toBeInTheDocument();
+ // Active-state Stop buttons must not render alongside a down agent.
+ expect(
+ screen.queryByRole("button", { name: "Stop" }),
+ ).not.toBeInTheDocument();
+ });
+});
diff --git a/panel/src/app/(dashboard)/agents/[agentId]/page.tsx b/panel/src/app/(dashboard)/agents/[agentId]/page.tsx
index e4257025..06a1c921 100644
--- a/panel/src/app/(dashboard)/agents/[agentId]/page.tsx
+++ b/panel/src/app/(dashboard)/agents/[agentId]/page.tsx
@@ -5,7 +5,6 @@ import { formatDistanceToNow } from "date-fns";
import {
useAgentStatus,
useStopAgent,
- useSpawnAgent,
useAgentDefinition,
} from "@/hooks/use-agents";
import { Button } from "@/components/ui/button";
@@ -33,6 +32,7 @@ import {
AgentStatusCards,
ResolveWaitDialog,
AgentStreamViewer,
+ SpawnAgentDialog,
} from "@/components/agents";
// Role display labels
@@ -66,7 +66,6 @@ export default function AgentDetailPage() {
const { data: agent, isLoading, error, refetch } = useAgentStatus(agentId);
const { data: definition } = useAgentDefinition(agentId);
const stopAgent = useStopAgent();
- const spawnAgent = useSpawnAgent();
// Get display values from definition or fallback
const displayName = definition?.name || agentId;
@@ -88,15 +87,6 @@ export default function AgentDetailPage() {
}
};
- const handleSpawn = async () => {
- try {
- await spawnAgent.mutateAsync({ agentId });
- toast.success("Agent spawned successfully");
- } catch {
- toast.error("Failed to spawn agent");
- }
- };
-
if (error) {
return (
@@ -113,10 +103,16 @@ export default function AgentDetailPage() {
The agent may not be running or the ID is invalid.
-
+
+
+ Spawn Agent
+
+ }
+ />
@@ -178,10 +174,16 @@ export default function AgentDetailPage() {
>
) : (
-
+
+
+ Spawn
+
+ }
+ />
)}
diff --git a/panel/src/components/a2a/__tests__/a2a-reply-composer.test.tsx b/panel/src/components/a2a/__tests__/a2a-reply-composer.test.tsx
index 2a232a5f..0366488f 100644
--- a/panel/src/components/a2a/__tests__/a2a-reply-composer.test.tsx
+++ b/panel/src/components/a2a/__tests__/a2a-reply-composer.test.tsx
@@ -119,14 +119,50 @@ describe("A2AReplyComposer", () => {
);
});
- it("states the pairwise seam honestly in the helper text", () => {
+ it("states the interjection semantics honestly in the helper text", () => {
renderComposer();
- // Guard the honesty note: the reply is a DIRECT CEO->participant message,
- // not an injection into the watched transcript.
+ // Guard the honesty note: the message posts into THIS conversation,
+ // visible to both participants — not a re-homed CEO<->target DM.
expect(
- screen.getByText(
- /direct A2A message from you to the selected participant/i,
- ),
+ screen.getByText(/posts into this conversation/i),
+ ).toBeInTheDocument();
+ expect(
+ screen.getByText(/visible to both participants/i),
).toBeInTheDocument();
});
});
+
+describe("A2AReplyComposer when one participant is the CEO", () => {
+ beforeEach(() => {
+ mutate.mockReset();
+ });
+
+ function renderCeoComposer(lastSender: string | null = "be-dev-1") {
+ return render(
+ ,
+ );
+ }
+
+ it("never offers the CEO itself as a recipient", () => {
+ const { container } = renderCeoComposer();
+ expect(container.querySelector('[data-value="ceo"]')).toBeNull();
+ expect(container.querySelector('[data-value="be-dev-1"]')).not.toBeNull();
+ });
+
+ it("always sends to the other participant, even when they spoke last as agent_a", () => {
+ renderCeoComposer("ceo");
+ fireEvent.change(screen.getByPlaceholderText(/chime in/i), {
+ target: { value: "Following up" },
+ });
+ fireEvent.click(screen.getByRole("button", { name: /send/i }));
+ expect(mutate).toHaveBeenCalledWith(
+ expect.objectContaining({ to_agent: "be-dev-1" }),
+ expect.anything(),
+ );
+ });
+});
diff --git a/panel/src/components/a2a/__tests__/a2a-utils.test.ts b/panel/src/components/a2a/__tests__/a2a-utils.test.ts
index 88f9b00f..51a926df 100644
--- a/panel/src/components/a2a/__tests__/a2a-utils.test.ts
+++ b/panel/src/components/a2a/__tests__/a2a-utils.test.ts
@@ -1,5 +1,9 @@
import { describe, it, expect } from "vitest";
-import { lastSenderOf, pickDefaultRecipient } from "../a2a-utils";
+import {
+ lastSenderOf,
+ pickDefaultRecipient,
+ recipientOptions,
+} from "../a2a-utils";
describe("lastSenderOf", () => {
it("returns null for an empty transcript", () => {
@@ -31,3 +35,20 @@ describe("pickDefaultRecipient", () => {
);
});
});
+
+describe("recipientOptions", () => {
+ it("returns both participants when neither is the CEO", () => {
+ expect(recipientOptions("be-dev-1", "be-qa")).toEqual([
+ "be-dev-1",
+ "be-qa",
+ ]);
+ });
+
+ it("excludes the CEO when it's agent_a", () => {
+ expect(recipientOptions("ceo", "be-dev-1")).toEqual(["be-dev-1"]);
+ });
+
+ it("excludes the CEO when it's agent_b", () => {
+ expect(recipientOptions("be-dev-1", "ceo")).toEqual(["be-dev-1"]);
+ });
+});
diff --git a/panel/src/components/a2a/a2a-reply-composer.tsx b/panel/src/components/a2a/a2a-reply-composer.tsx
index 016b16f2..cdf38d8d 100644
--- a/panel/src/components/a2a/a2a-reply-composer.tsx
+++ b/panel/src/components/a2a/a2a-reply-composer.tsx
@@ -15,7 +15,7 @@ import { toast } from "sonner";
import { getAgentDisplayName } from "@/lib/agent-utils";
import { getErrorMessage } from "@/lib/api/client";
import { useReplyAsCeo } from "@/hooks/use-a2a-live";
-import { pickDefaultRecipient } from "./a2a-utils";
+import { pickDefaultRecipient, recipientOptions } from "./a2a-utils";
interface A2AReplyComposerProps {
conversationId: string;
@@ -38,8 +38,14 @@ export function A2AReplyComposer({
const [chosenRecipient, setChosenRecipient] = useState(null);
const reply = useReplyAsCeo();
+ // Excludes "ceo" from the options: in the CEO's own conversation with an
+ // agent, one of {agentA, agentB} is "ceo" itself, and it must never be a
+ // selectable/default reply target.
+ const options = recipientOptions(agentA, agentB);
const recipient =
- chosenRecipient ?? pickDefaultRecipient(agentA, agentB, lastSender);
+ options.length === 1
+ ? options[0]
+ : (chosenRecipient ?? pickDefaultRecipient(agentA, agentB, lastSender));
const handleSubmit = (e: React.FormEvent) => {
e.preventDefault();
@@ -86,7 +92,7 @@ export function A2AReplyComposer({
- {[agentA, agentB].map((slug) => (
+ {options.map((slug) => (
{getAgentDisplayName(slug)}
@@ -104,9 +110,8 @@ export function A2AReplyComposer({
- Sends a direct A2A message from you to the selected participant. It
- lands in your own conversation with that agent, not inside this
- transcript.
+ Posts into this conversation — visible to both participants, addressed
+ to whoever you pick above.
);
diff --git a/panel/src/components/a2a/a2a-utils.ts b/panel/src/components/a2a/a2a-utils.ts
index f6977af8..d4808807 100644
--- a/panel/src/components/a2a/a2a-utils.ts
+++ b/panel/src/components/a2a/a2a-utils.ts
@@ -4,6 +4,10 @@
import type { A2AChatMessage } from "@/lib/api/a2a";
+/** The human CEO's fixed slug — never a valid reply target (the CEO composes
+ * as itself, so it can't be its own recipient). */
+export const CEO_SLUG = "ceo";
+
/**
* Slug of the sender of the chronologically latest message, or null when the
* transcript is empty. Sorts defensively — the API contract is oldest-first,
@@ -20,6 +24,16 @@ export function lastSenderOf(
return sorted[sorted.length - 1].from_agent;
}
+/**
+ * Valid reply recipients for a conversation: both participants, minus the
+ * CEO itself when it's a party. A CEO<->agent conversation has exactly one
+ * possible target (the agent) — no picker ambiguity, and critically no way
+ * to select "ceo" and have the CEO reply to itself.
+ */
+export function recipientOptions(agentA: string, agentB: string): string[] {
+ return [agentA, agentB].filter((slug) => slug !== CEO_SLUG);
+}
+
/**
* Default reply recipient: the participant who spoke last (the natural
* "answer them" target), falling back to agent_a when the transcript is empty
diff --git a/panel/src/components/agents/__tests__/spawn-agent-dialog.test.tsx b/panel/src/components/agents/__tests__/spawn-agent-dialog.test.tsx
new file mode 100644
index 00000000..84474691
--- /dev/null
+++ b/panel/src/components/agents/__tests__/spawn-agent-dialog.test.tsx
@@ -0,0 +1,131 @@
+import { describe, it, expect, vi, beforeEach, afterEach } from "vitest";
+import { fireEvent, render, screen, waitFor } from "@testing-library/react";
+
+// Production-triage bug: a manual spawn POSTed TWICE 2.5ms apart, both
+// rejected "Agent already running" with no visible reason. Covers the fix:
+// a synchronous re-entrancy guard against the double-fire, an
+// already_running-aware toast, and the real backend refusal message
+// reaching the CEO instead of a generic "Failed to spawn agent".
+
+const { mutateAsync, toastSuccess, toastError, toastInfo } = vi.hoisted(() => ({
+ mutateAsync: vi.fn(),
+ toastSuccess: vi.fn(),
+ toastError: vi.fn(),
+ toastInfo: vi.fn(),
+}));
+
+vi.mock("@/hooks/use-agents", () => ({
+ useSpawnAgent: () => ({ mutateAsync, isPending: false }),
+}));
+
+vi.mock("sonner", () => ({
+ toast: { success: toastSuccess, error: toastError, info: toastInfo },
+}));
+
+// client.ts registers axios interceptors that pull in the rate-limit store
+// at import time; stub it so importing the real getErrorMessage is side-effect
+// free (mirrors lib/__tests__/client.test.ts).
+vi.mock("@/store/rate-limit-store", () => ({
+ useRateLimitStore: { getState: vi.fn(() => ({ hitRateLimit: vi.fn() })) },
+}));
+
+import { SpawnAgentDialog } from "../spawn-agent-dialog";
+
+function openDialog() {
+ render(
+ Open Spawn}
+ />,
+ );
+ fireEvent.click(screen.getByRole("button", { name: "Open Spawn" }));
+}
+
+function submitButton() {
+ return screen.getByRole("button", { name: /Spawn Agent/i });
+}
+
+describe("SpawnAgentDialog", () => {
+ beforeEach(() => {
+ mutateAsync.mockReset();
+ toastSuccess.mockReset();
+ toastError.mockReset();
+ toastInfo.mockReset();
+ });
+ afterEach(() => vi.clearAllMocks());
+
+ it("submits task id and initial prompt from the form", async () => {
+ mutateAsync.mockResolvedValue({ already_running: false });
+ openDialog();
+
+ fireEvent.change(screen.getByLabelText(/Task ID/i), {
+ target: { value: "task-123" },
+ });
+ fireEvent.change(screen.getByLabelText(/Initial Prompt/i), {
+ target: { value: "go fix it" },
+ });
+ fireEvent.click(submitButton());
+
+ await waitFor(() => expect(mutateAsync).toHaveBeenCalledTimes(1));
+ expect(mutateAsync).toHaveBeenCalledWith({
+ agentId: "fe-dev-2",
+ request: { task_id: "task-123", initial_prompt: "go fix it" },
+ });
+ expect(toastSuccess).toHaveBeenCalled();
+ });
+
+ it("shows a distinct toast when the spawn was skipped as already-running", async () => {
+ mutateAsync.mockResolvedValue({ already_running: true });
+ openDialog();
+ fireEvent.click(submitButton());
+
+ await waitFor(() => expect(toastInfo).toHaveBeenCalledTimes(1));
+ expect(toastInfo.mock.calls[0][0]).toMatch(/already running/i);
+ expect(toastSuccess).not.toHaveBeenCalled();
+ });
+
+ it("surfaces the backend's actual refusal reason, not a generic message", async () => {
+ // Shape returned by axios on the 409 AgentReadinessError mapping.
+ const axiosLikeError = {
+ isAxiosError: true,
+ message: "Request failed with status code 409",
+ response: {
+ status: 409,
+ data: {
+ detail:
+ "spawn refused for fe-dev-2 (task=t1): state=awaiting_qa " +
+ "requires role in {'qa'} but agent fe-dev-2 is 'developer'",
+ },
+ },
+ };
+ mutateAsync.mockRejectedValue(axiosLikeError);
+ openDialog();
+ fireEvent.click(submitButton());
+
+ await waitFor(() => expect(toastError).toHaveBeenCalledTimes(1));
+ expect(toastError.mock.calls[0][0]).toContain("state=awaiting_qa");
+ expect(toastError.mock.calls[0][0]).not.toBe("Failed to spawn agent");
+ });
+
+ it("blocks a second submit fired before the first mutation settles", async () => {
+ let resolveSpawn: (v: { already_running: boolean }) => void = () => {};
+ mutateAsync.mockImplementation(
+ () =>
+ new Promise((resolve) => {
+ resolveSpawn = resolve;
+ }),
+ );
+ openDialog();
+
+ // Two synchronous clicks, mirroring the 2.5ms-apart double-fire from the
+ // production report — the second must never reach mutateAsync.
+ fireEvent.click(submitButton());
+ fireEvent.click(submitButton());
+
+ expect(mutateAsync).toHaveBeenCalledTimes(1);
+ resolveSpawn({ already_running: false });
+ await waitFor(() => expect(toastSuccess).toHaveBeenCalledTimes(1));
+ expect(mutateAsync).toHaveBeenCalledTimes(1);
+ });
+});
diff --git a/panel/src/components/agents/spawn-agent-dialog.tsx b/panel/src/components/agents/spawn-agent-dialog.tsx
index 7194be8a..7b604f42 100644
--- a/panel/src/components/agents/spawn-agent-dialog.tsx
+++ b/panel/src/components/agents/spawn-agent-dialog.tsx
@@ -1,7 +1,8 @@
"use client";
-import { useState } from "react";
+import { useRef, useState } from "react";
import { useSpawnAgent } from "@/hooks/use-agents";
+import { getErrorMessage } from "@/lib/api/client";
import { Button } from "@/components/ui/button";
import { Input } from "@/components/ui/input";
import { Label } from "@/components/ui/label";
@@ -32,21 +33,34 @@ export function SpawnAgentDialog({
const [taskId, setTaskId] = useState("");
const [initialPrompt, setInitialPrompt] = useState("");
const spawnAgent = useSpawnAgent();
+ // Synchronous re-entrancy guard: `spawnAgent.isPending` only flips on a
+ // re-render, which lags a fast double-click/double-fire by a tick or two —
+ // the guard below blocks a second call within the same synchronous burst
+ // regardless of render timing.
+ const submittingRef = useRef(false);
const handleSpawn = async () => {
+ if (submittingRef.current) return;
+ submittingRef.current = true;
try {
- await spawnAgent.mutateAsync({
+ const result = await spawnAgent.mutateAsync({
agentId,
request: {
task_id: taskId || undefined,
initial_prompt: initialPrompt || undefined,
},
});
- toast.success(`Agent ${agentName} spawned successfully`);
+ if (result.already_running) {
+ toast.info(`Agent ${agentName} already running — spawn skipped`);
+ } else {
+ toast.success(`Agent ${agentName} spawned successfully`);
+ }
setOpen(false);
resetForm();
- } catch {
- toast.error("Failed to spawn agent");
+ } catch (error) {
+ toast.error(getErrorMessage(error));
+ } finally {
+ submittingRef.current = false;
}
};
diff --git a/panel/src/lib/api/a2a.ts b/panel/src/lib/api/a2a.ts
index 4de0ab7a..842936de 100644
--- a/panel/src/lib/api/a2a.ts
+++ b/panel/src/lib/api/a2a.ts
@@ -88,9 +88,9 @@ export interface A2AChatMessage {
}
/**
- * CEO reply payload. The backend sends a DIRECT CEO -> to_agent message (it
- * lands in the CEO<->to_agent pairwise conversation), not an injection into
- * the watched transcript.
+ * CEO interjection payload. The backend posts this INTO the conversation
+ * being viewed (readable by both participants), addressed to `to_agent` —
+ * not a re-homed CEO<->to_agent pairwise DM.
*/
export interface AdminReplyRequest {
to_agent: string;
diff --git a/panel/src/types/index.ts b/panel/src/types/index.ts
index 0e547ff2..c4e23258 100644
--- a/panel/src/types/index.ts
+++ b/panel/src/types/index.ts
@@ -420,6 +420,9 @@ export interface AgentStatusResponse {
error_count: number;
started_at: string | null;
waiting_for: string | null;
+ // Only set on the spawn response: true when the spawn was a no-op because
+ // the agent was already active (see SpawnAgentResponse on the backend).
+ already_running?: boolean;
}
// OrchestratorStatusResponse from backend
diff --git a/roboco/agent_sdk/server.py b/roboco/agent_sdk/server.py
index 5bcb8676..8258a331 100644
--- a/roboco/agent_sdk/server.py
+++ b/roboco/agent_sdk/server.py
@@ -46,7 +46,10 @@ from roboco.agent_sdk.transcript_usage import (
)
from roboco.agents_config import get_agent_team
from roboco.foundation.policy.agent_loop import DEFAULT_BUDGET as _BUDGET
-from roboco.foundation.policy.agent_loop import retry_limit_for
+from roboco.foundation.policy.agent_loop import (
+ absolute_retry_limit_for,
+ retry_limit_for,
+)
from roboco.services.gateway.envelope import Envelope
logger = structlog.get_logger()
@@ -448,6 +451,12 @@ class _SessionState:
self.verb_attempts: dict[tuple[str, str | None], deque[float]] = defaultdict(
deque
)
+ # Session-scoped, never-pruned cumulative rejection count per (verb,
+ # task_id) — catches slow-drip retries that space out past the 60s
+ # window above (see foundation.agent_loop VERB_ABSOLUTE_RETRY_MULTIPLIER).
+ self.verb_absolute_attempts: dict[tuple[str, str | None], int] = defaultdict(
+ int
+ )
# Cumulative token usage for this session. Populated by /usage/sync,
# which parses the Claude Code transcript and *sets* these absolutely
# (the additive /usage/report path remains for explicit deltas).
@@ -568,6 +577,63 @@ def _check_verb_circuit(verb: str, task_id: str | None) -> dict[str, Any] | None
return env.as_dict()
+def _record_verb_attempt_absolute(verb: str, task_id: str | None) -> None:
+ """Bump the never-pruned cumulative rejection count for (verb, task_id).
+
+ Companion to `_record_verb_attempt` — same key, but this one never
+ decays, so it still accumulates when rejections are spaced past the
+ 60s window (the slow-drip case the sliding breaker alone misses).
+ """
+ _state.verb_absolute_attempts[(verb, task_id)] += 1
+
+
+def _verb_absolute_attempt_count(verb: str, task_id: str | None) -> int:
+ """Cumulative rejection count for (verb, task_id); 0 for unseen keys."""
+ return _state.verb_absolute_attempts.get((verb, task_id), 0)
+
+
+def _check_verb_absolute_circuit(
+ verb: str, task_id: str | None
+) -> dict[str, Any] | None:
+ """Return a circuit_open envelope dict if the ABSOLUTE session cap is hit.
+
+ Independent of window pruning — trips a slow-drip retry (one rejection
+ every few minutes) that never accumulates enough in any single 60s
+ window to trip `_check_verb_circuit`.
+ """
+ cap = absolute_retry_limit_for(verb)
+ if cap is None:
+ return None
+ count = _verb_absolute_attempt_count(verb, task_id)
+ if count < cap:
+ return None
+ env = Envelope.circuit_open(
+ verb=verb,
+ attempts=count,
+ window_seconds=_VERB_ATTEMPT_WINDOW_S,
+ message=(
+ f"verb {verb!r} rejected {count} times this session "
+ f"(absolute cap {cap}) — circuit breaker open"
+ ),
+ remediate=(
+ f"verb {verb!r} has been rejected {count} times this session "
+ f"(absolute cap {cap}, regardless of pacing). Stop retrying. Call "
+ "i_am_blocked(reason='unable to satisfy gate after N attempts') "
+ "or i_am_idle() to release the claim. The PM will pick it up."
+ ),
+ )
+ return env.as_dict()
+
+
+def _check_any_verb_circuit(verb: str, task_id: str | None) -> dict[str, Any] | None:
+ """Windowed breaker first (the common fast-storm case), then the
+ session-scoped absolute cap (catches the slow-drip case the window
+ empties between)."""
+ return _check_verb_circuit(verb, task_id) or _check_verb_absolute_circuit(
+ verb, task_id
+ )
+
+
@app.post("/verb/attempted", response_model=VerbCircuitStatus)
async def verb_attempted(req: VerbAttemptRequest) -> VerbCircuitStatus:
"""Record a verb-level rejection and report breaker state.
@@ -581,17 +647,17 @@ async def verb_attempted(req: VerbAttemptRequest) -> VerbCircuitStatus:
"""
if req.rejection_kind in _CIRCUIT_REJECTION_KINDS:
_record_verb_attempt(req.verb, req.task_id)
+ _record_verb_attempt_absolute(req.verb, req.task_id)
limit = retry_limit_for(req.verb)
count = _verb_attempt_count(req.verb, req.task_id)
- is_open = limit is not None and count >= limit
- envelope_dict = _check_verb_circuit(req.verb, req.task_id) if is_open else None
+ envelope_dict = _check_any_verb_circuit(req.verb, req.task_id)
return VerbCircuitStatus(
verb=req.verb,
task_id=req.task_id,
attempts=count,
limit=limit,
window_seconds=_VERB_ATTEMPT_WINDOW_S,
- open=is_open,
+ open=envelope_dict is not None,
circuit_envelope=envelope_dict,
)
@@ -603,15 +669,14 @@ async def verb_circuit_status(
"""Read-only breaker state for (verb, task_id) — does NOT record an attempt."""
limit = retry_limit_for(verb)
count = _verb_attempt_count(verb, task_id)
- is_open = limit is not None and count >= limit
- envelope_dict = _check_verb_circuit(verb, task_id) if is_open else None
+ envelope_dict = _check_any_verb_circuit(verb, task_id)
return VerbCircuitStatus(
verb=verb,
task_id=task_id,
attempts=count,
limit=limit,
window_seconds=_VERB_ATTEMPT_WINDOW_S,
- open=is_open,
+ open=envelope_dict is not None,
circuit_envelope=envelope_dict,
)
diff --git a/roboco/api/deps.py b/roboco/api/deps.py
index d3ebd557..1288b3ac 100644
--- a/roboco/api/deps.py
+++ b/roboco/api/deps.py
@@ -27,7 +27,7 @@ from roboco.api.auth.backend import (
from roboco.api.auth.session import resolve_session_user
from roboco.api.schemas.optimal import PaginationParams
from roboco.config import settings
-from roboco.db.base import get_db
+from roboco.db.base import get_db, get_db_committed
from roboco.db.tables import AgentTable, UserTable
from roboco.foundation.identity import BOARD_ROLES, DEV_ROLES, PM_ROLES, Role
from roboco.models import AgentRole, Team
@@ -53,8 +53,10 @@ logger = structlog.get_logger()
if TYPE_CHECKING:
from collections.abc import Callable, Coroutine
-# Type alias for database session dependency
-DbSession = Annotated[AsyncSession, Depends(get_db)]
+# Type alias for database session dependency. get_db_committed stashes the
+# session on request.state so DbCommitMiddleware can commit it before the
+# response reaches the client (see roboco/db/base.py, roboco/api/middleware.py).
+DbSession = Annotated[AsyncSession, Depends(get_db_committed)]
async def resolve_agent_id(agent_id_str: str, db: AsyncSession) -> UUID:
diff --git a/roboco/api/middleware.py b/roboco/api/middleware.py
index 77a62af7..dd90fa51 100644
--- a/roboco/api/middleware.py
+++ b/roboco/api/middleware.py
@@ -478,10 +478,16 @@ def setup_middleware(app: FastAPI) -> None:
app.add_exception_handler(Exception, generic_exception_handler)
# Middleware (added in reverse order due to LIFO): the LAST add_middleware
- # call is the OUTERMOST. FlowVerbTimeoutMiddleware is added FIRST so it is
- # the INNERMOST — closest to the routes — meaning correlation + logging
- # still wrap the 504 it returns, AND its asyncio.timeout cancels the route
- # coroutine + its get_db dependency directly (same task, reliable cancel).
+ # call is the OUTERMOST. DbCommitMiddleware is added FIRST so it is the
+ # INNERMOST of all four — closest to the routes, right next to CORS —
+ # and, critically, INSIDE FlowVerbTimeoutMiddleware: a hanging commit on a
+ # flow-verb request stays bounded by Flow's asyncio.timeout, and Flow's
+ # own synthesized 504 (sent via its own upstream `send`, never re-entering
+ # `self.app`) never reaches DbCommitMiddleware at all. FlowVerbTimeoutMiddleware
+ # is added next so correlation + logging still wrap the 504 it returns,
+ # AND its asyncio.timeout cancels the route coroutine + its get_db
+ # dependency directly (same task, reliable cancel).
+ app.add_middleware(DbCommitMiddleware)
app.add_middleware(FlowVerbTimeoutMiddleware)
app.add_middleware(RequestLoggingMiddleware)
app.add_middleware(CorrelationIdMiddleware)
@@ -499,9 +505,11 @@ class FlowVerbTimeoutMiddleware:
``kimi-k2.7-code:cloud`` agent on task 79d686f0). This wraps each
``/api/v1/flow/*`` request in ``asyncio.timeout``; on expiry the inner
app is cancelled (CancelledError propagates through ``get_db``, which now
- rolls back, releasing the lock) and a clean retryable 504 envelope is
- returned. Pure ASGI (not BaseHTTPMiddleware) so cancellation propagates
- into the route coroutine without the spawned-task gap.
+ invalidates the session — releasing the lock and discarding a connection
+ that may be mid-protocol rather than reusing it via a rollback, see
+ ``get_db``) and a clean retryable 504 envelope is returned. Pure ASGI (not
+ BaseHTTPMiddleware) so cancellation propagates into the route coroutine
+ without the spawned-task gap.
Reads (``evidence``) and journal writes (``note``) don't touch the task
row, so they are unaffected; only task-row writes route through ``claim``.
@@ -539,8 +547,9 @@ class FlowVerbTimeoutMiddleware:
async with asyncio.timeout(timeout):
await self.app(scope, receive, send_wrapper)
except TimeoutError:
- # The inner app was cancelled mid-verb; get_db has already rolled
- # back (releasing the FOR UPDATE lock) by the time we get here.
+ # The inner app was cancelled mid-verb; get_db has already
+ # invalidated the session (releasing the FOR UPDATE lock and
+ # discarding the connection) by the time we get here.
if started:
# The route had already begun a response before the timeout
# fired — the client owns whatever was sent; we cannot start
@@ -577,3 +586,74 @@ class FlowVerbTimeoutMiddleware:
# Client may have already disconnected (the original trigger);
# the lock is released regardless. Nothing to do.
pass
+
+
+class DbCommitMiddleware:
+ """Commits the request's DB session before the response reaches the client.
+
+ FastAPI resolves ``Depends(get_db)`` on the request-scoped ``AsyncExitStack``
+ and sends the response (``fastapi/routing.py``'s ``request_response``,
+ ``await response(scope, receive, send)``) BEFORE that stack unwinds and
+ runs ``get_db``'s post-yield ``await session.commit()``. So every write
+ endpoint that relies on it returns 200 while its commit is still pending —
+ a follow-up request (or a fresh connection) can read pre-commit state —
+ and a commit that later FAILS leaves the client told "ok" with nothing
+ persisted.
+
+ ``get_db_committed`` (``roboco/db/base.py`` — the ``roboco.api.deps.DbSession``
+ target every route depends on) stashes the live session on
+ ``request.state.db_session`` before yielding. This middleware wraps
+ ``send``: on the FIRST ``http.response.start`` it commits that session
+ BEFORE forwarding the event, so the client only ever sees the response
+ after the commit lands. A commit failure rolls back and re-raises — the
+ response hasn't started, so the surrounding exception-handling machinery
+ (FastAPI's handlers / Starlette's ``ServerErrorMiddleware``) turns it into
+ a clean 500 instead of the silent post-200 loss. ``session.in_transaction()``
+ makes the check idempotent and skips the exception path for free: an
+ exception rolls back (and closes) the session inside ``get_db`` before
+ its error response is built, so by the time THAT response's
+ ``http.response.start`` reaches here there is no open transaction left
+ to commit.
+
+ Added INSIDE (closer to the routes than) ``FlowVerbTimeoutMiddleware`` —
+ see ``setup_middleware`` — so a hanging commit on a flow-verb request
+ stays bounded by Flow's ``asyncio.timeout``. That timeout is scoped to
+ the WHOLE ``self.app(...)`` call including this middleware, so its
+ deadline can fire while ``await session.commit()`` below is itself
+ in flight (not just while a route handler hangs before responding) —
+ ``started`` in the outer middleware is still ``False`` at that point
+ (its own wrapped ``send`` hasn't been called yet), so it sends its 504
+ normally once this ``await`` raises ``CancelledError``. That
+ ``CancelledError`` is a ``BaseException`` the ``except Exception`` below
+ does not catch, so it propagates up through FastAPI's dependency
+ ``AsyncExitStack`` (still open here — ``response(scope, receive, send)``
+ is called from inside it, see ``get_db_committed``'s docstring) straight
+ into ``get_db``'s own ``except asyncio.CancelledError``, which invalidates
+ the session rather than rolling it back: a rollback would issue another
+ command over a connection whose wire-protocol state this cancellation may
+ have already left mid-flight, corrupting it further (SQLAlchemy's own
+ docs prescribe ``invalidate()``, not ``rollback()``, for this exact
+ external-cancellation case). Skipping that step is what let a later,
+ unrelated request's pool checkout crash on the poisoned connection.
+ """
+
+ def __init__(self, app: ASGIApp) -> None:
+ self.app = app
+
+ async def __call__(self, scope: Scope, receive: Receive, send: Send) -> None:
+ if scope["type"] != "http":
+ await self.app(scope, receive, send)
+ return
+
+ async def send_wrapper(message: Any) -> None:
+ if message["type"] == "http.response.start":
+ session = scope.get("state", {}).get("db_session")
+ if session is not None and session.in_transaction():
+ try:
+ await session.commit()
+ except Exception:
+ await session.rollback()
+ raise
+ await send(message)
+
+ await self.app(scope, receive, send_wrapper)
diff --git a/roboco/api/routes/a2a.py b/roboco/api/routes/a2a.py
index 8a5310c1..b1c2dc5e 100644
--- a/roboco/api/routes/a2a.py
+++ b/roboco/api/routes/a2a.py
@@ -1078,10 +1078,11 @@ async def reply_as_ceo(
agent: CurrentAgentContext,
data: AdminReplyRequest,
) -> MessageResponse:
- """CEO-only: chime into an existing A2A conversation as itself.
+ """CEO-only: interject into an existing A2A conversation as itself.
- The CEO addresses one of the conversation's two real participants (A2A
- conversations are strictly pairwise) on the conversation's linked task.
+ A one-directional interjection, not a CEO<->agent DM: the message is
+ inserted into THIS conversation (readable by both participants) and
+ addressed to one of its two real participants via ``interject_as_ceo``.
"""
_require_ceo(agent)
service = A2AService(db)
@@ -1094,23 +1095,12 @@ async def reply_as_ceo(
)
_resolve_reply_target(conv, data.to_agent)
- try:
- msg = await service.send(
- from_agent=agent.agent_id,
- to_agent=data.to_agent,
- task_id=require_uuid(conv.task_id),
- body=data.content,
- skill=data.skill,
- )
- except A2AAccessDeniedError as e:
- raise HTTPException(
- status_code=status.HTTP_403_FORBIDDEN,
- detail={
- "error": "A2A_ACCESS_DENIED",
- "message": e.message,
- "route_hint": e.route_hint,
- },
- ) from None
+ msg = await service.interject_as_ceo(
+ conversation_id=require_uuid(conversation_id),
+ to_agent=data.to_agent,
+ content=data.content,
+ skill=data.skill,
+ )
await db.commit()
diff --git a/roboco/api/routes/orchestrator.py b/roboco/api/routes/orchestrator.py
index 31d736a1..7a0f5314 100644
--- a/roboco/api/routes/orchestrator.py
+++ b/roboco/api/routes/orchestrator.py
@@ -6,6 +6,7 @@ API endpoints for managing the Agent Orchestrator.
from datetime import datetime
from typing import Annotated
+from uuid import UUID
from fastapi import APIRouter, Cookie, Depends, Header, HTTPException, status
from guard_core.handlers.behavior_handler import BehaviorRule
@@ -25,10 +26,16 @@ from roboco.api.schemas.orchestrator import (
OrchestratorStatusResponse,
ResolveWaitRequest,
SpawnAgentRequest,
+ SpawnAgentResponse,
WaitingAgentResponse,
)
from roboco.config import settings
+from roboco.db.base import get_db_context
+from roboco.db.tables import TaskTable
+from roboco.runtime import AgentState
+from roboco.runtime.orchestrator import AgentReadinessError
from roboco.security import guard_deco, prompt_injection_validator
+from roboco.services.task import get_task_service
_RUNAWAY_RULES = [
BehaviorRule(rule_type="frequency", threshold=120, window=60, action="log")
@@ -213,9 +220,55 @@ async def get_waiting_agents() -> list[WaitingAgentResponse]:
]
+def _build_manual_spawn_prompt(task: TaskTable, ceo_note: str | None) -> str:
+ """Build the initial prompt for a CEO-triggered manual (panel) spawn.
+
+ Mirrors the tone of dispatcher-built prompts (e.g. ``_build_pr_review_prompt``
+ in the orchestrator): point the agent at the task by id/title/status and
+ trust the gateway envelope's ``next`` / ``remediate`` to guide the actual
+ claim verb, rather than enumerating per-role verbs here.
+ """
+ lines = [
+ "You were manually spawned by the CEO to work a specific task.",
+ "",
+ f"TASK ID: {task.id}",
+ f"TITLE: {task.title}",
+ f"STATUS: {task.status.value}",
+ "",
+ "Claim it with the claim verb appropriate to your role and this "
+ "task's current state, then proceed. Trust the gateway envelope's "
+ "`next` / `remediate` fields to guide you rather than guessing.",
+ ]
+ if ceo_note:
+ lines += ["", "== CEO NOTE ==", ceo_note]
+ return "\n".join(lines)
+
+
+async def _resolve_manual_spawn_prompt(
+ task_id: str | None, ceo_message: str | None
+) -> str | None:
+ """Best-effort task-aware prompt for a manual panel spawn.
+
+ Falls back to ``ceo_message`` unchanged (current behavior) on any lookup
+ failure — bad ``task_id``, DB hiccup, task not found. Enrichment must
+ never block a spawn the CEO already asked for; ``spawn_agent``'s own
+ readiness gate is the real gatekeeper for an invalid/not-ready task.
+ """
+ if not task_id:
+ return ceo_message
+ try:
+ async with get_db_context() as db:
+ task = await get_task_service(db).get(UUID(task_id))
+ except Exception:
+ return ceo_message
+ if task is None:
+ return ceo_message
+ return _build_manual_spawn_prompt(task, ceo_message)
+
+
@router.post(
"/agents/{agent_id}/spawn",
- response_model=AgentStatusResponse,
+ response_model=SpawnAgentResponse,
status_code=status.HTTP_201_CREATED,
summary="Spawn agent",
description="Spawn a Claude Code instance for an agent.",
@@ -230,19 +283,42 @@ async def get_waiting_agents() -> list[WaitingAgentResponse]:
async def spawn_agent(
agent_id: str,
data: SpawnAgentRequest | None = None,
-) -> AgentStatusResponse:
+) -> SpawnAgentResponse:
"""Spawn an agent."""
agent_id = _validated_agent_id(agent_id)
orchestrator = get_orchestrator()
+ task_id = data.task_id if data else None
+ ceo_message = data.initial_prompt if data else None
+ prompt = await _resolve_manual_spawn_prompt(task_id, ceo_message)
+
+ # Pre-check for already-running signaling (see return below). Snapshot the
+ # instance identity BEFORE calling spawn_agent, which silently reuses a
+ # running instance rather than erroring — dispatchers rely on that no-op
+ # contract, so it stays untouched here.
+ pre_existing = orchestrator.get_instance(agent_id)
+ pre_active = pre_existing is not None and pre_existing.state not in (
+ AgentState.OFFLINE,
+ AgentState.WAITING_LONG,
+ )
+ pre_existing_id = getattr(pre_existing, "id", None)
try:
instance = await orchestrator.spawn_agent(
agent_id=agent_id,
- initial_prompt=data.initial_prompt if data else None,
- task_id=data.task_id if data else None,
+ initial_prompt=prompt,
+ task_id=task_id,
model=data.model if data else None,
spawned_by="api.orchestrator.spawn",
)
+ except AgentReadinessError as e:
+ # Expected, well-formed refusal (role/state mismatch, unmet
+ # dependency, missing readiness criteria) — not a server crash.
+ # 409 keeps it out of 5xx alerting and lets the panel surface the
+ # real reason instead of a generic "server error".
+ raise HTTPException(
+ status_code=status.HTTP_409_CONFLICT,
+ detail=str(e),
+ ) from e
except FileNotFoundError as e:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
@@ -254,13 +330,32 @@ async def spawn_agent(
detail=f"Failed to spawn agent: {e}",
) from e
- return AgentStatusResponse(
+ # already_running: the pre-existing instance was active AND spawn_agent
+ # handed back that exact same instance (identity, not state, since a
+ # freshly-launched instance can share the same STARTING state as a
+ # short-circuited one). AgentInstance.id is a fresh uuid4 per constructed
+ # object, so equality here means no new instance was built.
+ # ponytail: identity-compare across a pre/post HTTP-handler snapshot, not
+ # inside the orchestrator's own spawn lock — a genuinely simultaneous
+ # double-fire that races both pre-checks before either inserts its
+ # instance can still slip through undetected here. The client-side
+ # dedupe guard (SpawnAgentDialog) is the actual fix for that race;
+ # upgrade this to an orchestrator-native signal if that ever proves
+ # insufficient.
+ already_running = (
+ pre_active
+ and pre_existing_id is not None
+ and getattr(instance, "id", None) == pre_existing_id
+ )
+
+ return SpawnAgentResponse(
agent_id=instance.agent_id,
state=instance.state.value,
task_id=instance.current_task_id,
error_count=instance.error_count,
started_at=instance.started_at,
waiting_for=None,
+ already_running=already_running,
)
diff --git a/roboco/api/schemas/orchestrator.py b/roboco/api/schemas/orchestrator.py
index cc5a08db..91ee52fe 100644
--- a/roboco/api/schemas/orchestrator.py
+++ b/roboco/api/schemas/orchestrator.py
@@ -49,6 +49,19 @@ class SpawnAgentRequest(BaseModel):
model: str | None = None
+class SpawnAgentResponse(AgentStatusResponse):
+ """Response to a spawn request.
+
+ ``already_running`` is True when the spawn was a no-op because the agent
+ was already active — the route detected a pre-existing active instance
+ whose identity matches the instance returned by ``spawn_agent`` (which
+ silently reuses a running instance rather than erroring, so dispatchers
+ keep their existing no-op semantics). False for a genuine new spawn.
+ """
+
+ already_running: bool = False
+
+
class ResolveWaitRequest(BaseModel):
"""Request to resolve a wait condition."""
diff --git a/roboco/db/base.py b/roboco/db/base.py
index eb09cd31..eaefa7ba 100644
--- a/roboco/db/base.py
+++ b/roboco/db/base.py
@@ -6,10 +6,12 @@ import asyncio
from collections.abc import AsyncGenerator
from contextlib import asynccontextmanager
from pathlib import Path
+from typing import Annotated
import structlog
from alembic import command
from alembic.config import Config
+from fastapi import Depends, Request
from sqlalchemy import MetaData, text
from sqlalchemy.ext.asyncio import (
AsyncConnection,
@@ -82,22 +84,82 @@ async def get_db() -> AsyncGenerator[AsyncSession]:
@router.get("/items")
async def get_items(db: AsyncSession = Depends(get_db)):
...
+
+ Also called directly (no ``Depends``, outside HTTP request scope) by
+ ``websocket.py`` and a couple of read-only helpers — keep this signature
+ free of any required/HTTP-only parameter. For a request-scoped route that
+ wants its commit to land BEFORE the response reaches the client, depend
+ on ``get_db_committed`` instead (``roboco.api.deps.DbSession`` — the one
+ place every route already goes through — already does).
"""
session_factory = get_session_factory()
async with session_factory() as session:
try:
yield session
await session.commit()
- except (Exception, asyncio.CancelledError):
- # CancelledError is BaseException, so the bare `except Exception`
- # did not catch it — a server-side asyncio.timeout cancelling a
- # hung verb (FlowVerbTimeoutMiddleware) would otherwise leave the
- # request transaction unrolled-back, holding its FOR UPDATE row
- # lock. Roll back on cancellation too so the lock releases.
+ except asyncio.CancelledError:
+ await _discard_on_cancel(session)
+ raise
+ except Exception:
await session.rollback()
raise
+async def _discard_on_cancel(session: AsyncSession) -> None:
+ """Discard (never reuse) a session cancelled mid-flight.
+
+ A server-side ``asyncio.timeout`` (``FlowVerbTimeoutMiddleware``) can fire
+ while the session is mid ``await`` on a real DBAPI round-trip — not just
+ while idle holding a ``FOR UPDATE`` lock, but also mid-``commit()``
+ (``DbCommitMiddleware`` runs its own commit in the ASGI send path, still
+ inside the same cancellable scope). Cancelling a greenlet-bridged asyncpg
+ operation mid-flight leaves the connection's wire-protocol state
+ undefined; SQLAlchemy's own docs (``Session.invalidate``) prescribe
+ exactly this: on a Timeout/cancellation, invalidate rather than rollback,
+ since rollback() itself would issue another command over a connection
+ that may already be desynced, and a desynced connection returned to the
+ pool is what later corrupted a *different* request's checkout (the
+ uvloop/asyncpg segfault class this fixes). A plain hang (asyncio.sleep,
+ no DBAPI call in flight when cancelled) is also safe to invalidate — just
+ slightly more heavy-handed than the rollback it used to get.
+ """
+ try:
+ await session.invalidate()
+ except Exception as e:
+ # The connection is already being discarded; a failure tearing it
+ # down further (SQLAlchemy's own pool logs the underlying cause) must
+ # not mask the CancelledError the caller is propagating.
+ logger.debug("Session invalidate-on-cancel raised", error=str(e))
+
+
+async def get_db_committed(
+ request: Request, db: Annotated[AsyncSession, Depends(get_db)]
+) -> AsyncGenerator[AsyncSession]:
+ """FastAPI-only wrapper around ``get_db``: stashes the live session on
+ ``request.state.db_session`` so ``DbCommitMiddleware``
+ (``roboco/api/middleware.py``) can commit it BEFORE the response reaches
+ the client.
+
+ FastAPI resolves ``Depends(get_db)`` on the request-scoped exit stack, and
+ its routing sends the response to the client BEFORE that stack unwinds —
+ so ``get_db``'s post-yield ``commit()`` used to land after a 200 already
+ went out.
+
+ This is a separate function rather than a ``request`` parameter added to
+ ``get_db`` itself: FastAPI only special-cases a dependency parameter
+ typed exactly ``Request`` (``lenient_issubclass`` in
+ ``fastapi/dependencies/utils.py``) — a ``Request | None`` union is NOT
+ special-cased and instead gets validated as a Pydantic response field,
+ which crashes route registration outright (``Request`` isn't a valid
+ Pydantic field type). ``get_db`` is also called directly with no request
+ in scope, so its signature has to stay request-free; this wrapper is the
+ request-scoped variant, resolved once per request (FastAPI dependency
+ caching) so ``db`` here is the exact same session ``get_db`` yields.
+ """
+ request.state.db_session = db
+ yield db
+
+
@asynccontextmanager
async def get_db_context() -> AsyncGenerator[AsyncSession]:
"""
@@ -112,12 +174,10 @@ async def get_db_context() -> AsyncGenerator[AsyncSession]:
try:
yield session
await session.commit()
- except (Exception, asyncio.CancelledError):
- # CancelledError is BaseException, so the bare `except Exception`
- # did not catch it — a server-side asyncio.timeout cancelling a
- # hung verb (FlowVerbTimeoutMiddleware) would otherwise leave the
- # request transaction unrolled-back, holding its FOR UPDATE row
- # lock. Roll back on cancellation too so the lock releases.
+ except asyncio.CancelledError:
+ await _discard_on_cancel(session)
+ raise
+ except Exception:
await session.rollback()
raise
diff --git a/roboco/foundation/policy/agent_loop.py b/roboco/foundation/policy/agent_loop.py
index f9f9301d..4ab8408f 100644
--- a/roboco/foundation/policy/agent_loop.py
+++ b/roboco/foundation/policy/agent_loop.py
@@ -12,6 +12,13 @@ no per-verb retry cap — dogfooding showed i_am_done retried 5+ times in 2
minutes within the global budget. With the runtime tracker in place,
exceeding VERB_RETRY_LIMITS[verb] attempts in 60s returns
Envelope.circuit_open.
+
+VERB_ABSOLUTE_RETRY_MULTIPLIER is NEW. A 2026-07-08 production loop showed
+the 60s sliding window never trips on a slow drip — one rejected i_am_done
+every 3-4 minutes empties the window between attempts, so the agent ground
+for 30+ minutes without the breaker ever seeing more than 1 attempt at a
+time. absolute_retry_limit_for() adds a session-scoped, never-pruned
+cumulative cap alongside the window so pacing can't defeat the breaker.
"""
from __future__ import annotations
@@ -123,3 +130,22 @@ def retry_limit_for(verb: str) -> int | None:
if verb in UNLIMITED_RETRY_VERBS:
return None
return VERB_RETRY_LIMITS.get(verb, DEFAULT_BUDGET.verb_retry_max_per_minute)
+
+
+# Multiplier applied to retry_limit_for(verb) for the session-scoped
+# ABSOLUTE cap (see module docstring). i_am_done's windowed cap is 3, so its
+# absolute cap is 9 total rejections in one container session, regardless
+# of how the attempts are spaced.
+VERB_ABSOLUTE_RETRY_MULTIPLIER: int = 3
+
+
+def absolute_retry_limit_for(verb: str) -> int | None:
+ """Session-scoped cumulative cap: retry_limit_for(verb) * the multiplier.
+
+ None for verbs retry_limit_for treats as unlimited — a verb exempt from
+ the windowed breaker is exempt from the absolute one too.
+ """
+ limit = retry_limit_for(verb)
+ if limit is None:
+ return None
+ return limit * VERB_ABSOLUTE_RETRY_MULTIPLIER
diff --git a/roboco/runtime/orchestrator.py b/roboco/runtime/orchestrator.py
index 0b147c05..dadb2b07 100644
--- a/roboco/runtime/orchestrator.py
+++ b/roboco/runtime/orchestrator.py
@@ -3143,6 +3143,11 @@ class AgentOrchestrator:
# every gateway call 422s on header parse. Resolve via AGENT_UUIDS map;
# if the slug isn't in the map (custom agents), fall back to the slug
# and let the API surface the unknown-agent error.
+ # Also used as the CLI arg for the three ApiClient-based servers
+ # (optimal/docs/search) below — their spawn token (issue_agent_token)
+ # is signed over the UUID, so ApiClient's X-Agent-ID must match or
+ # verify_agent_token 401s with "signature mismatch" even though
+ # get_agent_role/get_agent_team resolve either form fine.
agent_uuid = AGENT_UUIDS.get(agent_id, agent_id)
mcp_env: dict[str, str] = {
@@ -3211,7 +3216,7 @@ class AgentOrchestrator:
"python",
"-m",
"roboco.mcp.optimal_server",
- agent_id,
+ agent_uuid,
],
"env": mcp_env,
},
@@ -3236,7 +3241,7 @@ class AgentOrchestrator:
"python",
"-m",
"roboco.mcp.docs_server",
- agent_id,
+ agent_uuid,
],
"env": mcp_env,
}
@@ -3259,7 +3264,7 @@ class AgentOrchestrator:
"python",
"-m",
"roboco.mcp.search_server",
- agent_id,
+ agent_uuid,
],
"env": mcp_env,
}
diff --git a/roboco/services/a2a.py b/roboco/services/a2a.py
index 13c229c7..7fbacf75 100644
--- a/roboco/services/a2a.py
+++ b/roboco/services/a2a.py
@@ -1838,6 +1838,66 @@ class A2AService:
)
return msg
+ async def interject_as_ceo(
+ self,
+ conversation_id: UUID,
+ to_agent: str,
+ content: str,
+ skill: str | None = None,
+ ) -> A2AChatMessage:
+ """CEO interjection: post a message directly into an existing
+ agent<->agent conversation, addressed to one of its participants.
+
+ One-directional and NOT a participant send: unlike
+ ``send_chat_message`` (which requires the sender to be a party to
+ the conversation), the CEO here is watching and interjecting into
+ someone else's thread, not conversing in its own — so the
+ participant check on the sender is deliberately bypassed rather
+ than weakened for every other caller. ``to_agent`` still must be
+ one of the conversation's two real participants.
+
+ Only ``to_agent``'s unread counter is bumped (a ping to whoever
+ it's addressed to); the other participant still sees the row via
+ the shared transcript / ``read_a2a``, just without a ping.
+
+ Direction is encoded as an ``@{to_agent}: `` content prefix —
+ ponytail: no ``to_agent`` column yet; add one (and stop parsing
+ the prefix) if the panel ever needs to render/filter by recipient
+ directly instead.
+ """
+ conv = await self.session.get(A2AConversationTable, conversation_id)
+ if conv is None:
+ raise ValueError(f"Conversation not found: {conversation_id}")
+ if to_agent not in (conv.agent_a, conv.agent_b):
+ raise ValueError(
+ f"{to_agent} is not a participant in this conversation "
+ f"(participants: {conv.agent_a}, {conv.agent_b})"
+ )
+
+ msg = A2AMessageTable(
+ conversation_id=conversation_id,
+ from_agent="ceo",
+ content=f"@{to_agent}: {content}",
+ message_kind=A2AMessageKind.MESSAGE,
+ skill=skill,
+ )
+ self.session.add(msg)
+
+ conv.message_count += 1
+ conv.last_message_at = datetime.now(UTC)
+ if to_agent == conv.agent_a:
+ conv.unread_by_a += 1
+ else:
+ conv.unread_by_b += 1
+
+ await self.session.flush()
+ await self.session.refresh(msg)
+
+ model = self._msg_to_model(msg)
+ task_id = str(conv.task_id) if conv.task_id else None
+ await self._publish_a2a_message_sent(model, task_id, "ceo", to_agent, skill)
+ return model
+
@staticmethod
async def _publish_a2a_message_sent(
msg: A2AChatMessage,
diff --git a/roboco/services/gateway/choreographer/_impl.py b/roboco/services/gateway/choreographer/_impl.py
index 0033daed..3028d5f1 100644
--- a/roboco/services/gateway/choreographer/_impl.py
+++ b/roboco/services/gateway/choreographer/_impl.py
@@ -654,6 +654,10 @@ class Choreographer:
"missing": env.missing or [],
"attempt_id": str(_uuid4()),
}
+ if env.remediate:
+ # Conventions-gate rejections carry the file:line violation
+ # listing ONLY here — without it the audit row is unactionable.
+ details["remediate"] = env.remediate
cid = structlog.contextvars.get_contextvars().get("correlation_id")
if cid is not None:
details["correlation_id"] = cid
diff --git a/roboco/services/gateway/content_actions.py b/roboco/services/gateway/content_actions.py
index af21e5a2..6beb1ba4 100644
--- a/roboco/services/gateway/content_actions.py
+++ b/roboco/services/gateway/content_actions.py
@@ -559,12 +559,15 @@ class ContentActions:
return reject
canonical_prefix = f"[{str(t.id)[:8]}]"
final_message = f"{canonical_prefix} {subject}"
- commit_result = await self.git.commit(
- branch_name=t.branch_name,
- message=final_message,
- task_id=t.id,
- files=files,
- )
+ try:
+ commit_result = await self.git.commit(
+ branch_name=t.branch_name,
+ message=final_message,
+ task_id=t.id,
+ files=files,
+ )
+ except GitError as exc:
+ return self._commit_git_error_envelope(exc, files=files)
sha = commit_result.get("sha", "")
await self.task.add_progress(
t.id, agent_id, f"committed {sha[:8]}: {final_message}"
@@ -577,6 +580,32 @@ class ContentActions:
context_briefing={},
)
+ @staticmethod
+ def _commit_git_error_envelope(
+ exc: GitError, *, files: list[str] | None
+ ) -> Envelope:
+ """Map a failed `git commit` onto an actionable invalid_state envelope.
+
+ A "no changes added to commit" / "nothing to commit" failure means
+ the passed `files` matched no modified paths; anything else is a
+ generic git failure the agent should inspect and retry.
+ """
+ text = str(exc)
+ if files and (
+ "no changes added to commit" in text or "nothing to commit" in text
+ ):
+ remediate = (
+ f"the files list {files!r} matched no modified paths; omit "
+ "files to stage all changes, or pass the exact modified paths"
+ )
+ else:
+ remediate = "inspect the git error above and retry"
+ return Envelope.invalid_state(
+ message=text,
+ remediate=remediate,
+ context_briefing={},
+ )
+
# Board roles co-review board/coordination tasks: a
# board/coordination task is dispatched to BOTH the Product Owner and the
# Head of Marketing, but it carries a single ``assigned_to``. The
diff --git a/roboco/services/gateway/envelope.py b/roboco/services/gateway/envelope.py
index b9cfdbd5..0d11f60d 100644
--- a/roboco/services/gateway/envelope.py
+++ b/roboco/services/gateway/envelope.py
@@ -183,6 +183,7 @@ class Envelope:
window_seconds: int,
remediate: str,
context_briefing: dict[str, Any] | None = None,
+ message: str | None = None,
) -> Envelope:
"""Per-verb retry circuit-breaker tripped — too many attempts in a window.
@@ -190,10 +191,14 @@ class Envelope:
a structured "stop hammering this verb" signal with a remediate hint
pointing to i_am_blocked() / i_am_idle() as graceful exits. Wired by
the agent_sdk runtime tracker — the gateway itself does not raise this.
+
+ `message` overrides the default windowed wording — used by the
+ session-scoped absolute breaker, whose trip isn't "in last Ns".
"""
return cls(
error="circuit_open",
- message=(
+ message=message
+ or (
f"verb {verb!r} rejected {attempts} times in last "
f"{window_seconds}s — circuit breaker open"
),
diff --git a/tests/integration/test_a2a_routes.py b/tests/integration/test_a2a_routes.py
index 3871d673..fcecff33 100644
--- a/tests/integration/test_a2a_routes.py
+++ b/tests/integration/test_a2a_routes.py
@@ -1237,6 +1237,8 @@ async def test_admin_reply_no_task_id_400(a2a_route_client: dict) -> None:
@pytest.mark.asyncio
async def test_admin_reply_success(a2a_route_client: dict) -> None:
+ """The route posts into the VIEWED conversation via interject_as_ceo —
+ not a re-homed CEO<->target DM (the prior, rejected behavior)."""
app = a2a_route_client["app"]
dev = a2a_route_client["dev"]
client = a2a_route_client["client"]
@@ -1249,7 +1251,7 @@ async def test_admin_reply_success(a2a_route_client: dict) -> None:
id=uuid4(),
conversation_id=conv_id,
from_agent="ceo",
- content="chiming in",
+ content="@be-dev-1: chiming in",
message_kind="message",
response_to_id=None,
requires_response=False,
@@ -1260,7 +1262,7 @@ async def test_admin_reply_success(a2a_route_client: dict) -> None:
with patch("roboco.api.routes.a2a.A2AService") as mock_service_cls:
instance = AsyncMock()
instance.get_conversation_admin = AsyncMock(return_value=conv)
- instance.send = AsyncMock(return_value=sent_msg)
+ instance.interject_as_ceo = AsyncMock(return_value=sent_msg)
mock_service_cls.return_value = instance
response = await client.post(
f"/api/a2a/chat/admin/conversations/{conv_id}/reply",
@@ -1268,44 +1270,13 @@ async def test_admin_reply_success(a2a_route_client: dict) -> None:
headers=_HDR,
)
assert response.status_code == HTTPStatus.CREATED
- instance.send.assert_awaited_once()
- call_kwargs = instance.send.await_args.kwargs
+ instance.interject_as_ceo.assert_awaited_once()
+ call_kwargs = instance.interject_as_ceo.await_args.kwargs
+ assert call_kwargs["conversation_id"] == conv_id
assert call_kwargs["to_agent"] == "be-dev-1"
- assert call_kwargs["task_id"] == task_id
- assert call_kwargs["body"] == "chiming in"
+ assert call_kwargs["content"] == "chiming in"
body = response.json()
- assert body["content"] == "chiming in"
-
-
-@pytest.mark.asyncio
-async def test_admin_reply_access_denied_maps_to_403(a2a_route_client: dict) -> None:
- """Defensive: if send() ever rejects a CEO-authored A2A, surface 403
- rather than crash — mirrors create_conversation's handling."""
- app = a2a_route_client["app"]
- dev = a2a_route_client["dev"]
- client = a2a_route_client["client"]
- _set_ceo_context(app, dev)
-
- conv_id = uuid4()
- task_id = uuid4()
- conv = _admin_conv_obj(conv_id=conv_id, task_id=task_id)
- with patch("roboco.api.routes.a2a.A2AService") as mock_service_cls:
- instance = AsyncMock()
- instance.get_conversation_admin = AsyncMock(return_value=conv)
- instance.send = AsyncMock(
- side_effect=A2AAccessDeniedError(
- from_agent="ceo",
- to_agent="be-dev-1",
- reason="denied",
- )
- )
- mock_service_cls.return_value = instance
- response = await client.post(
- f"/api/a2a/chat/admin/conversations/{conv_id}/reply",
- json={"to_agent": "be-dev-1", "content": "hi"},
- headers=_HDR,
- )
- assert response.status_code == HTTPStatus.FORBIDDEN
+ assert body["content"] == "@be-dev-1: chiming in"
# ---------------------------------------------------------------------------
diff --git a/tests/integration/test_a2a_service.py b/tests/integration/test_a2a_service.py
index 9570f313..15cb5ab2 100644
--- a/tests/integration/test_a2a_service.py
+++ b/tests/integration/test_a2a_service.py
@@ -731,6 +731,127 @@ async def test_get_conversation_admin_returns_none_for_unknown(
assert await svc.get_conversation_admin(uuid4()) is None
+# ---------------------------------------------------------------------------
+# interject_as_ceo — the CEO's one-directional interjection into a watched
+# agent<->agent conversation (not a re-homed CEO<->target DM).
+# ---------------------------------------------------------------------------
+
+
+@pytest.mark.asyncio
+async def test_interject_as_ceo_lands_in_viewed_conversation_with_prefix(
+ a2a_setup: dict,
+) -> None:
+ svc = a2a_setup["svc"]
+ conv = await svc.get_or_create_conversation("be-dev-1", "be-qa")
+
+ msg = await svc.interject_as_ceo(UUID(conv.id), "be-qa", "ship it")
+
+ assert msg.conversation_id == conv.id
+ assert msg.from_agent == "ceo"
+ assert msg.content == "@be-qa: ship it"
+
+
+@pytest.mark.asyncio
+async def test_interject_as_ceo_bumps_target_unread_when_target_is_agent_b(
+ a2a_setup: dict,
+) -> None:
+ """Canonical order makes "be-qa" agent_b — its counter, not agent_a's,
+ must move; the other participant gets no ping."""
+ svc = a2a_setup["svc"]
+ db = a2a_setup["db"]
+ conv = await svc.get_or_create_conversation("be-dev-1", "be-qa")
+ assert conv.agent_a == "be-dev-1"
+ assert conv.agent_b == "be-qa"
+
+ await svc.interject_as_ceo(UUID(conv.id), "be-qa", "ship it")
+
+ row = await db.get(A2AConversationTable, UUID(conv.id))
+ assert row is not None
+ assert row.unread_by_b == 1
+ assert row.unread_by_a == 0
+
+
+@pytest.mark.asyncio
+async def test_interject_as_ceo_bumps_target_unread_when_target_is_agent_a(
+ a2a_setup: dict,
+) -> None:
+ svc = a2a_setup["svc"]
+ db = a2a_setup["db"]
+ conv = await svc.get_or_create_conversation("be-dev-1", "be-qa")
+
+ await svc.interject_as_ceo(UUID(conv.id), "be-dev-1", "ship it")
+
+ row = await db.get(A2AConversationTable, UUID(conv.id))
+ assert row is not None
+ assert row.unread_by_a == 1
+ assert row.unread_by_b == 0
+
+
+@pytest.mark.asyncio
+async def test_interject_as_ceo_bumps_message_count_and_last_message_at(
+ a2a_setup: dict,
+) -> None:
+ svc = a2a_setup["svc"]
+ db = a2a_setup["db"]
+ conv = await svc.get_or_create_conversation("be-dev-1", "be-qa")
+ await svc.send_chat_message(UUID(conv.id), "be-dev-1", "hello")
+
+ await svc.interject_as_ceo(UUID(conv.id), "be-qa", "ship it")
+
+ row = await db.get(A2AConversationTable, UUID(conv.id))
+ assert row is not None
+ _EXPECTED_MESSAGE_COUNT = 2
+ assert row.message_count == _EXPECTED_MESSAGE_COUNT
+ assert row.last_message_at is not None
+
+
+@pytest.mark.asyncio
+async def test_interject_as_ceo_rejects_non_participant_target(
+ a2a_setup: dict,
+) -> None:
+ svc = a2a_setup["svc"]
+ conv = await svc.get_or_create_conversation("be-dev-1", "be-qa")
+
+ with pytest.raises(ValueError, match="not a participant"):
+ await svc.interject_as_ceo(UUID(conv.id), "ghost-agent", "hi")
+
+
+@pytest.mark.asyncio
+async def test_interject_as_ceo_unknown_conversation_raises(
+ a2a_setup: dict,
+) -> None:
+ svc = a2a_setup["svc"]
+ with pytest.raises(ValueError, match="Conversation not found"):
+ await svc.interject_as_ceo(uuid4(), "be-qa", "hi")
+
+
+@pytest.mark.asyncio
+async def test_interject_as_ceo_publishes_a2a_message_sent_event(
+ a2a_setup: dict,
+) -> None:
+ """Same operator-live-view chokepoint as send()/send_chat_message() —
+ the panel's /ws/system invalidation must fire for an interjection too."""
+ svc = a2a_setup["svc"]
+ task_id = a2a_setup["task_id"]
+ conv = await svc.get_or_create_conversation("be-dev-1", "be-qa", task_id=task_id)
+ mock_bus = AsyncMock()
+ mock_bus.is_connected = lambda: True
+ mock_bus.publish = AsyncMock(return_value=None)
+ with patch("roboco.services.a2a.get_event_bus", return_value=mock_bus):
+ sent = await svc.interject_as_ceo(UUID(conv.id), "be-qa", "ship it")
+
+ mock_bus.publish.assert_awaited_once()
+ published = mock_bus.publish.await_args.args[0]
+ assert published.type is EventType.A2A_MESSAGE_SENT
+ data = published.data
+ # Points at the VIEWED conversation, not a re-homed ceo<->target one.
+ assert data["conversation_id"] == conv.id
+ assert data["conversation_id"] == sent.conversation_id
+ assert data["task_id"] == str(task_id)
+ assert data["from_agent"] == "ceo"
+ assert data["to_agent"] == "be-qa"
+
+
# ---------------------------------------------------------------------------
# list_admin_pairs — the A2A switchboard's static-matrix + DB join
# ---------------------------------------------------------------------------
diff --git a/tests/integration/test_db_base.py b/tests/integration/test_db_base.py
index d1f4a7ce..ce590449 100644
--- a/tests/integration/test_db_base.py
+++ b/tests/integration/test_db_base.py
@@ -12,6 +12,7 @@ from typing import TYPE_CHECKING, cast
from unittest.mock import AsyncMock, MagicMock, patch
import pytest
+from fastapi import Request
from roboco.db.base import (
_db_has_alembic_version,
_db_has_tables,
@@ -19,6 +20,7 @@ from roboco.db.base import (
close_db,
drop_db,
get_db,
+ get_db_committed,
get_db_context,
get_engine,
get_session_factory,
@@ -108,6 +110,23 @@ async def test_get_db_yields_session_and_commits_on_success() -> None:
fake_session.rollback.assert_not_called()
+@pytest.mark.asyncio
+async def test_get_db_committed_stashes_session_on_request_state() -> None:
+ """DbCommitMiddleware (api/middleware.py) reads request.state.db_session
+ to commit it before the response reaches the client — get_db_committed
+ (the roboco.api.deps.DbSession target) must stash the session there
+ before yielding it back unchanged."""
+ fake_session = MagicMock()
+ request = Request({"type": "http"})
+
+ gen = get_db_committed(request, fake_session)
+ yielded = await gen.__anext__()
+ assert yielded is fake_session
+ assert request.state.db_session is fake_session
+ with pytest.raises(StopAsyncIteration):
+ await gen.__anext__()
+
+
@pytest.mark.asyncio
async def test_get_db_rolls_back_on_exception() -> None:
fake_session = MagicMock()
diff --git a/tests/unit/agent_sdk/test_verb_circuit_breaker.py b/tests/unit/agent_sdk/test_verb_circuit_breaker.py
index 1a08150a..0e6d8203 100644
--- a/tests/unit/agent_sdk/test_verb_circuit_breaker.py
+++ b/tests/unit/agent_sdk/test_verb_circuit_breaker.py
@@ -20,7 +20,11 @@ from unittest.mock import patch
import pytest
import roboco.agent_sdk.server as srv
from fastapi.testclient import TestClient
-from roboco.foundation.policy.agent_loop import VERB_RETRY_LIMITS, retry_limit_for
+from roboco.foundation.policy.agent_loop import (
+ VERB_RETRY_LIMITS,
+ absolute_retry_limit_for,
+ retry_limit_for,
+)
from roboco.services.gateway.envelope import Envelope
if TYPE_CHECKING:
@@ -428,3 +432,168 @@ def test_qa_handoff_retry_keys_match_mcp_verb_names() -> None:
assert "fail" in VERB_RETRY_LIMITS
assert retry_limit_for("pass") == VERB_RETRY_LIMITS["pass"]
assert retry_limit_for("fail") == VERB_RETRY_LIMITS["fail"]
+
+
+# ---------------------------------------------------------------------------
+# ABSOLUTE (session-scoped, never-pruned) breaker — the slow-drip fix
+#
+# Production, 2026-07-08: an agent's i_am_done was rejected once every 3-4
+# minutes for 30+ minutes. Each rejection arrived alone in an empty 60s
+# window, so `_check_verb_circuit` never saw more than 1 attempt at a time
+# and never tripped. `_verb_absolute_attempts` counts cumulatively across
+# the whole container session (never pruned) so pacing can't defeat it.
+# ---------------------------------------------------------------------------
+
+_I_AM_DONE_ABSOLUTE_CAP = 9 # VERB_RETRY_LIMITS["i_am_done"] * multiplier(3)
+
+
+def test_i_am_done_absolute_cap_matches_foundation() -> None:
+ """Pin the effective absolute cap so a foundation change is caught here too."""
+ assert absolute_retry_limit_for("i_am_done") == _I_AM_DONE_ABSOLUTE_CAP
+
+
+def test_absolute_tracker_keys_per_verb_task_pair() -> None:
+ """Different tasks accumulate independent absolute counts."""
+ a_count = 5
+ b_count = 2
+ for _ in range(a_count):
+ srv._record_verb_attempt_absolute("i_am_done", "task-A")
+ for _ in range(b_count):
+ srv._record_verb_attempt_absolute("i_am_done", "task-B")
+
+ assert srv._verb_absolute_attempt_count("i_am_done", "task-A") == a_count
+ assert srv._verb_absolute_attempt_count("i_am_done", "task-B") == b_count
+
+
+def test_absolute_tracker_keys_per_verb_independent_of_task() -> None:
+ """Different verbs on the same task don't share the cumulative counter."""
+ done_count = 4
+ submit_count = 1
+ for _ in range(done_count):
+ srv._record_verb_attempt_absolute("i_am_done", "task-A")
+ for _ in range(submit_count):
+ srv._record_verb_attempt_absolute("submit_up", "task-A")
+
+ assert srv._verb_absolute_attempt_count("i_am_done", "task-A") == done_count
+ assert srv._verb_absolute_attempt_count("submit_up", "task-A") == submit_count
+
+
+def test_absolute_counter_never_prunes_with_time() -> None:
+ """Unlike the windowed deque, a huge time jump does not reset the count."""
+ base = 1000.0
+ expected = 2
+ with patch("roboco.agent_sdk.server.time.monotonic") as mock_time:
+ mock_time.return_value = base
+ srv._record_verb_attempt_absolute("i_am_done", "task-A")
+ mock_time.return_value = base + 10_000.0 # far past any sliding window
+ srv._record_verb_attempt_absolute("i_am_done", "task-A")
+ assert srv._verb_absolute_attempt_count("i_am_done", "task-A") == expected
+
+
+def test_slow_drip_never_trips_window_but_trips_absolute_cap() -> None:
+ """Rejections spaced > 60s apart never trip `_check_verb_circuit`, but
+ the absolute cap still trips once the cumulative count reaches it —
+ the exact production scenario this breaker was added to close.
+ """
+ cap = absolute_retry_limit_for("i_am_done")
+ assert cap is not None
+ base = 1000.0
+ with patch("roboco.agent_sdk.server.time.monotonic") as mock_time:
+ for i in range(cap):
+ mock_time.return_value = base + i * 200.0 # always > 60s apart
+ srv._record_verb_attempt("i_am_done", "task-A")
+ srv._record_verb_attempt_absolute("i_am_done", "task-A")
+ # The sliding window is always empty when this attempt lands.
+ assert srv._check_verb_circuit("i_am_done", "task-A") is None
+
+ result = srv._check_verb_absolute_circuit("i_am_done", "task-A")
+ assert result is not None
+ assert result["error"] == "circuit_open"
+ assert "absolute cap" in result["message"]
+ assert "i_am_blocked" in result["remediate"]
+
+
+def test_absolute_check_returns_none_below_cap() -> None:
+ """One rejection short of the cap, the absolute breaker stays closed."""
+ cap = absolute_retry_limit_for("i_am_done")
+ assert cap is not None
+ for _ in range(cap - 1):
+ srv._record_verb_attempt_absolute("i_am_done", "task-A")
+ assert srv._check_verb_absolute_circuit("i_am_done", "task-A") is None
+
+
+def test_absolute_check_returns_none_for_unlimited_retry_verbs() -> None:
+ """give_me_work stays exempt from the absolute cap too."""
+ assert absolute_retry_limit_for("give_me_work") is None
+ for _ in range(50):
+ srv._record_verb_attempt_absolute("give_me_work", None)
+ assert srv._check_verb_absolute_circuit("give_me_work", None) is None
+
+
+def test_combined_check_still_trips_fast_storm_via_window() -> None:
+ """Windowed behavior is unchanged: 3 fast rejections (the existing
+ i_am_done cap) still trip via the combined check, well below the
+ absolute cap of 9 — and the message reads as a windowed trip, not a
+ session one.
+ """
+ limit = retry_limit_for("i_am_done")
+ assert limit is not None
+ for _ in range(limit):
+ srv._record_verb_attempt("i_am_done", "task-A")
+ srv._record_verb_attempt_absolute("i_am_done", "task-A")
+
+ result = srv._check_any_verb_circuit("i_am_done", "task-A")
+ assert result is not None
+ assert result["error"] == "circuit_open"
+ assert "this session" not in result["message"]
+ assert "in last" in result["message"]
+
+
+def test_verb_attempted_endpoint_trips_absolute_cap_on_slow_drip() -> None:
+ """End-to-end repro through the real /verb/attempted endpoint: rejections
+ paced far past 60s apart never open the windowed breaker but do open the
+ absolute one once the cumulative count reaches the cap.
+ """
+ client = TestClient(srv.app)
+ cap = absolute_retry_limit_for("i_am_done")
+ assert cap is not None
+ last_body: dict[str, object] | None = None
+ with patch("roboco.agent_sdk.server.time.monotonic") as mock_time:
+ for i in range(cap):
+ mock_time.return_value = 1000.0 + i * 200.0
+ resp = client.post(
+ "/verb/attempted",
+ json={
+ "verb": "i_am_done",
+ "task_id": "task-A",
+ "rejection_kind": "tracing_gap",
+ },
+ )
+ assert resp.status_code == _OK
+ last_body = resp.json()
+ if i < cap - 1:
+ assert last_body["open"] is False
+
+ assert last_body is not None
+ assert last_body["open"] is True
+ env = last_body["circuit_envelope"]
+ assert isinstance(env, dict)
+ assert env["error"] == "circuit_open"
+ assert "absolute cap" in env["message"]
+
+
+def test_state_reset_clears_verb_absolute_attempts() -> None:
+ """_state.reset() (a fresh container spawn) wipes the absolute tracker too."""
+ expected = 3
+ for _ in range(expected):
+ srv._record_verb_attempt_absolute("i_am_done", "task-A")
+ assert srv._verb_absolute_attempt_count("i_am_done", "task-A") == expected
+
+ srv._state.reset()
+ assert srv._verb_absolute_attempt_count("i_am_done", "task-A") == 0
+
+
+def test_verb_absolute_attempts_default_is_zero() -> None:
+ """defaultdict yields 0 for unseen keys — sanity check."""
+ fresh = srv._SessionState()
+ assert fresh.verb_absolute_attempts[("never_seen", None)] == 0
diff --git a/tests/unit/api/test_middleware.py b/tests/unit/api/test_middleware.py
index cd31c4a9..823c709c 100644
--- a/tests/unit/api/test_middleware.py
+++ b/tests/unit/api/test_middleware.py
@@ -4,22 +4,25 @@ from __future__ import annotations
import asyncio
from http import HTTPStatus
-from typing import Any
+from typing import TYPE_CHECKING, Annotated, Any, cast
# UUID annotates a Pydantic model field below, so it must stay a runtime import
# (Pydantic resolves the annotation when building the model) despite `from
# __future__ import annotations` making it look type-checking-only to ruff.
from uuid import UUID # noqa: TC003
-from fastapi import FastAPI, HTTPException
+import httpx
+from fastapi import Depends, FastAPI, HTTPException, Request
from fastapi.testclient import TestClient
from pydantic import BaseModel, field_validator
from roboco.api.middleware import (
+ DbCommitMiddleware,
_uuid_field_remediation,
get_status_code,
setup_middleware,
)
from roboco.config import settings
+from roboco.db.base import _discard_on_cancel
from roboco.exceptions import (
AuthenticationError,
InvalidStateError,
@@ -42,6 +45,9 @@ from roboco.services.base import (
)
from structlog.testing import capture_logs
+if TYPE_CHECKING:
+ from sqlalchemy.ext.asyncio import AsyncSession
+
# ---------------------------------------------------------------------------
# get_status_code
# ---------------------------------------------------------------------------
@@ -469,3 +475,247 @@ def test_flow_verb_timeout_slow_verb_uses_slow_budget(monkeypatch: Any) -> None:
response = client.post("/api/v1/flow/developer/i_am_done")
assert response.status_code == HTTPStatus.OK
assert response.json() == {"status": "ok"}
+
+
+# ---------------------------------------------------------------------------
+# DbCommitMiddleware — commits the request's DB session before the response
+# reaches the client. Reproduces the race: FastAPI sends the response before
+# a Depends(get_db)-with-yield dependency's post-yield commit runs.
+# ---------------------------------------------------------------------------
+
+
+class _OrderedSession:
+ """get_db-style fake session for ordering assertions.
+
+ Exposes ``in_transaction()`` / ``commit()`` / ``rollback()`` like the real
+ ``AsyncSession`` the middleware drives, recording call order in a shared
+ list so a test can assert the commit happens before the wire send.
+ """
+
+ def __init__(self, order: list[str], fail_commit: bool = False) -> None:
+ self._order = order
+ self._fail_commit = fail_commit
+ self._txn = True
+
+ def in_transaction(self) -> bool:
+ return self._txn
+
+ async def commit(self) -> None:
+ self._order.append("commit")
+ if self._fail_commit:
+ raise RuntimeError("commit failed")
+ self._txn = False
+
+ async def rollback(self) -> None:
+ self._order.append("rollback")
+ self._txn = False
+
+
+async def _fake_get_db(request: Request) -> Any:
+ """Module-level get_db-style dependency: stash the session on
+ request.state, yield, commit post-yield as the fallback — the exact
+ shape ``roboco.db.base.get_db`` uses and ``DbCommitMiddleware`` targets.
+
+ Reads its order-list/fail-flag from ``request.app.state`` rather than a
+ closure: ``Annotated[Any, Depends(...)]`` is stringified by this file's
+ ``from __future__ import annotations``, and ``typing.get_type_hints``
+ only resolves names from the function's module globals — a local
+ closure name would raise, silently downgrading the parameter to a plain
+ query param instead of a dependency.
+ """
+ order: list[str] = request.app.state.db_commit_order
+ session = _OrderedSession(order, fail_commit=request.app.state.db_commit_fail)
+ request.state.db_session = session
+ try:
+ yield session
+ await session.commit()
+ except Exception:
+ await session.rollback()
+ raise
+
+
+def _make_db_commit_app(order: list[str], fail_commit: bool = False) -> FastAPI:
+ app = FastAPI()
+ app.state.db_commit_order = order
+ app.state.db_commit_fail = fail_commit
+
+ @app.post("/write")
+ async def _write(_db: Annotated[Any, Depends(_fake_get_db)]) -> Any:
+ order.append("route_body")
+ return {"ok": True}
+
+ setup_middleware(app)
+ return app
+
+
+def _instrumented_transport(app: FastAPI, order: list[str]) -> httpx.ASGITransport:
+ """Wraps ``app`` so 'wire_response_start' marks the instant bytes would
+ leave the server — the outermost observation point, past every
+ middleware including DbCommitMiddleware. ``raise_app_exceptions=False``
+ lets the failing-commit test inspect the resulting 5xx response instead
+ of the exception ServerErrorMiddleware always re-raises after sending it."""
+
+ async def outer(scope: Any, receive: Any, send: Any) -> None:
+ async def capture(message: Any) -> None:
+ if message["type"] == "http.response.start":
+ order.append("wire_response_start")
+ await send(message)
+
+ await app(scope, receive, capture)
+
+ return httpx.ASGITransport(app=outer, raise_app_exceptions=False)
+
+
+async def test_db_commit_middleware_commits_before_response_reaches_client() -> None:
+ """The client only sees the response after the session commits — proving
+ the middleware, not get_db's post-yield fallback (which FastAPI's own
+ routing runs AFTER the response is already on the wire), commits in time."""
+ order: list[str] = []
+ app = _make_db_commit_app(order)
+ transport = _instrumented_transport(app, order)
+ async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
+ response = await client.post("/write")
+ assert response.status_code == HTTPStatus.OK
+ assert "commit" in order
+ assert "wire_response_start" in order
+ assert order.index("commit") < order.index("wire_response_start"), order
+
+
+async def test_db_commit_middleware_failing_commit_returns_5xx_not_200() -> None:
+ """A commit that fails after the route succeeded must not report 200 —
+ the response hasn't reached the wire yet, so it comes back as a 5xx."""
+ order: list[str] = []
+ app = _make_db_commit_app(order, fail_commit=True)
+ transport = _instrumented_transport(app, order)
+ async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
+ response = await client.post("/write")
+ assert response.status_code >= HTTPStatus.INTERNAL_SERVER_ERROR
+ assert "commit" in order
+
+
+async def test_db_commit_middleware_skips_non_http_scope() -> None:
+ """Websocket (and any non-http) scope passes straight through untouched —
+ no send wrapping, no state lookup."""
+ calls: list[dict[str, Any]] = []
+
+ async def inner_app(scope: Any, _receive: Any, _send: Any) -> None:
+ calls.append(scope)
+
+ middleware = DbCommitMiddleware(inner_app)
+ scope = {"type": "websocket"}
+
+ async def receive() -> dict[str, Any]:
+ return {}
+
+ async def send(_message: Any) -> None:
+ raise AssertionError("send should not be called for a websocket scope")
+
+ await middleware(scope, receive, send)
+ assert calls == [scope]
+
+
+def test_db_commit_middleware_passes_through_session_less_request() -> None:
+ """A request whose route never installs a get_db-style dependency (no
+ request.state.db_session stashed) reaches the client unmodified."""
+ app = FastAPI()
+
+ @app.get("/plain")
+ async def _plain() -> Any:
+ return {"ok": True}
+
+ setup_middleware(app)
+ client = TestClient(app)
+ response = client.get("/plain")
+ assert response.status_code == HTTPStatus.OK
+ assert response.json() == {"ok": True}
+
+
+# ---------------------------------------------------------------------------
+# FlowVerbTimeoutMiddleware x DbCommitMiddleware — cancellation landing
+# mid-commit (2026-07-08 CI segfault: FlowVerbTimeoutMiddleware's
+# asyncio.timeout is scoped to the whole request, so it can fire while
+# DbCommitMiddleware's own `await session.commit()` is in flight, not just
+# while a route handler hangs before responding).
+# ---------------------------------------------------------------------------
+
+
+class _CancelableCommitSession:
+ """get_db-style fake session whose ``commit()`` blocks forever on an
+ Event — it only ever exits via cancellation, reproducing the exact race
+ where FlowVerbTimeoutMiddleware's timeout fires mid-``commit()``."""
+
+ def __init__(self, order: list[str]) -> None:
+ self._order = order
+ self._txn = True
+ self._never_set = asyncio.Event()
+
+ def in_transaction(self) -> bool:
+ return self._txn
+
+ async def commit(self) -> None:
+ self._order.append("commit_start")
+ await self._never_set.wait()
+ self._order.append("commit_end") # unreachable: proves no commit-after-cancel
+
+ async def rollback(self) -> None:
+ self._order.append("rollback")
+ self._txn = False
+
+ async def invalidate(self) -> None:
+ self._order.append("invalidate")
+ self._txn = False
+
+
+async def _fake_get_db_cancel_safe(request: Request) -> Any:
+ """get_db-style dependency wired to the real ``_discard_on_cancel``
+ helper (``roboco.db.base``) so this test exercises the actual fix, not a
+ re-implementation of it."""
+ order: list[str] = request.app.state.db_commit_order
+ session = _CancelableCommitSession(order)
+ request.state.db_session = session
+ try:
+ yield session
+ await session.commit()
+ except asyncio.CancelledError:
+ # The double only duck-types AsyncSession's commit/rollback/invalidate
+ # surface — cast for the real helper's signature.
+ await _discard_on_cancel(cast("AsyncSession", session))
+ raise
+ except Exception:
+ await session.rollback()
+ raise
+
+
+def _make_cancel_during_commit_app(order: list[str]) -> FastAPI:
+ app = FastAPI()
+ app.state.db_commit_order = order
+
+ @app.post("/api/v1/flow/developer/give_me_work")
+ async def _write(_db: Annotated[Any, Depends(_fake_get_db_cancel_safe)]) -> Any:
+ order.append("route_body")
+ return {"status": "ok"}
+
+ setup_middleware(app)
+ return app
+
+
+def test_cancellation_mid_commit_invalidates_not_rollback(monkeypatch: Any) -> None:
+ """A flow-verb request that blows its server-side timeout WHILE
+ DbCommitMiddleware's commit is in flight must: propagate CancelledError
+ cleanly to a 504 (not hang, not a raw 500), discard the session via
+ ``invalidate()`` — NOT ``rollback()`` (SQLAlchemy's own docs: rolling
+ back a cancelled/timed-out operation risks issuing another command over
+ a connection whose wire-protocol state is now undefined, which is what
+ let a later request's pool checkout crash on the poisoned connection) —
+ and never resume/complete the cancelled commit.
+ """
+ monkeypatch.setattr(settings, "flow_verb_timeout_seconds", 0.05)
+ order: list[str] = []
+ app = _make_cancel_during_commit_app(order)
+
+ client = TestClient(app)
+ response = client.post("/api/v1/flow/developer/give_me_work")
+
+ assert response.status_code == HTTPStatus.GATEWAY_TIMEOUT
+ assert response.json()["error"] == "gateway_timeout"
+ assert order == ["route_body", "commit_start", "invalidate"], order
diff --git a/tests/unit/api/test_orchestrator_manual_spawn.py b/tests/unit/api/test_orchestrator_manual_spawn.py
new file mode 100644
index 00000000..134a694e
--- /dev/null
+++ b/tests/unit/api/test_orchestrator_manual_spawn.py
@@ -0,0 +1,277 @@
+"""Manual (panel) spawn: task-aware prompt helper + already-running signaling.
+
+Covers the CEO-facing spawn-refusal / double-fire triage: a task-aware
+initial prompt built server-side for a manual spawn (mirroring
+``_build_pr_review_prompt``'s tone), an ``AgentReadinessError`` refusal
+mapped to 409 (not an opaque 500) so the panel can show the real reason, and
+an ``already_running`` marker so a no-op spawn (agent already active) is
+distinguishable from a genuine new spawn.
+"""
+
+from __future__ import annotations
+
+from datetime import UTC, datetime
+from http import HTTPStatus
+from types import SimpleNamespace
+from typing import TYPE_CHECKING, cast
+from unittest.mock import AsyncMock, MagicMock
+from uuid import uuid4
+
+import pytest
+import pytest_asyncio
+import roboco.api.routes.orchestrator as orch_route
+from fastapi import FastAPI
+from httpx import ASGITransport, AsyncClient
+from roboco.api.deps import _ServiceHolder, set_orchestrator
+from roboco.api.routes.orchestrator import (
+ _build_manual_spawn_prompt,
+ _resolve_manual_spawn_prompt,
+)
+from roboco.api.routes.orchestrator import (
+ router as orch_router,
+)
+from roboco.runtime.orchestrator import AgentReadinessError, AgentState
+
+if TYPE_CHECKING:
+ from collections.abc import AsyncIterator
+
+ from roboco.db.tables import TaskTable
+
+_HDR = {"X-Agent-ID": str(uuid4()), "X-Agent-Role": "ceo"}
+
+
+def _fake_task(status_value: str = "pending") -> TaskTable:
+ # SimpleNamespace duck-types TaskTable's 3 fields the helper reads
+ # (id/title/status.value) without a real ORM row.
+ return cast(
+ "TaskTable",
+ SimpleNamespace(
+ id="task-123",
+ title="Fix the thing",
+ status=SimpleNamespace(value=status_value),
+ ),
+ )
+
+
+class _FakeDbCtx:
+ async def __aenter__(self) -> str:
+ return "fake-db"
+
+ async def __aexit__(self, *exc: object) -> bool:
+ return False
+
+
+class _FakeTaskService:
+ def __init__(self, task: object | None = None, error: Exception | None = None):
+ self._task = task
+ self._error = error
+
+ async def get(self, _task_id: object) -> object | None:
+ if self._error:
+ raise self._error
+ return self._task
+
+
+# ---------------------------------------------------------------------------
+# _build_manual_spawn_prompt — pure formatting
+# ---------------------------------------------------------------------------
+
+
+def test_build_manual_spawn_prompt_includes_task_fields() -> None:
+ prompt = _build_manual_spawn_prompt(_fake_task("awaiting_qa"), None)
+ assert "TASK ID: task-123" in prompt
+ assert "TITLE: Fix the thing" in prompt
+ assert "STATUS: awaiting_qa" in prompt
+ assert "claim verb" in prompt.lower()
+ assert "CEO NOTE" not in prompt
+
+
+def test_build_manual_spawn_prompt_appends_ceo_note() -> None:
+ prompt = _build_manual_spawn_prompt(_fake_task(), "Please prioritize this.")
+ assert "== CEO NOTE ==" in prompt
+ assert "Please prioritize this." in prompt
+ # CEO note comes after the task framing, not instead of it.
+ assert prompt.index("TASK ID") < prompt.index("CEO NOTE")
+
+
+# ---------------------------------------------------------------------------
+# _resolve_manual_spawn_prompt — best-effort enrichment
+# ---------------------------------------------------------------------------
+
+
+@pytest.mark.asyncio
+async def test_resolve_prompt_no_task_id_returns_message_unchanged() -> None:
+ result = await _resolve_manual_spawn_prompt(None, "hello")
+ assert result == "hello"
+
+
+@pytest.mark.asyncio
+async def test_resolve_prompt_enriches_when_task_found(
+ monkeypatch: pytest.MonkeyPatch,
+) -> None:
+ monkeypatch.setattr(orch_route, "get_db_context", _FakeDbCtx)
+ monkeypatch.setattr(
+ orch_route,
+ "get_task_service",
+ lambda _db: _FakeTaskService(task=_fake_task("verifying")),
+ )
+ result = await _resolve_manual_spawn_prompt(str(uuid4()), "Ship it")
+ assert result is not None
+ assert "STATUS: verifying" in result
+ assert "Ship it" in result
+
+
+@pytest.mark.asyncio
+async def test_resolve_prompt_falls_back_when_task_not_found(
+ monkeypatch: pytest.MonkeyPatch,
+) -> None:
+ monkeypatch.setattr(orch_route, "get_db_context", _FakeDbCtx)
+ monkeypatch.setattr(
+ orch_route, "get_task_service", lambda _db: _FakeTaskService(task=None)
+ )
+ result = await _resolve_manual_spawn_prompt(str(uuid4()), "hello")
+ assert result == "hello"
+
+
+@pytest.mark.asyncio
+async def test_resolve_prompt_falls_back_on_bad_task_id() -> None:
+ # Not a valid UUID — must not raise, must fall back unchanged.
+ result = await _resolve_manual_spawn_prompt("not-a-uuid", "hello")
+ assert result == "hello"
+
+
+@pytest.mark.asyncio
+async def test_resolve_prompt_falls_back_on_db_error(
+ monkeypatch: pytest.MonkeyPatch,
+) -> None:
+ monkeypatch.setattr(orch_route, "get_db_context", _FakeDbCtx)
+ monkeypatch.setattr(
+ orch_route,
+ "get_task_service",
+ lambda _db: _FakeTaskService(error=RuntimeError("db down")),
+ )
+ result = await _resolve_manual_spawn_prompt(str(uuid4()), "hello")
+ assert result == "hello"
+
+
+@pytest.mark.asyncio
+async def test_resolve_prompt_no_message_no_task_returns_none() -> None:
+ result = await _resolve_manual_spawn_prompt(None, None)
+ assert result is None
+
+
+# ---------------------------------------------------------------------------
+# Route: AgentReadinessError -> 409, already_running signaling
+# ---------------------------------------------------------------------------
+
+
+@pytest_asyncio.fixture
+async def orch_client() -> AsyncIterator[tuple[AsyncClient, MagicMock]]:
+ app = FastAPI()
+ app.include_router(orch_router, prefix="/api/orchestrator")
+ orchestrator = MagicMock()
+ set_orchestrator(orchestrator)
+ transport = ASGITransport(app=app)
+ async with AsyncClient(transport=transport, base_url="http://test") as client:
+ yield client, orchestrator
+ _ServiceHolder.orchestrator = None
+ app.dependency_overrides.clear()
+
+
+@pytest.mark.asyncio
+async def test_spawn_readiness_refusal_maps_to_409(
+ orch_client: tuple[AsyncClient, MagicMock],
+) -> None:
+ client, orch = orch_client
+ orch.get_instance = MagicMock(return_value=None)
+ orch.spawn_agent = AsyncMock(
+ side_effect=AgentReadinessError(
+ "spawn refused for fe-dev-2 (task=t1): state=awaiting_qa requires "
+ "role in {'qa'} but agent fe-dev-2 is 'developer'"
+ )
+ )
+ response = await client.post(
+ "/api/orchestrator/agents/fe-dev-2/spawn",
+ json={"agent_id": "fe-dev-2", "task_id": "t1"},
+ headers=_HDR,
+ )
+ assert response.status_code == HTTPStatus.CONFLICT
+ assert "requires role in" in response.json()["detail"]
+
+
+@pytest.mark.asyncio
+async def test_spawn_new_agent_not_flagged_already_running(
+ orch_client: tuple[AsyncClient, MagicMock],
+) -> None:
+ client, orch = orch_client
+ orch.get_instance = MagicMock(return_value=None)
+ instance = SimpleNamespace(
+ id=uuid4(),
+ agent_id="be-dev-1",
+ state=AgentState.STARTING,
+ current_task_id=None,
+ error_count=0,
+ started_at=datetime.now(UTC),
+ )
+ orch.spawn_agent = AsyncMock(return_value=instance)
+ response = await client.post(
+ "/api/orchestrator/agents/be-dev-1/spawn", headers=_HDR
+ )
+ assert response.status_code == HTTPStatus.CREATED
+ assert response.json()["already_running"] is False
+
+
+@pytest.mark.asyncio
+async def test_spawn_already_running_agent_is_flagged(
+ orch_client: tuple[AsyncClient, MagicMock],
+) -> None:
+ client, orch = orch_client
+ shared_id = uuid4()
+ existing = SimpleNamespace(
+ id=shared_id,
+ agent_id="ux-pm",
+ state=AgentState.STARTING,
+ current_task_id=None,
+ error_count=0,
+ started_at=datetime.now(UTC),
+ )
+ orch.get_instance = MagicMock(return_value=existing)
+ # spawn_agent's own no-op contract: hands back the SAME instance.
+ orch.spawn_agent = AsyncMock(return_value=existing)
+ response = await client.post("/api/orchestrator/agents/ux-pm/spawn", headers=_HDR)
+ assert response.status_code == HTTPStatus.CREATED
+ body = response.json()
+ assert body["already_running"] is True
+ assert body["state"] == "starting"
+
+
+@pytest.mark.asyncio
+async def test_spawn_offline_agent_not_flagged_already_running(
+ orch_client: tuple[AsyncClient, MagicMock],
+) -> None:
+ """A pre-existing OFFLINE instance is not "running" — a fresh spawn on top
+ of it must not be reported as a no-op."""
+ client, orch = orch_client
+ offline = SimpleNamespace(
+ id=uuid4(),
+ agent_id="be-dev-1",
+ state=AgentState.OFFLINE,
+ current_task_id=None,
+ error_count=0,
+ started_at=datetime.now(UTC),
+ )
+ orch.get_instance = MagicMock(return_value=offline)
+ new_instance = SimpleNamespace(
+ id=uuid4(),
+ agent_id="be-dev-1",
+ state=AgentState.STARTING,
+ current_task_id=None,
+ error_count=0,
+ started_at=datetime.now(UTC),
+ )
+ orch.spawn_agent = AsyncMock(return_value=new_instance)
+ response = await client.post(
+ "/api/orchestrator/agents/be-dev-1/spawn", headers=_HDR
+ )
+ assert response.status_code == HTTPStatus.CREATED
+ assert response.json()["already_running"] is False
diff --git a/tests/unit/gateway/test_audit_on_rejection.py b/tests/unit/gateway/test_audit_on_rejection.py
index 0fe0be3d..d4518f2e 100644
--- a/tests/unit/gateway/test_audit_on_rejection.py
+++ b/tests/unit/gateway/test_audit_on_rejection.py
@@ -26,6 +26,7 @@ from uuid import uuid4
import pytest
from roboco.services.gateway.choreographer import Choreographer, ChoreographerDeps
+from roboco.services.gateway.envelope import Envelope
def _make_deps(**overrides: Any) -> ChoreographerDeps:
@@ -176,3 +177,63 @@ async def test_audit_log_event_failure_does_not_propagate() -> None:
assert env.error == "not_found"
audit_svc.log_event.assert_awaited()
+
+
+# ---------------------------------------------------------------------------
+# remediate must ride along into the audit row — it's the only place a
+# conventions-gate rejection's file:line violation listing lives.
+# ---------------------------------------------------------------------------
+
+
+@pytest.mark.asyncio
+async def test_rejection_remediate_lands_in_audit_details() -> None:
+ """A rejection's `remediate` hint is copied into the audit row's details.
+
+ Without this, an operator reading `gateway.rejected` audit rows for a
+ conventions-gate rejection sees only the summary message — the
+ actionable detail lives solely in `remediate`.
+ """
+ aid = uuid4()
+ tid = uuid4()
+ code_task = MagicMock(
+ id=tid,
+ status="pending",
+ assigned_to=aid,
+ task_type="code",
+ priority=1,
+ parent_task_id=None,
+ sequence=0,
+ team="backend",
+ )
+ task_svc = AsyncMock()
+ task_svc.get.return_value = code_task
+ task_svc.agent_for.return_value = MagicMock(role="cell_pm", team="backend")
+ task_svc.list_in_progress_for_agent.return_value = []
+ task_svc.list_paused_for_agent.return_value = []
+ task_svc.get_subtasks.return_value = []
+ audit_svc = AsyncMock()
+ deps = _make_deps(task=task_svc, audit=audit_svc)
+ c = Choreographer(deps)
+
+ env = await c.i_will_work_on(aid, tid, plan="x")
+
+ assert env.error == "not_authorized"
+ assert env.remediate
+ args = audit_svc.log_event.await_args
+ assert args.kwargs["details"]["remediate"] == env.remediate
+
+
+@pytest.mark.asyncio
+async def test_rejection_without_remediate_omits_audit_key() -> None:
+ """A rejection with no remediate must not add a null key to the row."""
+ aid = uuid4()
+ tid = uuid4()
+ audit_svc = AsyncMock()
+ deps = _make_deps(audit=audit_svc)
+ c = Choreographer(deps)
+ env = Envelope(error="not_found", message="bare rejection, no remediate")
+
+ await c._emit_rejection(env, agent_id=aid, task_id=tid, verb="test_verb")
+
+ args = audit_svc.log_event.await_args
+ assert "remediate" not in args.kwargs["details"]
diff --git a/tests/unit/gateway/test_content_actions.py b/tests/unit/gateway/test_content_actions.py
index c992a906..47cc6b70 100644
--- a/tests/unit/gateway/test_content_actions.py
+++ b/tests/unit/gateway/test_content_actions.py
@@ -7,6 +7,7 @@ from uuid import uuid4
import pytest
from roboco.config import settings
+from roboco.exceptions import GitCommandError
from roboco.services.gateway.content_actions import ContentActions, ContentActionsDeps
@@ -162,6 +163,75 @@ async def test_commit_strips_existing_task_prefix() -> None:
assert "[ABC12345]" not in call_kwargs["message"]
+@pytest.mark.asyncio
+async def test_commit_no_changes_added_returns_files_remediate() -> None:
+ """A `files` list matching no modified paths → actionable remediate.
+
+ Regression: git.commit raising "no changes added to commit" (an agent
+ passed `files` that no-op'd the `git add`) used to propagate as a raw
+ GitCommandError instead of an envelope.
+ """
+ agent_id = uuid4()
+ task_id = uuid4()
+ task_obj = MagicMock(
+ id=task_id,
+ status="in_progress",
+ branch_name="feature/backend/abc",
+ active_claimant_id=agent_id,
+ )
+ task_svc = AsyncMock()
+ task_svc.get_active_task_for_agent.return_value = task_obj
+ task_svc.agent_for.return_value = MagicMock(role="developer")
+ git_svc = AsyncMock()
+ git_svc.commit.side_effect = GitCommandError("commit", "no changes added to commit")
+
+ deps = _make_deps(task=task_svc, git=git_svc)
+ ca = ContentActions(deps)
+
+ env = await ca.commit(
+ agent_id=agent_id,
+ message="feat(api): add /healthz endpoint for liveness checks",
+ files=["nonexistent.py"],
+ )
+ body = env.as_dict()
+
+ assert body["error"] == "invalid_state"
+ assert "nonexistent.py" in body["remediate"]
+ assert "omit files" in body["remediate"]
+ task_svc.add_progress.assert_not_awaited()
+
+
+@pytest.mark.asyncio
+async def test_commit_generic_git_failure_returns_envelope_not_exception() -> None:
+ """A generic git failure is caught and returned as an envelope, not raised."""
+ agent_id = uuid4()
+ task_id = uuid4()
+ task_obj = MagicMock(
+ id=task_id,
+ status="in_progress",
+ branch_name="feature/backend/abc",
+ active_claimant_id=agent_id,
+ )
+ task_svc = AsyncMock()
+ task_svc.get_active_task_for_agent.return_value = task_obj
+ task_svc.agent_for.return_value = MagicMock(role="developer")
+ git_svc = AsyncMock()
+ git_svc.commit.side_effect = GitCommandError("commit", "fatal: some other failure")
+
+ deps = _make_deps(task=task_svc, git=git_svc)
+ ca = ContentActions(deps)
+
+ env = await ca.commit(
+ agent_id=agent_id,
+ message="feat(api): add /healthz endpoint for liveness checks",
+ )
+ body = env.as_dict()
+
+ assert body["error"] == "invalid_state"
+ assert "inspect" in body["remediate"]
+ task_svc.add_progress.assert_not_awaited()
+
+
# ---------------------------------------------------------------------------
# note
# ---------------------------------------------------------------------------
diff --git a/tests/unit/runtime/test_mcp_cli_arg_uuid_auth.py b/tests/unit/runtime/test_mcp_cli_arg_uuid_auth.py
new file mode 100644
index 00000000..3545785f
--- /dev/null
+++ b/tests/unit/runtime/test_mcp_cli_arg_uuid_auth.py
@@ -0,0 +1,87 @@
+"""roboco-optimal/docs/search receive the agent's CLI arg as sys.argv[1] and
+forward it verbatim as X-Agent-ID via ApiClient/_get_agent_headers. The spawn
+token (_append_agent_auth_env) is signed over the agent's UUID, so that CLI
+arg must be the UUID too, or verify_agent_token 401s with a signature
+mismatch even though role/team resolve fine either way (get_agent_role/
+get_agent_team accept slug or UUID via _resolve_to_slug).
+"""
+
+from __future__ import annotations
+
+import json
+from pathlib import Path
+from typing import TYPE_CHECKING
+
+from roboco.agents_config import AGENT_UUIDS, verify_agent_token
+from roboco.config import settings
+from roboco.mcp import utils as mcp_utils
+from roboco.models.runtime import OrchestratorAgentConfig as AgentConfig
+from roboco.runtime.orchestrator import AgentOrchestrator
+
+if TYPE_CHECKING:
+ import pytest
+
+# main-pm carries roboco-optimal (always), roboco-docs (docs_roles) and
+# roboco-search (research_roles, research_enabled defaults True) all at once.
+_AGENT_SLUG = "main-pm"
+_CLI_ARG_SERVERS = ("roboco-optimal", "roboco-docs", "roboco-search")
+
+
+def _spawn_token(monkeypatch: pytest.MonkeyPatch) -> str:
+ """Mint the token exactly as _append_agent_auth_env does at spawn."""
+ monkeypatch.setenv("ROBOCO_AGENT_AUTH_SECRET", "spawn-secret")
+ monkeypatch.setattr(settings, "agent_token_ttl_seconds", 3600)
+ cmd: list[str] = []
+ config = AgentConfig(
+ agent_id=_AGENT_SLUG,
+ blueprint_path=Path("/app/blueprints/main-pm.md"),
+ provider_type="anthropic",
+ )
+ AgentOrchestrator._append_agent_auth_env(cmd, config)
+ for i, flag in enumerate(cmd):
+ if flag == "-e" and cmd[i + 1].startswith("ROBOCO_AGENT_TOKEN="):
+ return cmd[i + 1].split("=", 1)[1]
+ raise AssertionError("ROBOCO_AGENT_TOKEN not found in cmd")
+
+
+async def test_cli_arg_servers_get_uuid_not_slug() -> None:
+ """_generate_mcp_config passes the UUID (not the slug) as sys.argv[1]
+ to the three servers that identify their agent via CLI arg."""
+ orch = AgentOrchestrator.__new__(AgentOrchestrator)
+ config_path = await orch._generate_mcp_config(_AGENT_SLUG)
+ config = json.loads(Path(config_path).read_text())
+ servers = config["mcpServers"]
+ expected_uuid = AGENT_UUIDS[_AGENT_SLUG]
+ for name in _CLI_ARG_SERVERS:
+ assert name in servers, f"{name} should be mounted for {_AGENT_SLUG}"
+ cli_arg = servers[name]["args"][-1]
+ assert cli_arg == expected_uuid, (
+ f"{name} sys.argv[1] is {cli_arg!r}, expected the UUID "
+ f"{expected_uuid!r} — a slug here mismatches the UUID-signed "
+ f"spawn token and every call 401s."
+ )
+
+
+async def test_cli_arg_servers_headers_verify_against_spawn_token(
+ monkeypatch: pytest.MonkeyPatch,
+) -> None:
+ """The header tuple _get_agent_headers builds from the CLI arg
+ _generate_mcp_config hands these servers must verify against the token
+ the orchestrator actually injects into the container env."""
+ token = _spawn_token(monkeypatch)
+ monkeypatch.setenv("ROBOCO_AGENT_TOKEN", token)
+
+ orch = AgentOrchestrator.__new__(AgentOrchestrator)
+ config_path = await orch._generate_mcp_config(_AGENT_SLUG)
+ config = json.loads(Path(config_path).read_text())
+ servers = config["mcpServers"]
+
+ for name in _CLI_ARG_SERVERS:
+ cli_arg = servers[name]["args"][-1]
+ headers = mcp_utils._get_agent_headers(cli_arg)
+ assert verify_agent_token(
+ headers["X-Agent-Token"],
+ headers["X-Agent-ID"],
+ headers["X-Agent-Role"],
+ headers.get("X-Agent-Team", ""),
+ ), f"{name}'s header tuple ({headers}) failed verify_agent_token"