System Bus Worker
SkillCloud & infraDevelop, deploy, and debug the system-bus worker, joelclaw's 110+ Inngest durable function engine, webhook gateway, and observability pipeline. Triggers on 'add a function', 'new inngest function', 'system-bus', 'worker', 'add a webhook', 'deploy worker', 'restart worker', 'function failed', 'worker not working', 'register functions', or any task involving Inngest function development, webhook providers, or worker operations.
Available today. Use it from your connected AI after setup.
No other account needed.
Add ahel to your AI once: Claude, ChatGPT, Cursor, Claude Code or Codex. Then ask it to use this.
Then ask your AI: use the System Bus Worker skill
What this skill tells your AI
The instructions your AI receives, as published by joelhooks/joelclaw in skills/system-bus/SKILL.md and read by ahel’s review.
The system-bus worker (@joelclaw/system-bus) is joelclaw's event-driven backbone — 110+ Inngest durable functions, webhook ingestion, and observability. It runs as a Hono HTTP server registered with the self-hosted Inngest instance.
Architecture
packages/system-bus/
├── src/
│ ├── serve.ts # Hono server, Inngest registration, health endpoint
│ ├── inngest/
│ │ ├── client.ts # Inngest client + event type definitions
│ │ ├── middleware/ # Gateway injection, dependency injection
│ │ └── functions/
│ │ ├── index.ts # Combined exports
│ │ ├── index.host.ts # Functions for host-role worker (local Mac)
│ │ ├── index.cluster.ts # Functions for cluster-role worker (k8s)
│ │ └── <function-name>.ts # Individual functions
│ ├── lib/ # Shared utilities
│ │ ├── inference.ts # LLM calls via pi (CANONICAL — always use this)
│ │ ├── redis.ts # Redis client helper
│ │ ├── typesense.ts # Typesense client
│ │ ├── convex-content-sync.ts # Convex upsert for content pipeline
│ │ ├── pi-output.ts # pi JSON output / usage parsing
│ │ └── ...
│ ├── observability/
│ │ └── emit.ts # OTEL event emission
│ ├── webhooks/
│ │ ├── server.ts # Webhook router (mounted at /webhooks)
│ │ ├── types.ts # Provider interface
│ │ └── providers/ # Per-service webhook handlers
│ │ ├── front.ts
│ │ ├── github.ts
│ │ ├── vercel.ts
│ │ ├── todoist.ts
│ │ ├── mux.ts
│ │ └── joelclaw.ts
│ └── memory/ # Memory pipeline components
├── scripts/
│ └── sync-content-to-convex.ts # Manual full Convex sync
└── package.json
Worker Roles
Two deployment modes controlled by WORKER_ROLE env var:
| Role | Where | Functions |
|---|---|---|
host | Local Mac Mini via Talon supervisor (optional) | Agent loops, heartbeat checks, memory pipeline, content sync, video ingest, book download — anything needing local filesystem, pi CLI, or docker |
cluster | k8s pod (GHCR image) | Webhooks (Front, GitHub, Vercel, Todoist, Mux), approvals, notifications, Slack backfill — stateless, network-only |
Functions are split between index.host.ts and index.cluster.ts. The combined index.ts exports everything for tooling/tests.
Deployment Model
- Source of truth:
~/Code/joelhooks/joelclaw/packages/system-bus/ - Running host worker:
worker-supervisorprocess runningbun run src/serve.tsfrom~/Code/joelhooks/joelclaw/packages/system-bus/- verify with
lsof -iTCP:3111 -sTCP:LISTEN -n -Pandlsof -p <pid> | awk '$4=="cwd"{print}' - legacy clone
~/Code/system-bus-worker/may still exist, but it is not the active host worker when port 3111's cwd points at the monorepo
- verify with
- Cluster runtime:
system-bus-workerDeployment in the Talos/Colima k8s cluster for cluster-role workloads - Cluster deploy path:
~/Code/joelhooks/joelclaw/k8s/publish-system-bus-worker.sh
Host function rollout reality
Host worker registration supports an explicit INNGEST_SERVE_HOST override in ~/.config/system-bus.env. Set INNGEST_SERVE_HOST=connect to suppress SDK callback URL advertising when the self-hosted Inngest pod cannot route to the host network. Do not leave it pointing at a stale Tailscale or Docker-host address unless you have proven the Inngest pod can open that address from inside k8s; otherwise every host function fails with Unable to reach SDK URL.
After changing packages/system-bus/src/inngest/functions/* that run on the host worker:
- commit + push the monorepo change to
origin - confirm the live worker cwd:
pid=$(lsof -tiTCP:3111 -sTCP:LISTEN); lsof -p "$pid" | awk '$4=="cwd"{print}' - if cwd is
~/Code/joelhooks/joelclaw/packages/system-bus, kill the Bun worker PID and letworker-supervisorrespawn it - if cwd is the legacy
~/Code/system-bus-worker, inspect its status and divergence, preserve both sides, and use the current deployment source; never reset the legacy clone to discard work - verify
curl http://127.0.0.1:3111/shows functions andjoelclaw functionsreturns >0
The stale failure mode: a host worker can keep running old source for days. In that state, OTEL may show behavior that current monorepo code has already fixed. Always verify the live port-3111 process cwd and start time before debugging source that "should" already be deployed.
Queue pilot flags are evaluated inside the live worker process, not your shell. If a host-worker emitter like discovery-capture or /webhooks/github should switch to queue mode, put the flag in ~/.config/system-bus.env, then kickstart the worker and PUT-sync /api/inngest. Ad-hoc shell env only affects CLI-local emitters.
Queue triage flags follow the same rule. Current bounded admission contract:
QUEUE_TRIAGE_MODE=off|shadow|enforcesets the base triage mode.QUEUE_TRIAGE_FAMILIES=discovery,content,subscriptions,github(or exact event names) chooses which queue families participate at all.QUEUE_TRIAGE_ENFORCE_FAMILIES=discovery,githubis the narrow Story 4 override that upgrades onlydiscovery/notedandgithub/workflow_run.completedinto enforce.- Any non-eligible family is clamped back to
shadoweven if someone sets globalQUEUE_TRIAGE_MODE=enforce. - Handler routing always stays registry-derived; triage may only shape bounded admission fields.
content/updated is the odd one out: its ingress comes from the launchd watcher com.joel.content-sync-watcher, not from a worker-local function. The canonical watcher source now belongs in infra/launchd/com.joel.content-sync-watcher.plist plus scripts/content-sync-watcher.sh, and the script reads ~/.config/system-bus.env on each trigger so QUEUE_PILOTS=content can switch between joelclaw queue emit and legacy joelclaw send without hand-editing the live plist.
For Story 5 soak work, start from joelclaw jobs status for the first operator glance, then drop into joelclaw queue stats before spelunking raw OTEL or Redis. jobs status is the transitional runtime surface that rolls queue / Restate / Dkron / Inngest into one JSON snapshot without forcing the operator to jump between commands just to learn whether the substrate is healthy enough to take work. queue stats remains the queue-specific summary for Restate drainer health and queue triage behavior: it rolls up recent queue.dispatch.started|completed|failed telemetry plus the queue.triage.* lifecycle into live depth, terminal success/failure counts, waitTimeMs percentiles, dispatch-duration percentiles, fallback reasons, disagreement counts, applied-vs-suggested deltas, route mismatches, family rollups, and recent mismatch/fallback samples. Use joelclaw queue stats --since <iso|ms> when you need to anchor the sample to a known-clean point such as a supervised queue.drainer.started after rollout. Honest gotcha from the live Story 5 cleanup follow-through: global depth can lie because of unrelated historical backlog, so judge the supervised sample first with the anchored triage/dispatch window plus joelclaw queue inspect <stream-id> / joelclaw queue list --limit <n> on the fresh sample IDs. If old residue survives a supervised com.joel.restate-worker restart, clear it with a bounded @joelclaw/queue ack() pass only after confirming zero pending leases and an age filter on the orphaned stream IDs. If that command is broken or misleading, fix it before widening queue cutovers.
For ADR-0217 Phase 3 Story 2-4, the operator surfaces are joelclaw queue observe, joelclaw queue pause, joelclaw queue resume, and joelclaw queue control status. queue observe still answers “what would Sonnet do right now?” in dry-run, but its snapshot.control.activePauses plus top-level control block now reflect the shipped deterministic control plane: active pauses, expirations, and recent queue.control.* OTEL. It short-circuits to a deterministic noop when the backlog is fully explained by fresh active manual pauses and no recent failures suggest downstream trouble, and it now also short-circuits to deterministic resume_family when queued work is entirely held behind a settled observer pause with no fresh drainer/triage failures. That avoids wasting a 60s Sonnet call on an obvious operator hold state and stops the observer from mistaking its own pause for permanent downstream failure. queue pause / queue resume are the bounded manual apply path before any automatic observer mutation. queue control status is the direct operator truth source for active manual controls and recent queue.control.applied|expired|rejected events.
ADR-0217 Phase 3 Story 4 now has a live host-worker runtime in packages/system-bus/src/inngest/functions/queue-observer.ts. Durable cadence belongs in Inngest, not the gateway daemon: the cron controller stays on queue/observer, while manual queue/observer.requested probes now run through a separate queue/observer-requested function so operator requests do not sit behind the cron pass. Runtime flags live in ~/.config/system-bus.env and require the usual host-worker restart + PUT /api/inngest:
QUEUE_OBSERVER_MODE=off|dry-run|enforceQUEUE_OBSERVER_FAMILIES=discovery,content,subscriptions,githubQUEUE_OBSERVER_AUTO_FAMILIES=contentQUEUE_OBSERVER_INTERVAL_SECONDS(currently clamped to a 60s minimum on the durable cron path)
Both paths build the same bounded snapshot and use the shared Sonnet observer contract, but only the cron controller may auto-apply queue actions; manual probes are read-only even when the configured mode is enforce. The shared observer also short-circuits deterministic noops for both truly empty queues and empty queues that only still have active pauses hanging around, so the cron path does not waste a full model call on obvious nothing-to-do snapshots. Idle empty snapshots with no recent drainer failures or triage trouble now report downstreamState=healthy instead of inheriting a noisy degraded label from stale throughput/latency history. Settled observer-held backlogs now get the same treatment: once a queued family is already behind an observer pause for at least one cadence and no fresh drainer/triage failures exist, the shared observer deterministically emits resume_family instead of reading its own hold as proof that downstream is still down. The manual probe function also uses singleton-skip semantics so repeated operator requests do not pile up stale queued runs. Live canaries also exposed a prompt-contract gotcha: Sonnet was happy to emit escalate.reason, which blew up the strict schema. The shared parser now normalizes that legacy shape to { severity: "warn", message }, and the prompt explicitly tells Sonnet to prefer pause_family/resume_family over batch_family for the content/updated pilot when downstream is crook.
Current operator truth after the latest live canaries: dry-run is earned and the hardened enforce path has now completed one full automatic observer cycle on content/updated. The first supervised canary anchored at since=1772981290859 booted Restate out long enough to build a 30-item content backlog and auto-applied pause_family on snapshot cca656f7-a9ce-4ca2-9f6d-0ed332f56a4d. The patched follow-up canary anchored at since=1772985057594 then paused on snapshot 1cb24e7b-f0cd-4e0c-ae5d-27cb4934b49a, auto-resumed on snapshot 151aa03a-fced-41f0-9a54-2f3d1a70856d / run 01KK72HD0EMT3T34K8QP3SMEW9, and drained the held content item back to queue depth 0. The steady-state worker still belongs in QUEUE_OBSERVER_MODE=dry-run; enforce remains a supervised drill until more soak windows say otherwise. Dogfood follow-through also exposed that check/o11y-triage is a long-running current-state scan, not an irreplaceable per-event handler, so it now needs singleton-skip semantics to avoid piling duplicate queued runs behind one active pass.
Hard-won gotcha from the Story 3 live proof: queue operator commands must resolve Redis from the canonical CLI config (~/.config/system-bus.env → REDIS_URL) before ambient shell env. The first proof looked wrong because the shell had an unrelated Upstash REDIS_URL, so queue pause wrote control state to the wrong Redis while queue emit still hit the localhost worker/drainer queue. If the CLI and worker disagree about Redis, fix that first or your proof is bullshit.
Adding a New Inngest Function
- Create
packages/system-bus/src/inngest/functions/<name>.ts - Import
inngestfrom../client - Define the function:
import { inngest } from "../client";
export const myFunction = inngest.createFunction(
{
id: "system/my-function",
// NEVER set retries: 0 — let Inngest defaults handle retries
concurrency: { limit: 1, key: '"my-function"' },
},
{ event: "my/event.name" },
async ({ event, step, ...rest }) => {
const gateway = (rest as any).gateway as
import("../middleware/gateway").GatewayContext | undefined;
const result = await step.run("do-work", async () => {
// your logic here
return { done: true };
});
return result;
}
);
- Export from
index.host.tsorindex.cluster.ts(depending on role) - Add the export to
index.tsas well - Add the event type to
client.tsif it's a new event - TypeScript check:
bunx tsc --noEmit -p packages/system-bus/tsconfig.json - Deploy (see below)
Event Naming Convention
Events describe what happened, not commands: agent/memory.observed not agent/memory.write.
LLM Inference
ALWAYS use the shared utility:
import { infer } from "../../lib/inference";
const { text } = await infer("Your prompt", {
task: "classification",
model: "anthropic/claude-haiku",
system: "System prompt here",
component: "my-function",
action: "my-function.classify",
noTools: true,
print: true,
});
This shells to pi -p --no-session --no-extensions. Zero config, zero API cost. NEVER use OpenRouter, read auth.json, or use paid API keys directly.
If the function is doing long-form editorial or other large-context LLM work, set an explicit timeout on infer() instead of inheriting the shared 120s default. Current earned example: content/review.submitted rewrite/retry/verify runs use a 300s budget because long posts can blow past 120s after all preflight bookkeeping already succeeded.
Gateway Context
All functions receive gateway context via middleware (ADR-0144). Use it for notifications:
const gateway = (rest as any).gateway as
import("../middleware/gateway").GatewayContext | undefined;
await gateway?.notify("event.name", { details });
await gateway?.alert("Something broke", { error: String(err) });
await gateway?.progress("Step 3/5 complete");
Hard Rules
Shortened here. Read the whole file on GitHub.
Signals
- GitHub stars
- 64
- Forks
- 2
- Last commit
- Sep 2026
Advanced
- Item type
- skill
- Key
system-bus- Source
- github.com/joelhooks/joelclaw