Follows a command that prints events as JSON lines, shows status events on the status line, and injects flagged messages as a prompt between turns.

A Claude Code mod that follows an event source for the whole session.
The source is a command that prints one JSON record per line. The mod shows status records on its status line. It injects message records flagged inject: true into the session as a prompt, between turns only. After the prompt entered, the mod runs an ack command once per injected record.
The defaults follow a crabswarm chat room with crabswarm chat follow, inject the messages that mention this session, and mark each one read with crabswarm chat read --skip <seq>. SKILL.md describes the details.
\n and flushed as soon as it is written.type and a string message.{"type":"status","message":"..."} sets the mod's status line to message. An empty message clears it. The mod clears the status line itself 5 seconds after the latest entry.{"type":"message","message":"...","inject":true} is injected. A message whose inject is anything but true is ignored.{"type":"status","message":"following room-a"}
{"type":"message","message":"alice: please review the plan","inject":true,"seq":42}
<prefix> <message>, in arrival order. A line equal to an earlier line of the same prompt is left out.| Option | Default | Meaning |
|---|---|---|
source | ["crabswarm", "chat", "follow"] | The source command and its arguments. Empty turns the mod off. |
mapper | [] | A command run once per source line, with the line on standard input. Each non-empty line it prints is a record. A non-zero exit drops the line. Empty reads the source lines as records. |
prefix | "[crabswarm chat]" | The text before each injected message, separated by a space. Empty puts nothing. |
ack | ["crabswarm", "chat", "read", "--skip", "{seq}"] | The command run once per injected record. An item written exactly {name} is replaced by the record's top-level field name. Empty acks nothing. |
requireEnv | ["CRABSWARM_CHAT_TOKEN", "CMDMAN_CMD_ID"] | The mod runs only when the session's environment sets at least one of these. Empty always runs. |
source, mapper, ack and requireEnv are lists of strings ("type": "string", "multiple": true in plugin.json); prefix is a string.
["jq", "-c", "{type: \"message\", message: .text, inject: true, id}"].To override them, add pluginConfigs.<plugin id>.options to one of these settings files:
~/.claude/settings.json ($CLAUDE_CONFIG_DIR/settings.json when that variable is set);claude --settings <file>, for one session;Claude Code does not read pluginConfigs from a project's .claude/settings.json or .claude/settings.local.json.
The plugin id depends on how the mod was loaded:
handle-events@skills-dir when apm installed it into ~/.claude/skills/;handle-events or handle-events@inline under claude --plugin-dir.{
"pluginConfigs": {
"handle-events@skills-dir": {
"options": {
"source": ["my-tool", "events", "--json"],
"prefix": "[my-tool]",
"ack": ["my-tool", "events", "ack", "{id}"],
"requireEnv": []
}
}
}
}
claude --debug logs the keys it looked for when it finds none.
Every command runs in the session's environment and working directory.
crabswarm chat read --skip {seq}, also marks earlier unread mentions read.--skip moves it through seq.printenv, and the source, mapper and ack commands, on PATH of the Claude Code process.Install it globally with apm:
apm install -g ngicks/crabswarm/mods/handle-events
apm copies this folder to ~/.claude/skills/handle-events/. Claude Code adopts it as the plugin handle-events@skills-dir in the next session. Run /reload-plugins to load it into a running session.
For a one-off session from a checkout:
claude --plugin-dir ./mods/handle-events
claude plugin validate mods/handle-events
claude plugin test mods/handle-events
Claude Code writes the API types to .claude-plugin/types/ whenever it loads the mod from this folder. After that, tsc -p mods/handle-events type-checks it.
claude.claude plugin validate warns on a plugin name that reads as one of Anthropic's own.hooks/hooks.json carries a description key.hooks/*.json whose values are all lists into settings.json as settings hooks.modules stays out of the user's settings..claude-plugin/plugin.json makes apm treat this folder as a Claude plugin and deploy it whole into ~/.claude/skills/.hooks/register.tsx 129 lines1import type { EngineInterface, PluginOptions, Register, Timer } from 'claude-code'
2import { ackRecords } from './ack'
3import { type Run, type Spawn, anyEnvSet } from './host'
4import { actionOf } from './protocol'
5import { type InjectQueue, injectQueue } from './queue'
6import { runSource } from './source'
7
8// strings reads a `multiple` string option. The engine validates options
9// against the manifest before the module loads, so anything else is only an
10// unset option.
11function strings(v: PluginOptions[string] | undefined): readonly string[] {
12 return Array.isArray(v) ? v.filter(s => s !== '') : []
13}
14
15// A failure that repeats (a source that keeps dying, a source printing lines
16// the mod cannot read, an ack command that keeps failing) would otherwise put
17// a line in the debug log each time, so the first failure of a run is logged,
18// then one in every WARN_EVERY.
19const WARN_EVERY = 100
20
21// A status line entry is cleared this long after it was shown. A source
22// reports a moment rather than a lasting state, and an entry left in place
23// reads as current long after it stopped being true.
24const STATUS_HOLD_MS = 5000
25
26function throttledLog($: EngineInterface, counted: string) {
27 let failures = 0
28 return {
29 fail(message: string) {
30 failures++
31 if (failures === 1 || failures % WARN_EVERY === 0) {
32 $.ui.log(`handle-events: ${message} (${counted}: ${failures})`, { to: 'debug' })
33 }
34 },
35 ok() {
36 failures = 0
37 },
38 }
39}
40
41export const register: Register = (on, options) => {
42 const sourceArgv = strings(options.source)
43 const mapperArgv = strings(options.mapper)
44 const prefix = typeof options.prefix === 'string' ? options.prefix : ''
45 const ackArgv = strings(options.ack)
46 const requireEnv = strings(options.requireEnv)
47
48 let queue: InjectQueue | undefined
49 // session.start may fire again in the same load; one source is enough.
50 let started = false
51
52 on('session.start', async ($, e, next) => {
53 const result = await next(e)
54 if (started || sourceArgv.length === 0) return result
55 started = true
56 const run: Run = (argv, init) => $.process.run(argv, init)
57 // A session the source has no identity for stays quiet instead of
58 // restarting a failing source forever.
59 if (!(await anyEnvSet(run, requireEnv))) return result
60
61 const drops = throttledLog($, 'drops')
62 const ackFailures = throttledLog($, 'ack failures')
63 const restarts = throttledLog($, 'restarts')
64 let clearing: Timer | undefined
65 const status = (text: string | undefined) => {
66 clearing?.cancel()
67 clearing = text === undefined ? undefined : $.clock.after(STATUS_HOLD_MS, () => void $.ui.status(undefined))
68 void $.ui.status(text)
69 }
70
71 const q = injectQueue({
72 prefix,
73 submit: async text => {
74 const r = await $.prompt.submit({ text })
75 if (r.drop !== undefined) throw new Error(`the prompt was dropped: ${r.drop}`)
76 },
77 onSubmitted: records =>
78 ackRecords(records, { run, argv: ackArgv, onError: ackFailures.fail, onAcked: ackFailures.ok }),
79 onError: drops.fail,
80 })
81 queue = q
82
83 const spawn: Spawn = argv => $.process.spawn({ argv })
84 void runSource({
85 argv: sourceArgv,
86 mapper: mapperArgv,
87 run,
88 spawn,
89 now: () => $.clock.now(),
90 wait: ms => new Promise<void>(resolve => void $.clock.after(ms, () => resolve())),
91 onRecord: record => {
92 drops.ok()
93 switch (actionOf(record)) {
94 case 'status':
95 status(record.message === '' ? undefined : record.message)
96 break
97 case 'inject':
98 q.push(record)
99 break
100 }
101 },
102 onDrop: drops.fail,
103 onRestart: (reason, delayMs) => {
104 restarts.fail(reason)
105 status(`handle-events: ${reason}; restarting in ${Math.ceil(delayMs / 1000)}s`)
106 },
107 onRecovered: () => {
108 restarts.ok()
109 status(undefined)
110 },
111 })
112 return result
113 })
114
115 // A subagent's run raises no turn.start, so every one is the main loop's.
116 on('turn.start', async (_$, e, next) => {
117 queue?.turnStarted()
118 return next(e)
119 })
120
121 on('turn.complete', async (_$, e, next) => {
122 const result = await next(e)
123 // The flush is not awaited here: the prompt it submits starts its turn
124 // only once the session is idle, which is after this hook returns.
125 if (e.agentId === undefined) queue?.turnCompleted()
126 return result
127 })
128}
129hooks/ack.ts 60 lines1import { type Run, withStderr } from './host'
2import type { EventRecord } from './protocol'
3
4const ACK_TIMEOUT_MS = 10_000
5
6// An argv item that is exactly `{name}` stands for the record's field `name`.
7const PLACEHOLDER = /^\{([^{}]+)\}$/
8
9// fillArgv replaces each placeholder item of argv with the record's top-level
10// field of that name, written as text. It answers the name of the first field
11// that is missing or is not a string, a number or a boolean instead.
12export function fillArgv(argv: readonly string[], record: EventRecord): { argv: string[] } | { missing: string } {
13 const filled: string[] = []
14 for (const item of argv) {
15 const name = PLACEHOLDER.exec(item)?.[1]
16 if (name === undefined) {
17 filled.push(item)
18 continue
19 }
20 const v = Object.hasOwn(record, name) ? record[name] : undefined
21 if (typeof v !== 'string' && typeof v !== 'number' && typeof v !== 'boolean') return { missing: name }
22 filled.push(String(v))
23 }
24 return { argv: filled }
25}
26
27export type AckOptions = {
28 run: Run
29 argv: readonly string[]
30 onError?: (message: string) => void
31 onAcked?: () => void
32}
33
34// ackRecords runs the ack command once per record, in order, each run after
35// the one before it ended. A record that lacks a field the command names is
36// skipped; a run that exits non-zero is reported and the next one still runs.
37// An empty argv acks nothing.
38export async function ackRecords(records: readonly EventRecord[], opts: AckOptions): Promise<void> {
39 if (opts.argv.length === 0) return
40 for (const record of records) {
41 const filled = fillArgv(opts.argv, record)
42 if ('missing' in filled) {
43 opts.onError?.(`ack skipped: the record has no string, number or boolean "${filled.missing}": ${record.message}`)
44 continue
45 }
46 const command = filled.argv.join(' ')
47 try {
48 const r = await opts.run(filled.argv, { timeoutMs: ACK_TIMEOUT_MS })
49 if (r.exitCode !== 0) {
50 opts.onError?.(withStderr(`${command} exited ${r.exitCode}`, r.stderr))
51 continue
52 }
53 } catch (err) {
54 opts.onError?.(`${command}: ${err instanceof Error ? err.message : String(err)}`)
55 continue
56 }
57 opts.onAcked?.()
58 }
59}
60hooks/host.ts 43 lines1// The engine calls the helpers make. The engine refuses a $ passed across an
2// import, so register.tsx hands in closures over its own $ instead.
3export type Run = (
4 argv: readonly string[],
5 init?: { stdin?: string; timeoutMs?: number },
6) => Promise<{ exitCode: number; stdout: string; stderr: string }>
7
8export type SpawnChunk = { stream: 'stdout' | 'stderr'; text: string }
9export type SpawnEnd = { code: number | null; signal: string | null }
10
11// Spawn starts argv and streams its output; the iterator's return value is
12// how the child ended. Leaving the iteration early kills the child.
13export type Spawn = (argv: readonly string[]) => AsyncIterator<SpawnChunk, SpawnEnd | void>
14
15export type Now = () => Promise<number>
16
17// Wait resolves after ms. register.tsx builds it on $.clock.after rather than
18// $.clock.sleep: a sleep is charged to the budget of the hook that started
19// it, while a timer runs until the module reloads.
20export type Wait = (ms: number) => Promise<void>
21
22// withStderr appends what a command wrote to stderr to a message about it.
23export function withStderr(message: string, stderr: string): string {
24 const s = stderr.trim()
25 return s === '' ? message : `${message}: ${s}`
26}
27
28const PRINTENV_TIMEOUT_MS = 5000
29
30// anyEnvSet reports whether the session's environment sets at least one of
31// names to a non-empty value; no names at all is a yes.
32//
33// It asks printenv rather than $.env.get: the engine takes only a string
34// literal as a variable name there, and these names come from the options.
35export async function anyEnvSet(run: Run, names: readonly string[]): Promise<boolean> {
36 if (names.length === 0) return true
37 for (const name of names) {
38 const r = await run(['printenv', name], { timeoutMs: PRINTENV_TIMEOUT_MS })
39 if (r.exitCode === 0 && r.stdout.trim() !== '') return true
40 }
41 return false
42}
43hooks/protocol.ts 42 lines1// One record of the protocol: a JSON object carrying at least a string `type`
2// and a string `message`. Every other field rides along untouched; ack reads
3// them by name.
4export type EventRecord = {
5 readonly type: string
6 readonly message: string
7 readonly [field: string]: unknown
8}
9
10// What the mod does with a record.
11export type Action = 'status' | 'inject'
12
13// parseRecord reads one line of a source or a mapper as a record. It answers
14// undefined for anything that is not a JSON object with a string type and a
15// string message.
16export function parseRecord(line: string): EventRecord | undefined {
17 let v: unknown
18 try {
19 v = JSON.parse(line)
20 } catch {
21 return undefined
22 }
23 if (typeof v !== 'object' || v === null || Array.isArray(v)) return undefined
24 const r = v as Record<string, unknown>
25 if (typeof r.type !== 'string' || typeof r.message !== 'string') return undefined
26 return r as EventRecord
27}
28
29// actionOf picks what a record asks for. A type the mod does not know, and a
30// message that does not set inject to true, ask for nothing: the protocol lets
31// a source carry more than this mod handles.
32export function actionOf(record: EventRecord): Action | undefined {
33 switch (record.type) {
34 case 'status':
35 return 'status'
36 case 'message':
37 return record.inject === true ? 'inject' : undefined
38 default:
39 return undefined
40 }
41}
42hooks/queue.ts 101 lines1import type { EventRecord } from './protocol'
2
3// Submit hands one prompt to the session and rejects when it did not enter.
4export type Submit = (text: string) => Promise<void>
5
6export type QueueOptions = {
7 prefix: string
8 submit: Submit
9 // onSubmitted runs once the prompt holding the records entered.
10 onSubmitted?: (records: readonly EventRecord[]) => Promise<void>
11 onError?: (message: string) => void
12}
13
14// promptText writes records as one prompt: one line each, the prefix before
15// the message, in the order given. A line equal to an earlier one is left out.
16export function promptText(prefix: string, records: readonly EventRecord[]): string {
17 const seen = new Set<string>()
18 const lines: string[] = []
19 for (const r of records) {
20 const line = prefix === '' ? r.message : `${prefix} ${r.message}`
21 if (seen.has(line)) continue
22 seen.add(line)
23 lines.push(line)
24 }
25 return lines.join('\n')
26}
27
28export type InjectQueue = {
29 push: (record: EventRecord) => void
30 turnStarted: () => void
31 turnCompleted: () => void
32}
33
34// injectQueue holds records until the main loop is idle, then submits every
35// record it holds as one prompt. A record that arrives while a turn runs, or
36// while a prompt is on its way in, waits for the next turn to end.
37export function injectQueue(opts: QueueOptions): InjectQueue {
38 let pending: EventRecord[] = []
39 let turnRunning = false
40 // busy covers a submit and the onSubmitted after it, so two records that
41 // arrive a moment apart while idle become one prompt.
42 let busy = false
43 // starts counts turn.start events, so a submit that resolves can tell
44 // whether its turn already announced itself.
45 let starts = 0
46
47 const fail = (err: unknown) => opts.onError?.(errorText(err))
48
49 async function deliver(batch: readonly EventRecord[]) {
50 const startsBefore = starts
51 try {
52 await opts.submit(promptText(opts.prefix, batch))
53 } catch (err) {
54 opts.onError?.(`${batch.length} record(s) not injected: ${errorText(err)}`)
55 return
56 }
57 // $.prompt.submit resolves as the prompt's own turn starts. Until that
58 // turn's turn.start arrives, the queue counts the turn as running itself.
59 if (starts === startsBefore) turnRunning = true
60 try {
61 await opts.onSubmitted?.(batch)
62 } catch (err) {
63 fail(err)
64 }
65 }
66
67 async function flush(): Promise<void> {
68 if (busy || turnRunning || pending.length === 0) return
69 const batch = pending
70 pending = []
71 busy = true
72 try {
73 await deliver(batch)
74 } finally {
75 busy = false
76 }
77 await flush()
78 }
79
80 const kick = () => void flush().catch(fail)
81
82 return {
83 push(record) {
84 pending.push(record)
85 kick()
86 },
87 turnStarted() {
88 starts++
89 turnRunning = true
90 },
91 turnCompleted() {
92 turnRunning = false
93 kick()
94 },
95 }
96}
97
98function errorText(err: unknown): string {
99 return err instanceof Error ? err.message : String(err)
100}
101hooks/source.ts 154 lines1import { type Now, type Run, type Spawn, type SpawnEnd, type Wait, withStderr } from './host'
2import { lineSplitter, splitLines } from './lines'
3import { type EventRecord, parseRecord } from './protocol'
4
5// The pause before restarting a source starts here and doubles after each
6// restart up to BACKOFF_MAX_MS.
7export const BACKOFF_MIN_MS = 1000
8export const BACKOFF_MAX_MS = 30_000
9
10// A source that ran this long before it exited counts as healthy, so its next
11// restart waits the shortest pause again. A valid record counts the same.
12export const HEALTHY_RUN_MS = 30_000
13
14const MAPPER_TIMEOUT_MS = 10_000
15
16// How much of a dropped line a log message quotes.
17const QUOTE_MAX = 200
18
19export type RecordOptions = {
20 mapper: readonly string[]
21 run: Run
22 onDrop?: (message: string) => void
23}
24
25// recordsOf turns one line of a source into the records it carries. Without a
26// mapper the line is the record. With one, the mapper runs with the line on
27// its standard input, and each non-empty line it prints is a record; a mapper
28// that exits non-zero drops the line. A line that does not parse as a record
29// is dropped. Each drop is reported to onDrop.
30export async function recordsOf(line: string, opts: RecordOptions): Promise<EventRecord[]> {
31 if (line.trim() === '') return []
32 let lines = [line]
33 if (opts.mapper.length > 0) {
34 let r: Awaited<ReturnType<Run>>
35 try {
36 r = await opts.run(opts.mapper, { stdin: line + '\n', timeoutMs: MAPPER_TIMEOUT_MS })
37 } catch (err) {
38 opts.onDrop?.(`mapper failed on ${quote(line)}: ${errorText(err)}`)
39 return []
40 }
41 if (r.exitCode !== 0) {
42 opts.onDrop?.(withStderr(`mapper exited ${r.exitCode} on ${quote(line)}`, r.stderr))
43 return []
44 }
45 lines = splitLines(r.stdout)
46 }
47 const records: EventRecord[] = []
48 for (const l of lines) {
49 if (l.trim() === '') continue
50 const record = parseRecord(l)
51 if (record === undefined) {
52 opts.onDrop?.(`not a record: ${quote(l)}`)
53 continue
54 }
55 records.push(record)
56 }
57 return records
58}
59
60export type SourceOptions = RecordOptions & {
61 argv: readonly string[]
62 spawn: Spawn
63 now: Now
64 wait: Wait
65 onRecord: (record: EventRecord) => void
66 // onRestart is told why the source ended and how long the restart waits.
67 onRestart?: (reason: string, delayMs: number) => void
68 // onRecovered is called on the first valid record after a restart.
69 onRecovered?: () => void
70 signal?: AbortSignal
71}
72
73// runSource runs the source for as long as signal allows, restarting it each
74// time it ends. Records reach onRecord in the order the source printed them:
75// a line waits for the mapper run of the line before it.
76export async function runSource(opts: SourceOptions): Promise<void> {
77 let delay = BACKOFF_MIN_MS
78 let restarting = false
79 while (!opts.signal?.aborted) {
80 const startedAt = await opts.now()
81 let healthy = false
82 const reason = await runOnce(opts, record => {
83 healthy = true
84 if (restarting) {
85 restarting = false
86 opts.onRecovered?.()
87 }
88 opts.onRecord(record)
89 })
90 if (opts.signal?.aborted) return
91 if (healthy || (await opts.now()) - startedAt >= HEALTHY_RUN_MS) delay = BACKOFF_MIN_MS
92 restarting = true
93 opts.onRestart?.(reason, delay)
94 await opts.wait(delay)
95 delay = Math.min(delay * 2, BACKOFF_MAX_MS)
96 }
97}
98
99// runOnce runs the source until it ends and answers why it ended.
100async function runOnce(opts: SourceOptions, onRecord: (record: EventRecord) => void): Promise<string> {
101 const lines = lineSplitter()
102 let lastStderr = ''
103 const handle = async (line: string) => {
104 for (const record of await recordsOf(line, opts)) onRecord(record)
105 }
106 let end: SpawnEnd | void
107 try {
108 const child = opts.spawn(opts.argv)
109 for (;;) {
110 const step = await child.next()
111 if (step.done) {
112 end = step.value
113 break
114 }
115 if (step.value.stream === 'stderr') {
116 lastStderr = lastLine(step.value.text) ?? lastStderr
117 continue
118 }
119 for (const line of lines.push(step.value.text)) await handle(line)
120 if (opts.signal?.aborted) {
121 await child.return?.()
122 return 'stopped'
123 }
124 }
125 } catch (err) {
126 return `source failed: ${errorText(err)}`
127 }
128 const rest = lines.end()
129 if (rest !== undefined) await handle(rest)
130 return endReason(end, lastStderr)
131}
132
133function endReason(end: SpawnEnd | void, lastStderr: string): string {
134 let reason = 'source ended'
135 if (end?.code !== null && end?.code !== undefined) reason = `source exited with code ${end.code}`
136 else if (end?.signal) reason = `source ended by ${end.signal}`
137 return withStderr(reason, lastStderr)
138}
139
140function lastLine(text: string): string | undefined {
141 return splitLines(text)
142 .map(l => l.trim())
143 .filter(l => l !== '')
144 .at(-1)
145}
146
147function quote(line: string): string {
148 return JSON.stringify(line.length > QUOTE_MAX ? line.slice(0, QUOTE_MAX) + '...' : line)
149}
150
151function errorText(err: unknown): string {
152 return err instanceof Error ? err.message : String(err)
153}
154hooks/lines.ts 36 lines1// lineSplitter cuts a child's output into lines. A piece of output ends
2// wherever the child's write did, so a line may span pieces and a piece may
3// hold several lines; push keeps the unfinished tail for the next piece.
4//
5// A line loses its "\n" and a "\r" before it.
6export function lineSplitter() {
7 let tail = ''
8 return {
9 push(text: string): string[] {
10 const parts = (tail + text).split('\n')
11 tail = parts.pop() ?? ''
12 return parts.map(trimCR)
13 },
14 // end answers the text after the last "\n", once: what a child wrote
15 // before it exited without ending its last line.
16 end(): string | undefined {
17 const rest = tail
18 tail = ''
19 return rest === '' ? undefined : trimCR(rest)
20 },
21 }
22}
23
24// splitLines cuts a whole output into its lines, the last one whether or not
25// it ends in "\n".
26export function splitLines(text: string): string[] {
27 const s = lineSplitter()
28 const lines = s.push(text)
29 const rest = s.end()
30 return rest === undefined ? lines : [...lines, rest]
31}
32
33function trimCR(line: string): string {
34 return line.endsWith('\r') ? line.slice(0, -1) : line
35}
36