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
91 changes: 84 additions & 7 deletions plugins/dashboard/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ import {
ListToolsRequestSchema,
CallToolRequestSchema,
} from '@modelcontextprotocol/sdk/types.js'
import { readFileSync, mkdirSync, readdirSync, unlinkSync, watch, chmodSync, copyFileSync } from 'fs'
import { readFileSync, mkdirSync, readdirSync, unlinkSync, watch, chmodSync, copyFileSync, openSync, closeSync, writeSync, statSync } from 'fs'
import { homedir } from 'os'
import { join } from 'path'
import { installLifecycle, recordLifecycle } from './lifecycle.ts'
Expand All @@ -47,6 +47,16 @@ const AGENT_INBOX_DIR = join(STATE_DIR, 'agent-inbox')
// purpose: a nudge carries no payload and must never be mistaken for a message
// drop-file, so the two ingesters can never read each other's files.
const COLLECT_NOW_DIR = join(STATE_DIR, 'collect-now')
// DIVE-3809: the drain's `draining`/`rerun` guard below is PER-PROCESS. Two
// plugin processes for the same agent (an overlapping restart, a stray
// supervisor respawn) each fetch the SAME pending rows and push every message
// into the session twice, because the ack lands only after the notifications
// are sent. This file is the cross-process arm of that guard: O_EXCL create
// wins the drain, everyone else skips this pass and picks it up on the next
// sweep. Refuted as the CAUSE of the loss DIVE-3806 observed (lifecycle.log
// showed exactly one live process across that window) — it is still a real
// race, and it is scope 3 of this row.
const DRAIN_LOCK = join(STATE_DIR, 'pending-drain.lock')
mkdirSync(AGENT_INBOX_DIR, { recursive: true, mode: 0o700 })
mkdirSync(COLLECT_NOW_DIR, { recursive: true, mode: 0o700 })

Expand Down Expand Up @@ -302,10 +312,17 @@ function startCollectNowWatch(): void {
// DIVE-848 offline heal: a message sent while this box was unreachable never
// produced a drop file — it sits in the control plane with delivered_at NULL.
// Pull those on boot (and on a slow sweep), push them into the session, then
// ack so they stamp delivered. Ack only AFTER the notifications are sent; a
// ack so they stamp COLLECTED. Ack only AFTER the notifications are sent; a
// crash in between redelivers rather than losing the message. A row whose
// drop landed but whose delivered-stamp write failed may arrive twice — rare
// drop landed but whose collected-stamp write failed may arrive twice — rare
// and preferable to silence.
// DIVE-3809: the ack no longer stamps delivered_at. It could never attest
// delivery, and stamping the column `/pending` reads meant one wrong ack
// deleted the only copy. A collected row is now merely hidden for a TTL and
// comes back, bounded by an attempt count. Note the consequence for the
// empty-text branch below: it acks a row it never pushed, so such a row is
// re-offered until the attempt bound retires it — bounded and visible, where
// before it was silently destroyed.
// DIVE-3574: drainPending is now reachable from three places (boot, the 5-min
// timer, and a collect nudge that can fire several times a second while someone
// types in the dashboard) where it used to be reachable from two that could
Expand All @@ -316,11 +333,63 @@ function startCollectNowWatch(): void {
// landed after the fetch from waiting out the full timer.
let draining = false
let rerun = false

// A drain that dies mid-flight (SIGKILL, box reboot) leaves the lock file
// behind, and a stale lock that nothing clears would wedge the collect path
// permanently — the exact failure shape this row exists to remove. So the lock
// is TIME-BOUNDED: older than this and it is treated as abandoned and broken.
// One drain is a fetch + N notifications + an ack, all with short timeouts;
// two minutes is far past any healthy pass.
const DRAIN_LOCK_STALE_MS = 2 * 60_000

// Returns true if THIS process now holds the lock. Never throws: a filesystem
// that cannot support the lock must degrade to today's per-process-only
// behaviour, not stop the customer's message from being collected.
function acquireDrainLock(): boolean {
for (let attempt = 0; attempt < 2; attempt++) {
try {
const fd = openSync(DRAIN_LOCK, 'wx')
try { writeSync(fd, `${process.pid} ${new Date().toISOString()}\n`) } catch {}
closeSync(fd)
return true
} catch (err) {
if ((err as NodeJS.ErrnoException).code !== 'EEXIST') {
process.stderr.write(`dashboard channel: drain lock unavailable (${err}); per-process guard only\n`)
return true
}
// Held. Break it only if it is provably stale, then retry the create
// once — if another process wins that race, we simply skip this pass.
try {
const age = Date.now() - statSync(DRAIN_LOCK).mtimeMs
if (age <= DRAIN_LOCK_STALE_MS) return false
process.stderr.write(`dashboard channel: breaking stale drain lock (age ${Math.round(age / 1000)}s)\n`)
unlinkSync(DRAIN_LOCK)
} catch { return false }
}
}
return false
}

function releaseDrainLock(): void {
try { unlinkSync(DRAIN_LOCK) } catch {}
}

async function drainPending(): Promise<void> {
if (draining) { rerun = true; return }
draining = true
try {
await drainPendingOnce()
// Cross-process (DIVE-3809). Skipping is safe and NOT a lost message: the
// holder is draining the same rows right now, and anything it misses is
// re-offered by the control plane once the collect TTL expires.
if (!acquireDrainLock()) {
process.stderr.write('dashboard channel: another process holds the drain lock; skipping this pass\n')
return
}
try {
await drainPendingOnce()
} finally {
releaseDrainLock()
}
} finally {
draining = false
}
Expand Down Expand Up @@ -362,7 +431,7 @@ async function drainPendingOnce(): Promise<void> {
})
acked.push(m.id)
} catch (err) {
process.stderr.write(`dashboard channel: pending deliver failed for ${m.id}: ${err}\n`)
process.stderr.write(`dashboard channel: pending push failed for ${m.id}: ${err}\n`)
}
}
if (acked.length === 0) return
Expand All @@ -372,10 +441,18 @@ async function drainPendingOnce(): Promise<void> {
headers: { 'content-type': 'application/json' },
body: JSON.stringify({ agent: AGENT, ids: acked }),
})
// DIVE-3810: a non-2xx ack must reach the catch below. Without this the
// rows are logged as collected while the control plane still holds them
// uncollected — the exact split this row exists to close.
if (!ack.ok) throw new Error(`${ack.status}`)
process.stderr.write(`dashboard channel: healed ${acked.length} undelivered message(s)\n`)
// DIVE-3809: "collected", not "delivered" or "healed". This ack attests
// that the notification's bytes entered the stdout pipe — the SDK's send()
// has no reject path, and a client with nothing subscribed drops the
// notification silently — so it can never say the session displayed it.
// The control plane now re-offers a collected row whose TTL expires.
process.stderr.write(`dashboard channel: collected ${acked.length} pending message(s) (collection is not display)\n`)
} catch (err) {
process.stderr.write(`dashboard channel: pending ack failed (will redeliver next boot): ${err}\n`)
process.stderr.write(`dashboard channel: pending ack failed (row stays uncollected; re-offered next sweep): ${err}\n`)
}
}

