SLOPSHOPPER

factory

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

newguardstatusprompttoolprocess
v0.1.0no licenseupdated 2026-10-09synnaxlabs/foundation/.claude/plugins/factory
A shopper browsing a rack in a slop shop
README

factory

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).

Start a session

~/.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.

Broker

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.

Behavior

  • Messages wait in $.store, under the session's name, until their turn starts. A restarted session starts them.
  • A repeated message id gets an ack and no second turn.
  • When a message gets no ack in 60 s, the sender gets one notice turn for that receiver, and no other until the receiver acks again.
  • One turn takes all the messages waiting when the session goes idle. A message's later lines are indented two spaces, so no line passes for a header.
  • Messages start at most 30 turns an hour, or the session's turnCaps entry in roster.json. Past that they wait, and the status line says so.
  • A session reads roster.json again each minute, so a new name or cap applies with no restart.
  • The status line shows the name, the link, messages in and out, the queue, the cost, and the rate limits.
  • The 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.

Check it

claude plugin validate --strict .claude/plugins/factory
claude plugin test .claude/plugins/factory
Source 1 files
hooks/register.ts 391 lines
1// 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