fleetd #788: review fixes — cancel vs a collected offer, one MCP session
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 53s
CI / build (pull_request) Failing after 1m59s

Four items from the lead review of PR #789.

1. Injector.cancel touched the inbox offer unconditionally. A pane takes an
   offer on the MCP thread, so between its drain and the next poll the head
   entry is taken but still QUEUED. Cancelling it returned CANCELLED for text
   the pane already held, and cancelling a LATER entry nulled t.inboxOffer and
   so lost the only record that the head was taken, which made the next poll
   offer it a second time. cancel now touches the offer only when the entry is
   the offered head, withdraws first, and answers DELIVERED when the pane took
   it.

2. The mod sent initialize on every fleetTool call and never a DELETE, so
   fleetd's transport kept one session per call -- 20 a minute per pane at a
   3s poll. The mod now holds one MCP session and reopens it only on the 404
   "Session not found" fleetd answers for an id it no longer holds (measured
   against the live daemon, not assumed).

3. MessageService.collectInbox had been inserted between reply's javadoc and
   reply, leaving that block attached to nothing. Moved above it.

4. The timer's catch comment said a throw would kill the timer. It does not;
   the catch keeps every tick from writing an error to the debug log.

mvn clean install: Tests run: 2202, Failures: 0, Errors: 0, Skipped: 0,
BUILD SUCCESS. Counted again over 182 target/surefire-reports/TEST-*.xml:
2202/0/0/0. claude plugin validate plugin and claude plugin test plugin both
exit 0, 15 pass 0 fail.

Three mutation checks, each reverted:
- the exact pre-review cancel body: the two new cancel tests fail with
  "expected: <DELIVERED> but was: <CANCELLED>" and "expected: <true> but was:
  <false>", the other six pass;