Expand Down
168 changes: 168 additions & 0 deletions test/dive3809-drain-lock.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,168 @@
// DIVE-3809 scope 3: the drain guard must serialise ACROSS PROCESSES.
//
// server.ts's `draining`/`rerun` pair is module state — it serialises the three
// callers inside ONE process and is blind to a second one. Two plugin processes
// for the same agent (an overlapping restart, a stray respawn) each fetch the
// SAME pending rows and push every message into the session TWICE, because the
// ack only lands after the notifications are sent.
//
// DIVE-3806 REFUTED this as the cause of the loss it observed — lifecycle.log
// showed exactly one live process across that window — so this is not a fix for
// that. It is a real race, it is scope 3, and dashboard-collect-now.test.ts
// cannot see it: that harness runs a single process, so it grades the
// per-process half and would stay green with the lock deleted.
//
// The arm therefore runs TWO real server.ts processes against ONE stub control
// plane and ONE shared state dir, and counts pushes across BOTH.
import { describe, test, expect } from 'bun:test'
import { mkdtempSync, rmSync, mkdirSync, writeFileSync, existsSync, utimesSync } from 'node:fs'
import { tmpdir } from 'node:os'
import { join } from 'node:path'

const SERVER = join(import.meta.dir, '..', 'plugins', 'dashboard', 'server.ts')
const BOOT_MS = 5_000

