87f7c03535
- Delete plugin/.mcp.json; mounting fleet is now the instance's or the project's job, not the plugin's. The setup skill still writes the project .mcp.json entry. - register.js: read FLEETD_MCP_URL via $.env.get, defaulting to http://127.0.0.1:8765/mcp, documented at https://code.claude.com/docs/en/plugins/mods/api.md ("$.env — get and set environment variables"). - Gate the fleet_inbox poll by role: a worker or architect learns its role once from fleet_whoami and then never polls fleet_inbox; every other role polls as before. An unreachable daemon retries whoami on a later tick rather than deciding. - Bump 0.2.0 -> 0.3.0 in plugin.json and marketplace.json, and update README/SKILL.md to drop every claim that the plugin mounts fleetd. - Add 6 tests for the role gating; all 21 tests and `claude plugin validate` pass.
282 lines
11 KiB
JavaScript
282 lines
11 KiB
JavaScript
// fleet mod — session messaging for the claude-bridge fleet.
|
|
//
|
|
// $.store is kept under CLAUDE_CONFIG_DIR, so /fleet-peers and /fleet-mail reach
|
|
// only sessions that share this session's config dir. fleetd names a caller by
|
|
// its pane, not its account, so /fleet-whoami and anything else sent through
|
|
// fleetTool reach the fleet from either account.
|
|
//
|
|
// Delivery uses $.prompt.submit, so nothing is typed into a pane and no prompt
|
|
// box is read.
|
|
//
|
|
// There are two inboxes. The $.store one carries /fleet-mail between sessions
|
|
// that share this config dir. fleet_inbox carries what the fleet queued for
|
|
// this pane, from any account, and calling it is also what tells fleetd to
|
|
// queue here rather than type into the terminal.
|
|
|
|
const PRESENCE_PREFIX = 'presence:'
|
|
const INBOX_PREFIX = 'inbox:'
|
|
const PRESENCE_REFRESH_MS = 15_000
|
|
const INBOX_POLL_MS = 3_000
|
|
// fleetd stops queueing for a pane that goes quiet, so this must stay well
|
|
// under the window the daemon allows between calls.
|
|
const FLEETD_INBOX_POLL_MS = 3_000
|
|
// A session whose presence row is older than this is treated as gone. It must
|
|
// exceed PRESENCE_REFRESH_MS by enough that one missed refresh is not a death.
|
|
const PRESENCE_STALE_MS = 60_000
|
|
const MAX_INBOX = 50
|
|
|
|
/** The store key holding one session's queued messages. */
|
|
function inboxKey(sessionId) {
|
|
return INBOX_PREFIX + sessionId
|
|
}
|
|
|
|
/** The store key holding one session's presence row. */
|
|
function presenceKey(sessionId) {
|
|
return PRESENCE_PREFIX + sessionId
|
|
}
|
|
|
|
/**
|
|
* Append one message to a target's inbox.
|
|
*
|
|
* $.store has no compare-and-swap, so two senders writing in the same instant
|
|
* can lose a message. Callers that need delivery confirmed should read the
|
|
* inbox back.
|
|
*/
|
|
async function deliver($, target, message) {
|
|
const key = inboxKey(target)
|
|
const queued = (await $.store.get(key)) || []
|
|
queued.push(message)
|
|
// Keep the newest: an unread inbox must not grow without bound.
|
|
const kept = queued.slice(-MAX_INBOX)
|
|
await $.store.set(key, kept)
|
|
return kept.length
|
|
}
|
|
|
|
/** Every session that refreshed its presence row recently, newest first. */
|
|
async function livePeers($, now) {
|
|
const keys = await $.store.keys()
|
|
const rows = []
|
|
for (const key of keys) {
|
|
if (!key.startsWith(PRESENCE_PREFIX)) continue
|
|
const row = await $.store.get(key)
|
|
if (!row || typeof row.at !== 'number') continue
|
|
if (now - row.at > PRESENCE_STALE_MS) continue
|
|
rows.push(row)
|
|
}
|
|
rows.sort((a, b) => b.at - a.at)
|
|
return rows
|
|
}
|
|
|
|
const FLEETD_MCP_DEFAULT = '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, and the URL it was opened against. 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
|
|
let fleetdUrl = FLEETD_MCP_DEFAULT
|
|
|
|
/**
|
|
* Open an MCP session on the local fleetd and hold it for later calls.
|
|
*
|
|
* Cleared first, so a failure here leaves no dead id behind for the next call to reuse. Reads
|
|
* FLEETD_MCP_URL fresh on every open, so a session opened after the daemon moves uses the new
|
|
* address.
|
|
*/
|
|
async function openFleetSession($) {
|
|
mcpSessionId = null
|
|
fleetdUrl = (await $.env.get('FLEETD_MCP_URL')) || FLEETD_MCP_DEFAULT
|
|
const init = await $.http.fetch(fleetdUrl, {
|
|
method: 'POST',
|
|
headers: MCP_HEADERS,
|
|
body: JSON.stringify({
|
|
jsonrpc: '2.0', id: 1, method: 'initialize',
|
|
params: { protocolVersion: '2025-06-18', capabilities: {}, clientInfo: { name: 'fleet-mod', version: '0' } },
|
|
}),
|
|
})
|
|
if (!init.ok) throw new Error('fleetd initialize failed with status ' + init.status)
|
|
const opened = init.headers['mcp-session-id']
|
|
await $.http.fetch(fleetdUrl, {
|
|
method: 'POST',
|
|
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(fleetdUrl, {
|
|
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)
|
|
if (body.error) throw new Error(body.error.message)
|
|
return body.result.content.map((c) => c.text).join('\n')
|
|
}
|
|
|
|
export function register(on) {
|
|
on('session.start', async ($, e, next) => {
|
|
const self = await $.session.id()
|
|
|
|
// Announce before the first refresh is due, or a session shorter than one
|
|
// refresh interval never appears to its peers at all.
|
|
await $.store.set(presenceKey(self), {
|
|
sessionId: self,
|
|
cwd: await $.session.cwd(),
|
|
at: await $.clock.now(),
|
|
})
|
|
|
|
// Keep announcing: a row that stops being refreshed is how another session
|
|
// learns this one is gone.
|
|
$.clock.every(PRESENCE_REFRESH_MS, async () => {
|
|
const now = await $.clock.now()
|
|
await $.store.set(presenceKey(self), {
|
|
sessionId: self,
|
|
cwd: await $.session.cwd(),
|
|
at: now,
|
|
})
|
|
})
|
|
|
|
// Collect this session's mail and hand it to Claude. $.prompt.submit waits
|
|
// for the session to be idle, so this never lands mid-turn.
|
|
$.clock.every(INBOX_POLL_MS, async () => {
|
|
const key = inboxKey(self)
|
|
const queued = (await $.store.get(key)) || []
|
|
if (queued.length === 0) return
|
|
await $.store.set(key, [])
|
|
for (const message of queued) {
|
|
await $.prompt.submit({
|
|
text: 'Message from fleet session ' + message.from + ':\n\n' + message.text,
|
|
})
|
|
}
|
|
})
|
|
|
|
// A spawned worker or architect already gets its brief pasted into its pane, so this
|
|
// poll would be a second, redundant delivery path for it. Every other role collects
|
|
// its own mail through this poll.
|
|
let role = null
|
|
|
|
// Collect what the fleet queued for this pane and hand each message to
|
|
// Claude. Every call also renews fleetd's record that this pane collects its
|
|
// own mail, so an empty answer still has to be asked for.
|
|
$.clock.every(FLEETD_INBOX_POLL_MS, async () => {
|
|
if (role === null) {
|
|
try {
|
|
role = JSON.parse(await fleetTool($, 'fleet_whoami', {})).role
|
|
} catch {
|
|
// A daemon that is down, or a pane fleetd cannot place, is the ordinary
|
|
// case on a host with no fleet running. Retry on the next tick.
|
|
return
|
|
}
|
|
}
|
|
if (role === 'worker' || role === 'architect') return
|
|
|
|
let collected
|
|
try {
|
|
collected = JSON.parse(await fleetTool($, 'fleet_inbox', {}))
|
|
} catch {
|
|
// 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 || []) {
|
|
await $.prompt.submit({ text: 'Message from the fleet, via fleetd:\n\n' + text })
|
|
}
|
|
})
|
|
|
|
await $.command.register({
|
|
name: 'fleet-peers',
|
|
description: 'List fleet sessions on this machine, including other accounts',
|
|
})
|
|
await $.command.register({
|
|
name: 'fleet-whoami',
|
|
description: 'Show who fleetd says this session is',
|
|
})
|
|
await $.command.register({
|
|
name: 'fleet-mail',
|
|
description: 'Send a message to a fleet session on this machine',
|
|
argumentHint: '<sessionId> <text>',
|
|
// Runs even while Claude is working, so a correction is never queued
|
|
// behind the turn it is meant to correct.
|
|
immediate: true,
|
|
})
|
|
return next(e)
|
|
})
|
|
|
|
on('command.run', { command: 'fleet-peers' }, async ($) => {
|
|
const now = await $.clock.now()
|
|
const self = await $.session.id()
|
|
const peers = await livePeers($, now)
|
|
if (peers.length === 0) return { text: 'No fleet sessions have announced themselves yet.' }
|
|
const lines = peers.map((p) => {
|
|
const age = Math.round((now - p.at) / 1000)
|
|
const mark = p.sessionId === self ? ' (this session)' : ''
|
|
return p.sessionId + ' ' + p.cwd + ' seen ' + age + 's ago' + mark
|
|
})
|
|
return { text: 'Fleet sessions on this machine:\n' + lines.join('\n') }
|
|
})
|
|
|
|
on('command.run', { command: 'fleet-mail' }, async ($, e) => {
|
|
const args = (e.args || '').trim()
|
|
const split = args.indexOf(' ')
|
|
if (split < 1) return { text: 'Usage: /fleet-mail <sessionId> <text>' }
|
|
const target = args.slice(0, split)
|
|
const text = args.slice(split + 1).trim()
|
|
if (text === '') return { text: 'Usage: /fleet-mail <sessionId> <text>' }
|
|
|
|
const now = await $.clock.now()
|
|
const peers = await livePeers($, now)
|
|
if (!peers.some((p) => p.sessionId === target)) {
|
|
return { text: 'No live fleet session ' + target + '. Run /fleet-peers.' }
|
|
}
|
|
const self = await $.session.id()
|
|
const depth = await deliver($, target, { from: self, text: text, at: now })
|
|
return { text: 'Queued for ' + target + ' (' + depth + ' in its inbox).' }
|
|
})
|
|
|
|
on('command.run', { command: 'fleet-whoami' }, async ($) => {
|
|
try {
|
|
return { text: 'fleetd says: ' + (await fleetTool($, 'fleet_whoami', {})) }
|
|
} catch (err) {
|
|
return { text: 'fleetd unreachable: ' + err.message }
|
|
}
|
|
})
|
|
|
|
// Record what arrives over Claude Code's own channel, so a message delivered
|
|
// by the fleet and one delivered by SendMessage can be told apart.
|
|
//
|
|
// This hook gates delivery, so it must never decide the message's fate. The
|
|
// .catch handler passes the message on when the logging above throws.
|
|
on('session.receive', async ($, e, next) => {
|
|
$.ui.log('fleet: inbound ' + (e.origin && e.origin.kind) + ', ' + String(e.text).length + ' chars')
|
|
return next(e)
|
|
}).catch(async ($, e, next) => {
|
|
if (next.called) return undefined
|
|
return next(e)
|
|
})
|
|
}
|