Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
234 changes: 234 additions & 0 deletions .genie/wishes/genie-serve-stability/WISH.md

Large diffs are not rendered by default.

8 changes: 6 additions & 2 deletions src/__tests__/migrate.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import {
readFileSync,
readdirSync,
readlinkSync,
realpathSync,
rmSync,
symlinkSync,
writeFileSync,
Expand All @@ -27,8 +28,11 @@ import {
let testDir: string;

beforeEach(() => {
testDir = join(tmpdir(), `genie-migrate-test-${Date.now()}-${Math.random().toString(36).slice(2)}`);
mkdirSync(testDir, { recursive: true });
const raw = join(tmpdir(), `genie-migrate-test-${Date.now()}-${Math.random().toString(36).slice(2)}`);
mkdirSync(raw, { recursive: true });
// macOS: tmpdir() lives under /var → /private/var symlink; realpath so
// computed paths match what migrate.ts returns after resolving links.
testDir = realpathSync(raw);
});

afterEach(() => {
Expand Down
94 changes: 94 additions & 0 deletions src/lib/agent-registry.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -516,6 +516,100 @@ describe.skipIf(!DB_AVAILABLE)('pg', () => {
const a = await get('sdk-worker');
expect(a!.state).toBe('working');
});

test('flips workers registered on a dead socket to error (Bug 4)', async () => {
// Bug 4 repro: workers recorded on a tmux socket that no longer exists
// were stuck in 'idle'/'working' forever because reconcileStaleSpawns
// catches the TmuxUnreachableError from isPaneAlive and skips the worker.
// After the fix, a dead socket is detected once up-front and every
// worker on it is transitioned to 'error' with pane_id cleared.
//
// We pick a socket name that we can guarantee does NOT exist under
// /tmp/tmux-<uid>/ — any random UUID works.
const deadSocketName = `genie-dead-${Date.now()}-${Math.floor(Math.random() * 1e9)}`;
process.env.GENIE_TMUX_SOCKET = deadSocketName;
try {
const oldChange = new Date(Date.now() - 5_000).toISOString();
await register(
makeAgent({
id: 'dead-socket-worker-1',
paneId: '%9999',
state: 'idle',
startedAt: oldChange,
lastStateChange: oldChange,
}),
);
await register(
makeAgent({
id: 'dead-socket-worker-2',
paneId: '%9998',
state: 'working',
startedAt: oldChange,
lastStateChange: oldChange,
}),
);

const reset = await reconcileStaleSpawns(2);
expect(reset).toContain('dead-socket-worker-1');
expect(reset).toContain('dead-socket-worker-2');

const a1 = await get('dead-socket-worker-1');
expect(a1!.state).toBe('error');
expect(a1!.paneId).toBe('');
const a2 = await get('dead-socket-worker-2');
expect(a2!.state).toBe('error');
expect(a2!.paneId).toBe('');
} finally {
// biome-ignore lint/performance/noDelete: assigning undefined would set the string "undefined"
delete process.env.GENIE_TMUX_SOCKET;
}
Comment on lines +520 to +565

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟡 Minor

Restore the original GENIE_TMUX_SOCKET.

These tests overwrite a process-global env var and then always delete it. If the runner already had GENIE_TMUX_SOCKET set, later tests lose that value. Snapshot and restore it in both tests.

Proposed cleanup pattern
       const deadSocketName = `genie-dead-${Date.now()}-${Math.floor(Math.random() * 1e9)}`;
+      const originalSocketName = process.env.GENIE_TMUX_SOCKET;
       process.env.GENIE_TMUX_SOCKET = deadSocketName;
       try {
         const oldChange = new Date(Date.now() - 5_000).toISOString();
@@
       } finally {
-        // biome-ignore lint/performance/noDelete: assigning undefined would set the string "undefined"
-        delete process.env.GENIE_TMUX_SOCKET;
+        if (originalSocketName === undefined) {
+          // biome-ignore lint/performance/noDelete: assigning undefined would set the string "undefined"
+          delete process.env.GENIE_TMUX_SOCKET;
+        } else {
+          process.env.GENIE_TMUX_SOCKET = originalSocketName;
+        }
       }

Apply the same restore pattern around the live-socket test.

Also applies to: 493-512

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/lib/agent-registry.test.ts` around lines 430 - 475, The test overwrites
process.env.GENIE_TMUX_SOCKET and blindly deletes it in finally, which clobbers
any preexisting value; update the dead-socket-worker-1/2 test (and the other
live-socket test around the reconcileStaleSpawns tests) to snapshot const prev =
process.env.GENIE_TMUX_SOCKET before assigning a test socket, and in the finally
restore it with process.env.GENIE_TMUX_SOCKET = prev (or delete if prev is
undefined) so existing env is preserved; touch the tests referencing
reconcileStaleSpawns and the GENIE_TMUX_SOCKET manipulation to follow this
pattern.

});

test('dead socket reconciliation preserves existing behavior on live sockets', async () => {
// Sanity check: when the tmux socket DOES exist, the dead-socket
// fast path must not fire — workers must go through the normal
// per-pane isPaneAlive check (which remains transient-blip-safe).
// We force the socket to appear alive by pointing GENIE_TMUX_SOCKET
// at a socket file we create in /tmp/tmux-<uid>/.
const { existsSync, mkdirSync, writeFileSync, unlinkSync } = await import('node:fs');
const { join } = await import('node:path');
const { tmpdir } = await import('node:os');
const uid = process.getuid?.() ?? 501;
const socketDir = `/tmp/tmux-${uid}`;
const stubSocketName = `genie-live-stub-${Date.now()}`;
const stubPath = join(socketDir, stubSocketName);
mkdirSync(socketDir, { recursive: true });
writeFileSync(stubPath, '');
process.env.GENIE_TMUX_SOCKET = stubSocketName;
try {
// A fresh (recent) idle row: must NOT be reset because the socket
// is alive AND lastStateChange is inside the threshold. This
// exercises the preserved transient-retry path.
await register(
makeAgent({
id: 'live-socket-fresh-idle',
paneId: '%8888',
state: 'idle',
lastStateChange: new Date().toISOString(),
}),
);
const reset = await reconcileStaleSpawns(2);
expect(reset).not.toContain('live-socket-fresh-idle');
const a = await get('live-socket-fresh-idle');
expect(a!.state).toBe('idle');
} finally {
// biome-ignore lint/performance/noDelete: assigning undefined would set the string "undefined"
delete process.env.GENIE_TMUX_SOCKET;
try {
unlinkSync(stubPath);
} catch {
/* best-effort cleanup */
}
// silence unused-import warnings if any
void existsSync;
void tmpdir;
}
});
});

describe('templates', () => {
Expand Down
128 changes: 111 additions & 17 deletions src/lib/agent-registry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,20 @@ import { recordAuditEvent } from './audit.js';
import { type Sql, getConnection } from './db.js';
import type { AgentIdentity, ExecutorState } from './executor-types.js';
import type { ProviderName } from './provider-adapters.js';
import { isPaneAlive } from './tmux.js';
import { isPaneAlive, isTmuxSocketAlive } from './tmux.js';

/**
* Resolve the tmux socket name a worker row is expected to live on.
*
* The agents schema does not record a per-worker socket column — every
* genie worker runs on the global `-L <socket>` server whose name comes
* from `GENIE_TMUX_SOCKET` (default `genie`). This helper exists so the
* reconciliation loop has a single source of truth and legacy rows
* missing any future per-row socket column fall back safely.
*/
function resolveWorkerSocketName(): string {
return process.env.GENIE_TMUX_SOCKET || 'genie';
}

export type AgentState = 'spawning' | 'working' | 'idle' | 'permission' | 'question' | 'done' | 'error' | 'suspended';
export type TransportType = 'tmux' | 'inline';
Expand Down Expand Up @@ -284,6 +297,7 @@ export async function list(): Promise<Agent[]> {
*
* @param thresholdSeconds - How long an agent must be stuck before reset (default: 60)
*/
// biome-ignore lint/complexity/noExcessiveCognitiveComplexity: three sequential reconciliation passes with socket-grouping fast path
export async function reconcileStaleSpawns(thresholdSeconds = 60): Promise<string[]> {
try {
const sql = await getConnection();
Expand Down Expand Up @@ -313,7 +327,102 @@ export async function reconcileStaleSpawns(thresholdSeconds = 60): Promise<strin
AND pane_id IS NOT NULL AND pane_id != ''
AND started_at < now() - interval '1 second' * ${thresholdSeconds}
`;

// Third pass (query): dead-pane zombies in active (non-spawning) states.
// Rows whose state is idle/working/permission/question but whose tmux
// pane no longer exists (e.g. user killed the session, machine rebooted,
// process crashed without state update). These count toward the resume
// concurrency cap and permanently block auto-resume once accumulated.
// Only rows matching the tmux pane pattern `%\d+` are candidates —
// synthetic paneIds ('sdk', 'inline', etc.) are non-tmux transports
// with their own liveness source and must not be touched here.
const activeDeadCandidates = await sql<{ id: string; pane_id: string; state: string }[]>`
SELECT id, pane_id, state FROM agents
WHERE state IN ('idle', 'working', 'permission', 'question')
AND pane_id ~ '^%[0-9]+$'
AND last_state_change < now() - interval '1 second' * ${thresholdSeconds}
`;

// Dead-socket short-circuit (Bug 4).
//
// Before touching tmux, group the 2nd+3rd pass candidates by the tmux
// socket they live on and check the socket file once per unique name.
// If the socket is gone, every worker on it is permanently dead —
// no need to probe `isPaneAlive` (which would just raise
// `TmuxUnreachableError` and cause the per-row catch to skip the
// worker forever, as reported in the Bug 4 scheduler-log spam).
//
// Workers on a LIVE socket keep the existing per-worker
// isPaneAlive + try/catch behaviour — that path is correct for
// "socket up but pane gone" and for transient tmux blips.
type SpawningRow = { id: string; pane_id: string };
type ActiveRow = { id: string; pane_id: string; state: string };
const socketBuckets = new Map<string, { spawning: SpawningRow[]; active: ActiveRow[] }>();
for (const row of staleWithPane) {
const socket = resolveWorkerSocketName();
const bucket = socketBuckets.get(socket) ?? { spawning: [], active: [] };
bucket.spawning.push(row);
socketBuckets.set(socket, bucket);
}
for (const row of activeDeadCandidates) {
const socket = resolveWorkerSocketName();
const bucket = socketBuckets.get(socket) ?? { spawning: [], active: [] };
bucket.active.push(row);
socketBuckets.set(socket, bucket);
}

const liveSpawning: SpawningRow[] = [];
const liveActive: ActiveRow[] = [];
for (const [socketName, bucket] of socketBuckets) {
if (isTmuxSocketAlive(socketName)) {
liveSpawning.push(...bucket.spawning);
liveActive.push(...bucket.active);
continue;
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
// Socket is dead — mark every candidate on it as error in one pass.
const allIds = [...bucket.spawning.map((r) => r.id), ...bucket.active.map((r) => r.id)];
if (allIds.length === 0) continue;

// Per-row state-guarded UPDATE preserves the concurrent-transition
// race protection from the original code.
for (const row of bucket.spawning) {
const updated = await sql<{ id: string }[]>`
UPDATE agents
SET state = 'error', last_state_change = now(), pane_id = ''
WHERE id = ${row.id} AND state = 'spawning'
RETURNING id
`;
if (updated.length > 0) {
console.error(
`[reconcile] Reset agent ${row.id} (dead socket ${socketName}, pane ${row.pane_id}) from spawning → error`,
);
resetIds.push(row.id);
}
}
for (const row of bucket.active) {
const prevState = row.state;
const updated = await sql<{ id: string }[]>`
UPDATE agents
SET state = 'error', last_state_change = now(), pane_id = ''
WHERE id = ${row.id} AND state = ${prevState}
RETURNING id
`;
if (updated.length > 0) {
console.error(
`[reconcile] Reset agent ${row.id} (dead socket ${socketName}, pane ${row.pane_id}) from ${prevState} → error`,
);
resetIds.push(row.id);
}
}
Comment on lines +388 to +416

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

The sequential UPDATE queries inside these loops can be optimized by performing a single batch update per bucket using WHERE id IN ${sql(ids)}. This significantly reduces database round-trips, which is especially important when reconciling a large number of stale agents. Additionally, resolveWorkerSocketName() should be called once outside the loops to avoid redundant environment variable lookups and object creation on every iteration.

recordAuditEvent('worker', socketName, 'recovery_socket_dead', 'reconciler', {
socket: socketName,
worker_ids: allIds,
worker_count: allIds.length,
}).catch(() => {});
}
Comment on lines +376 to +422

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟡 Minor

Dead-socket path skips per-worker state_changed audit events.

The live-socket branch at Lines 405–432 emits a per-worker state_changed audit event when a zombie is reset, matching the pattern used everywhere else in this file (see also Lines 274 and 394). The dead-socket branch records only one aggregate recovery_socket_dead event keyed by the socket name — so workers killed via the fast path are invisible to per-worker state-transition audit queries and dashboards that pivot on entity_id.

Also, using socketName as the entity_id under entity_type = 'worker' is a type/id mismatch; consider entity_type = 'tmux_socket' (or similar) for the aggregate event, and still emit per-worker state_changed rows inside the loops at Lines 347 and 361.

🔧 Suggested patch
       for (const row of bucket.spawning) {
         const updated = await sql<{ id: string }[]>`
           UPDATE agents
           SET state = 'error', last_state_change = now(), pane_id = ''
           WHERE id = ${row.id} AND state = 'spawning'
           RETURNING id
         `;
         if (updated.length > 0) {
           console.error(
             `[reconcile] Reset agent ${row.id} (dead socket ${socketName}, pane ${row.pane_id}) from spawning → error`,
           );
+          recordAuditEvent('worker', row.id, 'state_changed', 'reconciler', {
+            state: 'error',
+            reason: 'socket_dead',
+            previous_state: 'spawning',
+            socket: socketName,
+          }).catch(() => {});
           resetIds.push(row.id);
         }
       }
       for (const row of bucket.active) {
         const prevState = row.state;
         const updated = await sql<{ id: string }[]>`
           UPDATE agents
           SET state = 'error', last_state_change = now(), pane_id = ''
           WHERE id = ${row.id} AND state = ${prevState}
           RETURNING id
         `;
         if (updated.length > 0) {
           console.error(
             `[reconcile] Reset agent ${row.id} (dead socket ${socketName}, pane ${row.pane_id}) from ${prevState} → error`,
           );
+          recordAuditEvent('worker', row.id, 'state_changed', 'reconciler', {
+            state: 'error',
+            reason: 'socket_dead',
+            previous_state: prevState,
+            socket: socketName,
+          }).catch(() => {});
           resetIds.push(row.id);
         }
       }
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
for (const [socketName, bucket] of socketBuckets) {
if (isTmuxSocketAlive(socketName)) {
liveSpawning.push(...bucket.spawning);
liveActive.push(...bucket.active);
continue;
}
// Socket is dead — mark every candidate on it as error in one pass.
const allIds = [...bucket.spawning.map((r) => r.id), ...bucket.active.map((r) => r.id)];
if (allIds.length === 0) continue;
// Per-row state-guarded UPDATE preserves the concurrent-transition
// race protection from the original code.
for (const row of bucket.spawning) {
const updated = await sql<{ id: string }[]>`
UPDATE agents
SET state = 'error', last_state_change = now(), pane_id = ''
WHERE id = ${row.id} AND state = 'spawning'
RETURNING id
`;
if (updated.length > 0) {
console.error(
`[reconcile] Reset agent ${row.id} (dead socket ${socketName}, pane ${row.pane_id}) from spawning → error`,
);
resetIds.push(row.id);
}
}
for (const row of bucket.active) {
const prevState = row.state;
const updated = await sql<{ id: string }[]>`
UPDATE agents
SET state = 'error', last_state_change = now(), pane_id = ''
WHERE id = ${row.id} AND state = ${prevState}
RETURNING id
`;
if (updated.length > 0) {
console.error(
`[reconcile] Reset agent ${row.id} (dead socket ${socketName}, pane ${row.pane_id}) from ${prevState} → error`,
);
resetIds.push(row.id);
}
}
recordAuditEvent('worker', socketName, 'recovery_socket_dead', 'reconciler', {
socket: socketName,
worker_ids: allIds,
worker_count: allIds.length,
}).catch(() => {});
}
for (const [socketName, bucket] of socketBuckets) {
if (isTmuxSocketAlive(socketName)) {
liveSpawning.push(...bucket.spawning);
liveActive.push(...bucket.active);
continue;
}
// Socket is dead — mark every candidate on it as error in one pass.
const allIds = [...bucket.spawning.map((r) => r.id), ...bucket.active.map((r) => r.id)];
if (allIds.length === 0) continue;
// Per-row state-guarded UPDATE preserves the concurrent-transition
// race protection from the original code.
for (const row of bucket.spawning) {
const updated = await sql<{ id: string }[]>`
UPDATE agents
SET state = 'error', last_state_change = now(), pane_id = ''
WHERE id = ${row.id} AND state = 'spawning'
RETURNING id
`;
if (updated.length > 0) {
console.error(
`[reconcile] Reset agent ${row.id} (dead socket ${socketName}, pane ${row.pane_id}) from spawning → error`,
);
recordAuditEvent('worker', row.id, 'state_changed', 'reconciler', {
state: 'error',
reason: 'socket_dead',
previous_state: 'spawning',
socket: socketName,
}).catch(() => {});
resetIds.push(row.id);
}
}
for (const row of bucket.active) {
const prevState = row.state;
const updated = await sql<{ id: string }[]>`
UPDATE agents
SET state = 'error', last_state_change = now(), pane_id = ''
WHERE id = ${row.id} AND state = ${prevState}
RETURNING id
`;
if (updated.length > 0) {
console.error(
`[reconcile] Reset agent ${row.id} (dead socket ${socketName}, pane ${row.pane_id}) from ${prevState} → error`,
);
recordAuditEvent('worker', row.id, 'state_changed', 'reconciler', {
state: 'error',
reason: 'socket_dead',
previous_state: prevState,
socket: socketName,
}).catch(() => {});
resetIds.push(row.id);
}
}
recordAuditEvent('worker', socketName, 'recovery_socket_dead', 'reconciler', {
socket: socketName,
worker_ids: allIds,
worker_count: allIds.length,
}).catch(() => {});
}
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/lib/agent-registry.ts` around lines 335 - 381, The dead-socket branch
(inside the socketBuckets loop using isTmuxSocketAlive) currently only records
an aggregate recordAuditEvent('worker', socketName, 'recovery_socket_dead', ...)
and skips emitting per-worker state_changed audit rows like the live-socket
branch does; update the code so that inside the loops iterating bucket.spawning
and bucket.active (where you already update agent rows and push resetIds) you
also call recordAuditEvent with the per-worker event (use entity_type = 'worker'
and entity_id = row.id, event type 'state_changed' and include previous/new
state metadata) for every agent actually transitioned, and change the aggregate
event to use an appropriate entity_type such as 'tmux_socket' (keep entity_id =
socketName) so the aggregate recovery_socket_dead event no longer mislabels
workers.


// Workers on live sockets: run the original per-row liveness probe.
for (const row of liveSpawning) {
try {
const alive = await isPaneAlive(row.pane_id);
if (!alive) {
Expand All @@ -334,22 +443,7 @@ export async function reconcileStaleSpawns(thresholdSeconds = 60): Promise<strin
// when we can't verify pane status
}
}

// Third pass: dead-pane zombies in active (non-spawning) states.
// Rows whose state is idle/working/permission/question but whose tmux
// pane no longer exists (e.g. user killed the session, machine rebooted,
// process crashed without state update). These count toward the resume
// concurrency cap and permanently block auto-resume once accumulated.
// Only rows matching the tmux pane pattern `%\d+` are candidates —
// synthetic paneIds ('sdk', 'inline', etc.) are non-tmux transports
// with their own liveness source and must not be touched here.
const activeDeadCandidates = await sql<{ id: string; pane_id: string; state: string }[]>`
SELECT id, pane_id, state FROM agents
WHERE state IN ('idle', 'working', 'permission', 'question')
AND pane_id ~ '^%[0-9]+$'
AND last_state_change < now() - interval '1 second' * ${thresholdSeconds}
`;
for (const row of activeDeadCandidates) {
for (const row of liveActive) {
try {
const alive = await isPaneAlive(row.pane_id);
if (!alive) {
Expand Down
Loading
Loading