Skip to content

Commit 84b3d4e

Browse files
committed
fix(review): address coderabbit findings on dispatch collapse PR
- Strict TriageResponseSchema rejects legacy {mode,complexity} fields - Scaler rollbackSpawn() clears reservation only if it still matches - Router uses rollbackSpawn on K8s failure to avoid trampling concurrent reservations - Migration 004 UPDATEs now guard with IS DISTINCT FROM to be idempotent - job-queue getQueueLength() moves valkey client acquisition inside try/catch - Docs: fix fallback-reason count, parameterize RBAC namespaces, scope daemon-secrets, correct architecture scale-up diagram edges - Tests: align router mock contract with EphemeralSpawnErrorKind taxonomy, add daemon-mode DAEMON_AUTH_TOKEN coverage, cover legacy triage shape rejection
1 parent afbe7ff commit 84b3d4e

11 files changed

Lines changed: 117 additions & 38 deletions

File tree

‎docs/ARCHITECTURE.md‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ flowchart TD
1212
ROUTE["Router<br/>idempotency + allowlist + concurrency"]:::guard
1313
TR["Haiku triage<br/>binary heavy classifier"]:::decide
1414
QUEUE["Orchestrator job queue<br/>Valkey list"]:::store
15-
SCALE{{"Scale-up decision<br/>heavy OR queue ≥ threshold?<br/>and cooldown elapsed?"}}:::fork
15+
SCALE{{"Scale-up decision<br/>heavy OR queue ≥ threshold<br/>AND no persistent slots?<br/>and cooldown elapsed?"}}:::fork
1616
SPAWN["K8s API<br/>create bare Pod<br/>DAEMON_EPHEMERAL=true"]:::decide
1717
FLEET["Daemon fleet<br/>persistent + ephemeral<br/>WebSocket connections"]:::target
1818
PIPE["runPipeline<br/>clone → prompt → Claude Agent SDK"]:::work
@@ -22,7 +22,7 @@ flowchart TD
2222
ACK -. async .-> ROUTE
2323
ROUTE --> TR
2424
TR --> QUEUE
25-
ROUTE --> SCALE
25+
QUEUE --> SCALE
2626
SCALE -->|yes| SPAWN
2727
SPAWN --> FLEET
2828
SCALE -->|no, or cooldown active| FLEET

