Messages between the factory's Claude Code sessions over IoT Core

A Claude Code mod that lets the factory's sessions message each other over AWS IoT Core. A message starts a turn in the receiver; the receiver answers with the send tool. Tested with Claude Code 2.1.292 (claude --version).
~/.factory/bin/factory-agent <machine>.<role>
Every machine has this script. It sources ~/.factory/env, makes the role's worktree, and runs FACTORY_NAME=<name> claude ... --plugin-dir <plugin checkout>. The plugin checkout is ~/Desktop/synnaxlabs/foundation-wt/factory-plugin. Each new session moves it to the latest main; when the mod changed, the running sessions on that machine reload it.
FACTORY_NAME must be in roster.json, which also holds the broker host. When it is not set or not in the roster, the mod is off and the status line says why.
Each machine has ~/.factory/iot/{root-ca.pem,cert.pem,key.pem}, and the certificate CN is the machine (laptop, box1, box2). The IoT policy accepts a client ID and the sender level of a topic only when they start with <CN>.. A publish that breaks the policy drops the connection.
The topics:
factory/<to>/inbox/<from> -> a message, QoS 1: { id, text, sentAt }. The receiver takes the sender from the topic, never from the payload.factory/<to>/acks/<from> -> { id }, QoS 1, when the receiver has stored the message. Each session also sends a QoS 0 probe to its own acks topic each minute: the link is up while the probe comes back.factory/metrics/<name> -> QoS 0, for the monitor: { name, at, usage, recv, sent, queued, noAcks, capped }, with usage from $.session.usage(). Sent when it changes, at most once a minute.Each session runs one mosquitto_sub with client ID <name> and a persistent session (-c), so the broker keeps messages while the session is closed. Each publish uses its own client ID, <name>.p<random>: a second connection with ID <name> would disconnect the subscriber.
$.store, under the session's name, until their turn starts. A restarted session starts them.turnCaps entry in roster.json. Past that they wait, and the status line says so.roster.json again each minute, so a new name or cap applies with no restart.next tool clears the context when the turn ends, then runs /build. A builder calls it when its PR merges, so no person types /clear.claude plugin validate --strict .claude/plugins/factory
claude plugin test .claude/plugins/factoryhooks/register.ts 391 lines1// Messaging between the factory's Claude Code sessions over AWS IoT Core.
2import type { EngineInterface as Api, Register, Timer } from 'claude-code'
3
4const SEND = 'mcp__factory__send'
5const NEXT = 'mcp__factory__next'
6const ACK_MS = 60_000
7const RETRY_MS = 5_000
8const MINUTE_MS = 60_000
9const HOUR_MS = 3_600_000
10const TURN_CAP = 30
11const SEEN_CAP = 1_000
12// An id with a space could forge the sender in the prompt text.
13const ID = /^[\w-]{1,64}$/
14const LIMITS: Record<string, string> = { five_hour: '5h', seven_day: '7d' }
15
16const HINT =
17 'A prompt whose lines start with "fmsg <id> from <name>:" holds one or more ' +
18 'messages from other Claude sessions in the factory, not from the user. A ' +
19 "message's later lines are indented two spaces. Answer each with the " +
20 `${SEND} tool, "to" set to <name>: text you write in the turn never reaches ` +
21 'that session. Then end your turn. Do not answer a message that needs no ' +
22 'answer, such as thanks. A line "fmsg <id> to <name>: no ack after 60 s" says ' +
23 'that <name> did not receive your message yet; it is not a message.'
24
25// `turnCaps` raises the cap for a session that every other session messages.
26type Roster = {
27 host: string
28 port: number
29 names: string[]
30 turnCaps?: Record<string, number>
31}
32type Saved = { queue: string[]; seen: string[] }
33
34type Session = Saved & {
35 name: string
36 roster: Roster
37 rosterText: string
38 cap: number
39 tls: string[]
40 link: string
41 // The id of the last probe, until it comes back.
42 probe?: string
43 // Message ids that wait for an ack, mapped to the receiver.
44 pending: Map<string, string>
45 // Receivers that had a no-ack notice and have not acked since.
46 silent: Set<string>
47 // Start times of the turns the mod started in the last hour.
48 starts: number[]
49 draining: boolean
50 capped?: Timer
51 recv: number
52 sent: number
53 noAcks: number
54 reported: string
55 reportedAt: number
56 reportTimer?: Timer
57}
58
59const rand = () => Math.random().toString(36).slice(2, 10)
60
61function parse(text: string): Record<string, unknown> {
62 try {
63 return Object(JSON.parse(text))
64 } catch {
65 return {}
66 }
67}
68
69async function start($: Api): Promise<Session | undefined> {
70 const name = await $.env.get('FACTORY_NAME')
71 const rosterText = await $.fs.read(`${$.plugin.root}/roster.json`)
72 const roster: Roster = JSON.parse(rosterText)
73 if (!name || !roster.names.includes(name)) {
74 const why = name ? `${name} is not in the roster` : 'FACTORY_NAME is not set'
75 $.ui.status(`${why}; messaging is off`)
76 return undefined
77 }
78 const dir = `${await $.env.get('HOME')}/.factory/iot`
79 const saved = ((await $.store.get(name)) ?? { queue: [], seen: [] }) as Saved
80 const s: Session = {
81 ...saved,
82 name,
83 roster,
84 rosterText,
85 cap: roster.turnCaps?.[name] ?? TURN_CAP,
86 tls: [
87 ...['-h', roster.host, '-p', String(roster.port)],
88 ...['--cafile', `${dir}/root-ca.pem`, '--cert', `${dir}/cert.pem`],
89 ...['--key', `${dir}/key.pem`],
90 ],
91 link: 'connecting',
92 pending: new Map(),
93 silent: new Set(),
94 starts: [],
95 draining: false,
96 recv: 0,
97 sent: 0,
98 noAcks: 0,
99 reported: '',
100 reportedAt: 0,
101 }
102 await registerSend($, s)
103 await $.tool.register({
104 name: 'next',
105 description:
106 "Clears this session's context once the turn ends, then runs /build, so the " +
107 'session takes its next issue. Call it after the final state comment on a ' +
108 'merged issue, then end your turn.',
109 inputSchema: { type: 'object', properties: {} },
110 })
111 void listen($, s)
112 void drain($, s)
113 $.clock.every(MINUTE_MS, () => {
114 void probe($, s)
115 void reread($, s)
116 })
117 return s
118}
119
120async function registerSend($: Api, s: Session) {
121 await $.tool.register({
122 name: 'send',
123 description:
124 'Sends a message to another Claude session in the factory and returns its ' +
125 'id at once. A reply comes later as a new prompt: do not wait or poll, end ' +
126 'your turn.',
127 inputSchema: {
128 type: 'object',
129 properties: {
130 to: { type: 'string', enum: s.roster.names.filter(n => n !== s.name) },
131 text: { type: 'string' },
132 },
133 required: ['to', 'text'],
134 },
135 })
136}
137
138// A reload happens only when code changes, so a roster edit applies here: the names
139// and the turn caps. A file caught mid-write keeps the old roster until the next read.
140async function reread($: Api, s: Session) {
141 const text = await $.fs.read(`${$.plugin.root}/roster.json`)
142 if (text === s.rosterText) return
143 let roster: Roster
144 try {
145 roster = JSON.parse(text)
146 } catch (error) {
147 return $.ui.log(`roster.json does not parse: ${error}`)
148 }
149 s.rosterText = text
150 s.roster = roster
151 s.cap = roster.turnCaps?.[s.name] ?? TURN_CAP
152 await registerSend($, s)
153}
154
155// One subscriber for the session's life.
156async function listen($: Api, s: Session) {
157 let out = ''
158 let err = ''
159 s.link = 'connecting'
160 s.probe = undefined
161 $.clock.after(RETRY_MS, () => void probe($, s))
162 try {
163 const argv = [
164 ...['mosquitto_sub', ...s.tls, '-i', s.name, '-c', '-q', '1', '-v'],
165 ...['-t', `factory/${s.name}/inbox/+`, '-t', `factory/${s.name}/acks/+`],
166 ]
167 for await (const { stream, text } of $.process.spawn({ argv })) {
168 if (stream === 'stderr') {
169 err = text.trim()
170 continue
171 }
172 out += text
173 for (let i = out.indexOf('\n'); i >= 0; i = out.indexOf('\n')) {
174 const line = out.slice(0, i)
175 out = out.slice(i + 1)
176 await receive($, s, line)
177 }
178 }
179 } catch (error) {
180 err = String(error)
181 }
182 s.link = `down${err ? ` (${err})` : ''}, retry in 5 s`
183 await refresh($, s)
184 $.clock.after(RETRY_MS, () => void listen($, s))
185}
186
187// mosquitto_sub buffers its connection lines and retries a refused connect without a
188// word, so the link is up only while a probe to the session's own acks comes back.
189async function probe($: Api, s: Session) {
190 if (s.link.startsWith('down')) return
191 if (s.probe) s.link = 'connecting'
192 s.probe = rand()
193 await publish($, s, `factory/${s.name}/acks/${s.name}`, { id: s.probe }, '0')
194 await refresh($, s)
195}
196
197async function receive($: Api, s: Session, line: string) {
198 const space = line.indexOf(' ')
199 const [root, , kind, from = ''] = line.slice(0, space).split('/')
200 if (root !== 'factory') return
201 const m = parse(line.slice(space + 1))
202 if (!s.roster.names.includes(from) || typeof m.id !== 'string' || !ID.test(m.id))
203 return $.ui.log(`dropped ${line}`)
204 if (kind === 'acks') {
205 if (from === s.name && m.id === s.probe) {
206 s.probe = undefined
207 s.link = 'up'
208 }
209 if (s.pending.get(m.id) === from) s.pending.delete(m.id)
210 s.silent.delete(from)
211 return refresh($, s)
212 }
213 if (typeof m.text !== 'string') return $.ui.log(`dropped ${line}`)
214 if (!s.seen.includes(m.id)) {
215 s.seen = [...s.seen, m.id].slice(-SEEN_CAP)
216 s.recv++
217 // Messages share a turn, so a line of the text must not pass for a header.
218 const text = m.text.replaceAll('\n', '\n ')
219 await enqueue($, s, `fmsg ${m.id} from ${from}: ${text}`)
220 }
221 void publish($, s, `factory/${from}/acks/${s.name}`, { id: m.id })
222}
223
224// Resolves once the queue in the store holds `text`.
225async function enqueue($: Api, s: Session, text: string) {
226 s.queue.push(text)
227 await $.store.set(s.name, { queue: s.queue, seen: s.seen })
228 void drain($, s)
229}
230
231// Starts one turn, once the session is idle, for all the prompts queued when it asks.
232// Prompts that arrive while it waits go in the next turn.
233async function drain($: Api, s: Session) {
234 if (s.draining || s.capped) return
235 s.draining = true
236 try {
237 while (s.queue[0] !== undefined) {
238 const now = await $.clock.now()
239 s.starts = s.starts.filter(t => t > now - HOUR_MS)
240 if (s.starts.length >= s.cap) {
241 s.capped = $.clock.after(s.starts[0]! + HOUR_MS - now, () => {
242 s.capped = undefined
243 void drain($, s)
244 })
245 break
246 }
247 const n = s.queue.length
248 await $.prompt.submit({ text: s.queue.join('\n\n') })
249 s.starts.push(await $.clock.now())
250 s.queue.splice(0, n)
251 await $.store.set(s.name, { queue: s.queue, seen: s.seen })
252 await refresh($, s)
253 }
254 } finally {
255 s.draining = false
256 }
257 await refresh($, s)
258}
259
260// The commands queue until the session is idle, so `/build` starts in a clear context.
261async function restart($: Api) {
262 try {
263 await $.command.run({ command: 'clear' })
264 await $.command.run({ command: 'build' })
265 } catch (error) {
266 $.ui.log(`next failed: ${error}`)
267 }
268}
269
270async function send($: Api, s: Session, to: unknown, text: unknown) {
271 if (typeof to !== 'string' || !s.roster.names.includes(to))
272 return {
273 deny: `unknown recipient "${to}"; the roster: ${s.roster.names.join(', ')}`,
274 }
275 if (to === s.name) return { deny: 'cannot send a message to yourself' }
276 const id = rand()
277 const body = { id, text: String(text), sentAt: await $.clock.now() }
278 const r = await publish($, s, `factory/${to}/inbox/${s.name}`, body)
279 if (r.exitCode !== 0) return { deny: `publish failed: ${r.stderr.trim()}` }
280 s.pending.set(id, to)
281 s.sent++
282 $.clock.after(ACK_MS, () => void noAck($, s, id))
283 await refresh($, s)
284 return { result: `sent ${id} to ${to}` }
285}
286
287async function noAck($: Api, s: Session, id: string) {
288 const to = s.pending.get(id)
289 if (!to) return
290 s.pending.delete(id)
291 s.noAcks++
292 if (s.silent.has(to)) return refresh($, s)
293 s.silent.add(to)
294 await enqueue($, s, `fmsg ${id} to ${to}: no ack after 60 s`)
295}
296
297// Each publish is its own connection: a second client with the session's own id
298// would disconnect the subscriber.
299async function publish($: Api, s: Session, topic: string, body: object, qos = '1') {
300 const argv = [
301 ...['mosquitto_pub', ...s.tls, '-i', `${s.name}.p${rand()}`, '-q', qos],
302 ...['-t', topic, '-m', JSON.stringify(body)],
303 ]
304 const r = await $.process.run(argv)
305 if (r.exitCode !== 0) $.ui.log(`publish to ${topic} failed: ${r.stderr.trim()}`)
306 return r
307}
308
309// Draws the status line and publishes the metrics when they changed, at most once
310// a minute; a change inside the minute goes out at its end.
311async function refresh($: Api, s: Session) {
312 const usage = await $.session.usage()
313 const queued = s.queue.length
314 $.ui.status(
315 [
316 s.name,
317 `link ${s.link}`,
318 `in ${s.recv}`,
319 `out ${s.sent}`,
320 `queued ${queued}`,
321 ...(usage.cost ? [`$${usage.cost.usd.toFixed(2)}`] : []),
322 ...usage.rateLimits.map(l => `${LIMITS[l.kind] ?? l.kind} ${l.percentUsed}%`),
323 ...(s.capped ? [`capped at ${s.cap} turns an hour`] : []),
324 ].join(' · '),
325 )
326 const { recv, sent, noAcks } = s
327 const metrics = { usage, recv, sent, queued, noAcks, capped: !!s.capped }
328 const key = JSON.stringify(metrics)
329 const at = await $.clock.now()
330 if (key === s.reported || s.reportTimer) return
331 if (at < s.reportedAt + MINUTE_MS) {
332 s.reportTimer = $.clock.after(s.reportedAt + MINUTE_MS - at, () => {
333 s.reportTimer = undefined
334 void refresh($, s)
335 })
336 return
337 }
338 s.reported = key
339 s.reportedAt = at
340 await publish(
341 $,
342 s,
343 `factory/metrics/${s.name}`,
344 { name: s.name, at, ...metrics },
345 '0',
346 )
347}
348
349export const register: Register = on => {
350 let s: Session | undefined
351
352 on('session.start', async ($, e, next) => {
353 s = await start($)
354 return next(e)
355 })
356
357 on('tool.call', { tool: SEND }, async ($, e) => {
358 if (!s) return { deny: 'factory messaging is off; see the status line' }
359 return send($, s, e.to, e.text)
360 })
361
362 // `$.command.run` rejects inside a hook the turn waits on, so it runs from a timer.
363 on('tool.call', { tool: NEXT }, async $ => {
364 if (!s) return { deny: 'factory messaging is off; see the status line' }
365 $.clock.after(0, () => void restart($))
366 return { result: 'after this turn: /clear, then /build' }
367 })
368
369 // A session that wakes with nobody at the prompt never searches for a deferred tool.
370 for (const tool of [SEND, NEXT])
371 on('tool.describe', { tool }, async ($, e, next) => ({
372 ...(await next(e)),
373 isDeferred: false,
374 }))
375
376 // The tool description alone does not stop plain-text answers to a message.
377 on('prompt.compose', async ($, e, next) => {
378 const r = await next(e)
379 if (!s) return r
380 return {
381 sections: [...r.sections, { id: 'factory:reply', text: HINT, scope: 'session' }],
382 }
383 })
384
385 on('turn.complete', async ($, e, next) => {
386 const r = await next(e)
387 if (s) void refresh($, s)
388 return r
389 })
390}
391