- initialize on every call: the two session tests fail (1 vs 2 initializes);
- no 404 retry: only the retry test fails (2 vs 1 initializes).
This commit is contained in:
Dai Ha
2026-10-06 05:48:34 +02:00
parent ac535503ff
commit d9fedfbc5a
5 changed files with 206 additions and 35 deletions
@@ -388,7 +388,8 @@ public final class Injector {
/**
* Cancel this exact queued delivery. The target monitor serializes this operation with
* {@link #onStatus}: if delivery wins that race, this returns {@link Cancellation#DELIVERED}
* rather than claiming the message remained queued.
* rather than claiming the message remained queued. A message a mod-served pane has already
* collected answers the same way, even though no poll has recorded that delivery yet.
*/
public Cancellation cancel(Delivery delivery) {
Pending p = delivery.pending;
@@ -397,15 +398,25 @@ public final class Injector {
return cancellationOf(p);
}
synchronized (t) {
if (p.state != Pending.State.QUEUED || !t.queue.remove(p)) {
if (p.state != Pending.State.QUEUED) {
return cancellationOf(p);
}
if (t.queue.peek() == p && t.inboxOffer != null) {
// This exact entry is the one offered to a mod-served pane. The pane takes an
// offer on its own thread, so withdraw first and then read the outcome: a taken
// offer means the pane already holds this text, and the next poll records that
// delivery. Cancelling it would tell the caller nothing arrived while the pane
// acts on it.
paneInbox.withdrawAll(p.target);
if (t.inboxOffer.taken()) {
return Cancellation.DELIVERED;
}
t.inboxOffer = null;
}
if (!t.queue.remove(p)) {
return cancellationOf(p);
}
p.state = Pending.State.CANCELLED;
// The cancelled entry may be the one currently offered to a mod-served pane, so take
// the offer back: a cancelled message the caller was told never arrived must not stay
// collectable. The next poll re-offers whatever is at the head by then.
paneInbox.withdrawAll(p.target);
t.inboxOffer = null;
if (isQuiescent(t)) {
targets.remove(p.target, t);
}
@@ -514,6 +514,20 @@ public final class MessageService {
return false;
}
/**
* Hand {@code session} every message queued for it that it has not collected yet, and record
* that it collects its own mail. While that record is fresh, delivery to that session is
* offered for collection instead of typed into its terminal; once it goes stale, the terminal
* route takes over again with nothing lost.
*
* <p>The messages are returned in the order they were queued, and are removed by this call.
* An empty list is an ordinary answer: a session polling on a timer keeps itself collecting
* between messages.
*/
public List<String> collectInbox(String session) {
return injector.collectInbox(session);
}
/**
* Route a worker's explicit {@code fleet_reply}: resolve an open send, complete an async ticket
* still parked waiting on this exact turn's answer, or — only once neither applies — queue it in
@@ -552,20 +566,6 @@ public final class MessageService {
* anything
* @return which of the three ways (fleetd #365) the reply actually landed — never {@code null}
*/
/**
* Hand {@code session} every message queued for it that it has not collected yet, and record
* that it collects its own mail. While that record is fresh, delivery to that session is
* offered for collection instead of typed into its terminal; once it goes stale, the terminal
* route takes over again with nothing lost.
*
* <p>The messages are returned in the order they were queued, and are removed by this call.
* An empty list is an ordinary answer: a session polling on a timer keeps itself collecting
* between messages.
*/
public List<String> collectInbox(String session) {
return injector.collectInbox(session);
}
public ReplyOutcome reply(String session, String content) {
if (content == null || content.isBlank()) {
throw new IllegalArgumentException("content is required");
@@ -136,6 +136,59 @@ class InjectorModServedDeliveryTest {
assertEquals(List.of("kept"), injector.collectInbox(MOD), "control: an offer is collectable");
}
@Test
void aMessageThePaneAlreadyCollectedCannotBeCancelled() {
injector.collectInbox(MOD);
// The control first: an offer the pane has not taken really is cancellable, so the
// different answer below is the collection and not a cancel that gave up on this route.
Injector.Delivery untaken = injector.enqueue(MOD, "retracted", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(Injector.Cancellation.CANCELLED, injector.cancel(untaken),
"control: an uncollected offer is still cancellable");
Injector.Delivery taken = injector.enqueue(MOD, "do the task", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(List.of("do the task"), injector.collectInbox(MOD), "the pane takes the offer");
// No poll has run since the pane took it, so the entry is still at the head and still
// QUEUED: the state alone cannot tell this case from an uncollected offer.
assertEquals(Injector.Cancellation.DELIVERED, injector.cancel(taken),
"the pane holds this text and will act on it, so nothing can be cancelled");
injector.onStatus(MOD, AgentStatus.IDLE);
assertTrue(taken.completion().isDone(), "the next poll records the delivery");
assertEquals(List.of(), typed(), "and nothing was typed into the pane");
}
@Test
void cancellingALaterMessageLeavesACollectedOneDeliveredOnce() {
injector.collectInbox(MOD);
Injector.Delivery first = injector.enqueue(MOD, "first", TestTurnTokens.inert(MOD));
Injector.Delivery second = injector.enqueue(MOD, "second", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE); // offers the head
assertEquals(List.of("first"), injector.collectInbox(MOD), "the pane takes the head");
assertEquals(Injector.Cancellation.CANCELLED, injector.cancel(second),
"a message behind the collected one was never offered, so it cancels");
injector.onStatus(MOD, AgentStatus.IDLE);
assertTrue(first.completion().isDone(),
"cancelling a later message must not lose the record that the head was taken");
assertEquals(List.of(), injector.collectInbox(MOD),
"and the head must not be offered a second time");
// The control: the same injector still hands a later message over, so the empty
// collection above is this one not being re-offered rather than the route going quiet.
injector.onStatus(MOD, AgentStatus.WORKING); // the pane picks the collected message up
injector.onStatus(MOD, AgentStatus.IDLE); // and that turn ends
injector.enqueue(MOD, "third", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(List.of("third"), injector.collectInbox(MOD), "control: a later message is offered");
assertEquals(List.of(), typed(), "nothing took the terminal route");
}
@Test
void aPaneThatNeverCollectedIsTypedIntoFromTheStart() {
injector.enqueue(PTY, "do the task", TestTurnTokens.inert(PTY));
+42 -11
View File
@@ -69,14 +69,20 @@ async function livePeers($, now) {
const FLEETD_MCP = 'http://127.0.0.1:8765/mcp'
const MCP_HEADERS = { 'Content-Type': 'application/json', Accept: 'application/json, text/event-stream' }
// The status fleetd answers, with "Session not found", for an Mcp-Session-Id it no longer holds.
const MCP_SESSION_GONE = 404
// The MCP session every call below shares. fleetd keeps a server-side session per initialize and
// drops it only on a DELETE, so one initialize per call would leave a session behind every time.
let mcpSessionId = null
/**
* Call one fleet_* tool on the local fleetd, and return its text result.
* Open an MCP session on the local fleetd and hold it for later calls.
*
* fleetd names the caller from the TCP connection, so the answer is about this
* session's own pane, whichever Claude account the session runs on.
* Cleared first, so a failure here leaves no dead id behind for the next call to reuse.
*/
async function fleetTool($, tool, args) {
async function openFleetSession($) {
mcpSessionId = null
const init = await $.http.fetch(FLEETD_MCP, {
method: 'POST',
headers: MCP_HEADERS,
@@ -86,15 +92,40 @@ async function fleetTool($, tool, args) {
}),
})
if (!init.ok) throw new Error('fleetd initialize failed with status ' + init.status)
const headers = { ...MCP_HEADERS, 'Mcp-Session-Id': init.headers['mcp-session-id'] }
const opened = init.headers['mcp-session-id']
await $.http.fetch(FLEETD_MCP, {
method: 'POST', headers, body: JSON.stringify({ jsonrpc: '2.0', method: 'notifications/initialized' }),
})
const call = await $.http.fetch(FLEETD_MCP, {
method: 'POST',
headers,
headers: { ...MCP_HEADERS, 'Mcp-Session-Id': opened },
body: JSON.stringify({ jsonrpc: '2.0', method: 'notifications/initialized' }),
})
mcpSessionId = opened
}
/** Send one tools/call on the session this mod holds, and return the raw HTTP answer. */
function sendFleetToolCall($, tool, args) {
return $.http.fetch(FLEETD_MCP, {
method: 'POST',
headers: { ...MCP_HEADERS, 'Mcp-Session-Id': mcpSessionId },
body: JSON.stringify({ jsonrpc: '2.0', id: 2, method: 'tools/call', params: { name: tool, arguments: args } }),
})
}
/**
* Call one fleet_* tool on the local fleetd, and return its text result.
*
* fleetd names the caller from the TCP connection on every request, not from the MCP session, so
* the answer is about this session's own pane whichever Claude account the session runs on, and
* reusing one session never changes whose call it is.
*/
async function fleetTool($, tool, args) {
if (mcpSessionId === null) await openFleetSession($)
let call = await sendFleetToolCall($, tool, args)
if (call.status === MCP_SESSION_GONE) {
// A daemon restart drops every session it held. Open a new one and retry once.
await openFleetSession($)
call = await sendFleetToolCall($, tool, args)
}
if (!call.ok) throw new Error('fleetd ' + tool + ' failed with status ' + call.status)
// A tool call answers as a server-sent event: the JSON is on the data: line.
const line = call.text.split('\n').find((l) => l.startsWith('data:'))
const body = JSON.parse(line ? line.slice(5) : call.text)
@@ -148,8 +179,8 @@ export function register(on) {
collected = JSON.parse(await fleetTool($, 'fleet_inbox', {}))
} catch {
// A daemon that is down, or a pane fleetd cannot place, is the ordinary
// case on a host with no fleet running. Throwing here would kill the
// timer and with it every later message.
// case on a host with no fleet running. The timer survives a throw, so
// this only keeps every tick from writing an error to the debug log.
return
}
for (const text of collected.messages || []) {
+79 -3
View File
@@ -191,15 +191,35 @@ function stubSessionStart(on: any, submitted: string[]): Map<string, any> {
return store
}
/** Answer fleetd's three MCP requests, with the tool result read fresh on every call. */
function stubFleetdDynamic(on: any, toolText: () => string, fetches: string[]) {
/**
* Answer fleetd's MCP requests, with the tool result read fresh on every call.
*
* `fetches` collects every method sent, and `sessions` the Mcp-Session-Id of each tools/call.
* Each initialize hands out the next id, so a reused session and a reopened one differ.
* `sessionGone` makes a tools/call answer the way fleetd answers for a session id it no longer
* holds.
*/
function stubFleetdDynamic(
on: any,
toolText: () => string,
fetches: string[],
sessions: string[] = [],
sessionGone: () => boolean = () => false,
) {
let opened = 0
on('http.fetch', (_$: any, e: any) => {
const body = JSON.parse(e.init.body)
fetches.push(body.method)
if (body.method === 'initialize') {
return { value: { ok: true, status: 200, headers: { 'mcp-session-id': 'sid-1' }, text: '{}' } }
opened += 1
return { value: { ok: true, status: 200, headers: { 'mcp-session-id': 'sid-' + opened }, text: '{}' } }
}
if (body.method === 'tools/call') {
sessions.push(e.init.headers['Mcp-Session-Id'])
if (sessionGone()) {
const gone = '{"jsonRpcError":{"code":-32603,"message":"Session not found"}}'
return { value: { ok: false, status: 404, headers: {}, text: gone } }
}
const result = { jsonrpc: '2.0', id: 2, result: { content: [{ type: 'text', text: toolText() }] } }
return { value: { ok: true, status: 200, headers: {}, text: 'data: ' + JSON.stringify(result) + '\n' } }
}
@@ -207,6 +227,11 @@ function stubFleetdDynamic(on: any, toolText: () => string, fetches: string[]) {
})
}
/** How many of `fetches` were an initialize. */
function initializes(fetches: string[]): number {
return fetches.filter((m) => m === 'initialize').length
}
test('the fleetd inbox poll submits each collected message and names the sender', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
@@ -271,3 +296,54 @@ test('a fleetd that is down leaves the poll timer running', async ($, on) => {
expect(fetches.length).toBeGreaterThan(afterFirst)
expect(submitted.length).toBe(0)
})
test('a second inbox poll reuses the first MCP session', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const sessions: string[] = []
const inbox = { sessionId: 'term_self', count: 0, messages: [] as string[] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), fetches, sessions)
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
await clock.advance(3_000)
// fleetd holds a session per initialize and drops it only on a DELETE, so one initialize per
// poll would leave one behind every three seconds.
expect(sessions.length).toBe(2) // control: both polls really reached fleetd
expect(initializes(fetches)).toBe(1)
expect(sessions[1]).toBe(sessions[0])
})
test('a session fleetd no longer holds is opened again and the call retried', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const sessions: string[] = []
let gone = false
const inbox = { sessionId: 'term_self', count: 1, messages: ['the daemon restarted'] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), fetches, sessions, () => {
// Only the first call after the flag is set is refused; the retry succeeds.
const refuse = gone
gone = false
return refuse
})
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
// The control: an ordinary poll opens one session and needs no second one.
await clock.advance(3_000)
expect(initializes(fetches)).toBe(1)
expect(submitted.length).toBe(1)
gone = true
await clock.advance(3_000)
expect(initializes(fetches)).toBe(2)
expect(sessions[sessions.length - 1]).toBe('sid-2')
expect(submitted.length).toBe(2)
})