‎docs/CONFIGURATION.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -104,7 +104,7 @@ Required whenever the orchestrator role is active (i.e. the webhook server proce
104104
| `TRIAGE_TIMEOUT_MS` | `5000` | Per-call wall clock. Beyond this, the circuit-breaker counter increments. |
105105
| `DEFAULT_MAXTURNS` | `30` | Agent turn cap. Applied on every execution — triage no longer influences `maxTurns`. |
106106

107-
See [Triage](TRIAGE.md) for the binary `heavy` signal, circuit breaker, and the five fallback reasons that appear in logs.
107+
See [Triage](TRIAGE.md) for the binary `heavy` signal, circuit breaker, and the six fallback reasons that appear in logs.
108108

109109
## Mode matrix — what's required when
110110

‎docs/DEPLOYMENT.md‎

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -314,7 +314,8 @@ apiVersion: rbac.authorization.k8s.io/v1
314314
kind: Role
315315
metadata:
316316
name: github-app-playground-ephemeral-spawner
317-
namespace: default
317+
# Role lives in the namespace where ephemeral daemon Pods will be created.
318+
namespace: ${EPHEMERAL_DAEMON_NAMESPACE}
318319
rules:
319320
- apiGroups: [""]
320321
resources: ["pods"]
@@ -324,11 +325,13 @@ apiVersion: rbac.authorization.k8s.io/v1
324325
kind: RoleBinding
325326
metadata:
326327
name: github-app-playground-ephemeral-spawner
327-
namespace: default
328+
# Must match the Role namespace above.
329+
namespace: ${EPHEMERAL_DAEMON_NAMESPACE}
328330
subjects:
329331
- kind: ServiceAccount
330332
name: github-app-playground
331-
namespace: default
333+
# Namespace where the orchestrator ServiceAccount actually lives.
334+
namespace: ${ORCHESTRATOR_NAMESPACE}
332335
roleRef:
333336
kind: Role
334337
name: github-app-playground-ephemeral-spawner
@@ -339,4 +342,4 @@ Without these verbs, every scale-up attempt yields `dispatch_reason=ephemeral-sp
339342

340343
### `daemon-secrets` Secret
341344

342-
Spawned ephemeral daemon Pods receive their configuration via `envFrom: secretRef: daemon-secrets`. Create this Secret once in `EPHEMERAL_DAEMON_NAMESPACE` with the GitHub App credentials used by the orchestrator (so the daemon can reuse the installation-token path), Claude provider keys, and the data-layer URLs (`DATABASE_URL`, `VALKEY_URL`). See [DAEMON.md](DAEMON.md) for the full key list.
345+
Spawned ephemeral daemon Pods receive their configuration via `envFrom: secretRef: daemon-secrets`. Create this Secret once in `EPHEMERAL_DAEMON_NAMESPACE` with only the daemon runtime values it needs — `DAEMON_AUTH_TOKEN`, Claude provider keys, and the daemon-side data-layer URLs (`DATABASE_URL`, `VALKEY_URL`). Do **not** copy GitHub App private-key material into this Secret: the orchestrator mints installation tokens and hands them to the daemon per job, so expanding the blast radius to every ephemeral Pod is unnecessary. See [DAEMON.md](DAEMON.md) for the full key list.

‎src/db/migrations/004_collapse_dispatch_to_daemon.sql‎

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,10 +19,16 @@ ALTER TABLE executions DROP CONSTRAINT IF EXISTS executions_dispatch_reason_chec
1919
ALTER TABLE executions DROP CONSTRAINT IF EXISTS executions_dispatch_mode_check;
2020
ALTER TABLE triage_results DROP CONSTRAINT IF EXISTS triage_results_mode_check;
2121

22+
-- Mirror the open-ended guard used for executions.dispatch_reason below:
23+
-- rewrite any row whose target/mode is not already 'daemon', regardless of
24+
-- whether it comes from the documented legacy set. Without this, a stray
25+
-- value from an operator backfill or an older-branch deploy would survive
26+
-- the UPDATE and fail the CHECK (= 'daemon') added later in this file,
27+
-- aborting the whole transaction.
2228
UPDATE executions SET dispatch_target = 'daemon'
23-
WHERE dispatch_target IN ('inline', 'shared-runner', 'isolated-job');
29+
WHERE dispatch_target IS DISTINCT FROM 'daemon';
2430
UPDATE executions SET dispatch_mode = 'daemon'
25-
WHERE dispatch_mode IN ('inline', 'shared-runner', 'isolated-job');
31+
WHERE dispatch_mode IS DISTINCT FROM 'daemon';
2632

2733
-- Rewrite any reason not already in the post-collapse set. Keeping this
2834
-- unconditional (rather than enumerating legacy values) defends against

‎src/orchestrator/ephemeral-daemon-scaler.ts‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,19 @@ export function markSpawn(now: number): void {
6868
lastSpawnAtMs = now;
6969
}
7070

71+
/**
72+
* Release a cooldown reservation *only* if it still matches the caller's
73+
* attempt timestamp. If a newer reservation has already overtaken it (e.g.
74+
* a later spawn won the race while this one's K8s call was still in
75+
* flight), leave the newer timestamp in place. An unconditional reset
76+
* would reopen the thundering-herd window the cooldown is meant to close.
77+
*/
78+
export function rollbackSpawn(expectedTimestampMs: number): void {
79+
if (lastSpawnAtMs === expectedTimestampMs) {
80+
lastSpawnAtMs = 0;
81+
}
82+
}
83+
7184
/** Test-only: reset the module-level cooldown state between test cases. */
7285
export function _resetEphemeralScalerForTests(): void {
7386
lastSpawnAtMs = 0;

‎src/orchestrator/job-queue.ts‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -41,8 +41,11 @@ export async function enqueueJob(job: QueuedJob): Promise<void> {
4141
* persistent-daemon routing instead of trapping a job.
4242
*/
4343
export async function getQueueLength(): Promise<number> {
44-
const valkey = requireValkeyClient();
4544
try {
45+
// `requireValkeyClient()` itself throws when the client isn't ready, so
46+
// acquire inside the try block — otherwise the "return 0 on any read
47+
// error" contract silently loses to the client acquisition throw.
48+
const valkey = requireValkeyClient();
4649
// eslint-disable-next-line @typescript-eslint/no-unsafe-assignment -- Valkey LLEN returns number
4750
const len: number = await valkey.send("LLEN", [QUEUE_KEY]);
4851
return typeof len === "number" && Number.isFinite(len) ? len : 0;

‎src/orchestrator/triage.ts‎

Lines changed: 16 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -32,12 +32,22 @@ import { logger } from "../logger";
3232
import { CircuitBreaker } from "../utils/circuit-breaker";
3333
import { sanitizeContent } from "../utils/sanitize";
3434

35-
/** Zod schema — binary heavy classifier. */
36-
export const TriageResponseSchema = z.object({
37-
heavy: z.boolean(),
38-
confidence: z.number().min(0).max(1),
39-
rationale: z.string().min(1).max(500),
40-
});
35+
/**
36+
* Zod schema — binary heavy classifier.
37+
*
38+
* `.strict()` so legacy pre-collapse fields (`mode`, `complexity`) on a
39+
* response from a stale provider/proxy are a hard parse error instead of
40+
* being silently stripped. That forces the `parse-error` fallback branch
41+
* in `triageRequest`, which is safer than silently absorbing a shape we
42+
* no longer model.
43+
*/
44+
export const TriageResponseSchema = z
45+
.object({
46+
heavy: z.boolean(),
47+
confidence: z.number().min(0).max(1),
48+
rationale: z.string().min(1).max(500),
49+
})
50+
.strict();
4151
export type TriageResponse = z.infer<typeof TriageResponseSchema>;
4252

4353
export interface TriageResult extends TriageResponse {

‎src/webhook/router.ts‎

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,11 @@ import {
99
isAtCapacity,
1010
} from "../orchestrator/concurrency";
1111
import { getPersistentPoolFreeSlots } from "../orchestrator/daemon-registry";
12-
import { decideEphemeralSpawn, markSpawn } from "../orchestrator/ephemeral-daemon-scaler";
12+
import {
13+
decideEphemeralSpawn,
14+
markSpawn,
15+
rollbackSpawn,
16+
} from "../orchestrator/ephemeral-daemon-scaler";
1317
import { createExecution } from "../orchestrator/history";
1418
import { dispatchJob } from "../orchestrator/job-dispatcher";
1519
import { enqueueJob, getQueueLength, type QueuedJob } from "../orchestrator/job-queue";
@@ -240,8 +244,11 @@ export async function decideDispatch(ctx: BotContext): Promise<DispatchDecision>
240244
...(triageAttempted && { triageAttempted: true }),
241245
};
242246
} catch (err) {
243-
// Spawn failed — release the cooldown so the next request can retry.
244-
markSpawn(0);
247+
// Spawn failed — release only *our* reservation. Unconditionally
248+
// zeroing the timestamp would stomp on a newer concurrent spawn
249+
// that already won the cooldown race while this call was in flight,
250+
// reopening the thundering-herd window the cooldown exists to close.
251+
rollbackSpawn(spawnAttemptAt);
245252
const kind = err instanceof EphemeralSpawnError ? err.kind : undefined;
246253
const message =
247254
err instanceof EphemeralSpawnError

‎test/config.test.ts‎

Lines changed: 19 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -109,6 +109,19 @@ describe("configSchema — data layer validation", () => {
109109
});
110110
expect(result.success).toBe(true);
111111
});
112+
113+
it("still requires DAEMON_AUTH_TOKEN in daemon mode (ORCHESTRATOR_URL set)", () => {
114+
// The DB/Valkey waiver must not cascade into waiving daemon auth —
115+
// an orchestrator-connected daemon without a shared token would
116+
// accept unauthenticated connections on restart.
117+
const { daemonAuthToken, ...withoutToken } = ANTHROPIC_BASE;
118+
expect(daemonAuthToken).toBeDefined();
119+
const result = configSchema.safeParse({
120+
...withoutToken,
121+
orchestratorUrl: "wss://orchestrator.example.com",
122+
});
123+
expect(result.success).toBe(false);
124+
});
112125
});
113126

114127
describe("configSchema — ephemeral-daemon defaults", () => {
@@ -194,7 +207,12 @@ describe("assertOauthRequiresAllowlist", () => {
194207
});
195208

196209
it("throws when OAuth is set without an allowlist", () => {
197-
const cfg = { ...baseOauthCfg, allowedOwners: undefined };
210+
// With `exactOptionalPropertyTypes`, the absence of a property is
211+
// distinct from an explicit `undefined`. Destructure the property
212+
// out so this test actually models "no allowlist configured" — not
213+
// "allowlist is undefined-valued".
214+
const { allowedOwners, ...cfg } = baseOauthCfg;
215+
expect(allowedOwners).toBeDefined();
198216
expect(() => {
199217
assertOauthRequiresAllowlist(cfg);
200218
}).toThrow(/ALLOWED_OWNERS/);

‎test/contract/triage-response.test.ts‎

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -56,13 +56,25 @@ describe("TriageResponse schema — invalid inputs", () => {
5656
});
5757

5858
it("rejects legacy pre-collapse shape (mode/complexity)", () => {
59-
const r = TriageResponseSchema.safeParse({
59+
// Schema is `.strict()` — missing required `heavy` *or* extra
60+
// legacy keys is enough to fail. Cover both so future contributors
61+
// know the unknown-key rejection is intentional.
62+
const missingHeavy = TriageResponseSchema.safeParse({
6063
mode: "daemon",
6164
confidence: 0.5,
6265
complexity: "trivial",
6366
rationale: "legacy shape",
6467
});
65-
expect(r.success).toBe(false);
68+
expect(missingHeavy.success).toBe(false);
69+
70+
const legacyKeysPresent = TriageResponseSchema.safeParse({
71+
heavy: true,
72+
confidence: 0.5,
73+
rationale: "legacy fields present",
74+
mode: "daemon",
75+
complexity: "trivial",
76+
});
77+
expect(legacyKeysPresent.success).toBe(false);
6678
});
6779

6880
it("rejects confidence outside [0, 1]", () => {

0 commit comments

Comments
 (0)