function spawnPlugin(dir: string, port: number, delivered: string[]) {
const proc = Bun.spawn(['bun', SERVER], {
env: {
...process.env,
DASHBOARD_STATE_DIR: dir,
DASHBOARD_API_BASE: `http://127.0.0.1:${port}`,
CONNECTORD_TOKEN: 'test-token-abcdefghijkl',
USER: 'agent-dev',
},
stdin: 'pipe',
stdout: 'pipe',
stderr: 'pipe',
})
proc.stdin.write(
JSON.stringify({
jsonrpc: '2.0', id: 1, method: 'initialize',
params: { protocolVersion: '2024-11-05', capabilities: {}, clientInfo: { name: 'lock-test', version: '0' } },
}) + '\n'
)
proc.stdin.flush()
void (async () => {
const reader = proc.stdout.getReader()
const dec = new TextDecoder()
let buf = ''
for (;;) {
const { done, value } = await reader.read()
if (done) return
buf += dec.decode(value, { stream: true })
const lines = buf.split('\n')
buf = lines.pop() ?? ''
for (const line of lines) {
if (!line.trim()) continue
try {
const msg = JSON.parse(line)
if (msg.method === 'notifications/claude/channel') delivered.push(String(msg.params?.content ?? ''))
if (msg.id === 1) {
proc.stdin.write(JSON.stringify({ jsonrpc: '2.0', method: 'notifications/initialized' }) + '\n')
proc.stdin.flush()
}
} catch {}
}
}
})()
return proc
}

async function waitFor(pred: () => boolean, ms: number) {
const deadline = Date.now() + ms
while (Date.now() < deadline) {
if (pred()) return true
await Bun.sleep(25)
}
return pred()
}

describe('cross-process drain lock (DIVE-3809)', () => {
test('two plugin processes on one state dir push each message EXACTLY once', async () => {
const dir = mkdtempSync(join(tmpdir(), 'drain-lock-'))
const delivered: string[] = []
let queue = [
{ id: 1, text: 'first' },
{ id: 2, text: 'second' },
]
let pendingHits = 0
// Snapshot at REQUEST time then delay, exactly as dashboard-collect-now
// does: a stub that re-reads after the delay hands the second drain the
// state AFTER the first one's ack and silently removes the race.
const api = Bun.serve({
port: 0,
async fetch(req) {
const url = new URL(req.url)
if (url.pathname === '/server/messages/pending') {
pendingHits++
const snapshot = queue
await Bun.sleep(900)
return Response.json({ pending: snapshot })
}
if (url.pathname === '/server/messages/pending/ack') {
const body = (await req.json()) as { ids: number[] }
queue = queue.filter(m => !body.ids.includes(m.id))
return Response.json({ ok: true })
}
return new Response('not found', { status: 404 })
},
})
const a = spawnPlugin(dir, api.port, delivered)
const b = spawnPlugin(dir, api.port, delivered)
try {
// Both boot drains fire ~5s after start, i.e. genuinely overlapping.
await waitFor(() => delivered.length >= 2, BOOT_MS + 10_000)
// Settle: give a second (unserialised) drain every chance to double-push.
await Bun.sleep(3_000)
const counts = new Map<string, number>()
for (const d of delivered) counts.set(d, (counts.get(d) ?? 0) + 1)
expect(delivered.length).toBeGreaterThanOrEqual(2)
expect([...counts.entries()].filter(([, n]) => n > 1)).toEqual([])
expect(pendingHits).toBeGreaterThan(0)
} finally {
a.kill(); b.kill(); api.stop(true)
rmSync(dir, { recursive: true, force: true })
}
}, 40_000)

test('a stale lock left by a killed drain is broken, not honoured forever', async () => {
// The failure this forecloses: a lock file is the classic way to turn an
// intermittent loss into a permanent one. A drain killed mid-flight leaves
// the file behind, and if nothing breaks it the collect path is wedged for
// the life of the box — strictly worse than the bug being fixed.
const dir = mkdtempSync(join(tmpdir(), 'stale-lock-'))
mkdirSync(dir, { recursive: true })
const lock = join(dir, 'pending-drain.lock')
writeFileSync(lock, '999999 stale\n')
// Age it well past the 2-minute staleness bound. utimesSync, not `touch
// -d` — bun's built-in shell does not implement that flag.
const old = new Date(Date.now() - 10 * 60_000)
utimesSync(lock, old, old)
const delivered: string[] = []
let queue = [{ id: 7, text: 'after a crash' }]
const api = Bun.serve({
port: 0,
async fetch(req) {
const url = new URL(req.url)
if (url.pathname === '/server/messages/pending') return Response.json({ pending: queue })
if (url.pathname === '/server/messages/pending/ack') {
const body = (await req.json()) as { ids: number[] }
queue = queue.filter(m => !body.ids.includes(m.id))
return Response.json({ ok: true })
}
return new Response('not found', { status: 404 })
},
})
const p = spawnPlugin(dir, api.port, delivered)
try {
const got = await waitFor(() => delivered.length >= 1, BOOT_MS + 10_000)
expect(got).toBe(true)
expect(delivered[0]).toBe('after a crash')
// And the drain cleaned up after itself, so the next pass is not blocked.
expect(existsSync(lock)).toBe(false)
} finally {
p.kill(); api.stop(true)
rmSync(dir, { recursive: true, force: true })
}
}, 40_000)
})
Loading