Architecture
This is the engineering reference — how a turn actually runs, why the pieces are shaped the way they are, and where to look in the code. It is deliberately separate from the capability guides; those tell you how to use the runtime, this tells you how it works.
Five principles run through everything below:
- RubyLLM does the model work. Chat, streaming, the tool-loop, and provider retries are never reimplemented — the engine is the operational shell around them.
- One pipeline. Every conversational turn goes through the same ordered stages. There is no fast path that skips policy or persistence.
- Everything is a Command. Every state-changing interaction is a typed command through one bus; reads go straight to the stores, never through the runtime.
- Durable by default. Every turn checkpoints; a killed process resumes from the last checkpoint on reboot.
- Config over code. Agents, tools, skills, and policies are data in stores, editable hot — not classes you redeploy.
The engine at a glance
A request enters as a Command, becomes a Task running on its own fiber, and streams its progress out through the Event Stream as it moves down the pipeline.
flowchart TD
client([HTTP client]) -->|"POST /v1/responses"| bus[Command Bus]
bus -->|turn command| task[Task actor<br/>Async fiber]
task --> cb[Context Builder]
cb --> pol[Policy Engine]
pol --> mw[Middleware<br/>edge limit · input guardrail]
mw --> ex[Executor<br/>chat + tool-loop]
ex --> rll[(RubyLLM ⇄ provider)]
ex --> persist[Persistence<br/>checkpoint · session · task]
task -.emits.-> es[Event Stream]
es -->|SSE| client
cb -.reads.-> stores[(SQLite stores)]
persist -.writes.-> stores
- Command Bus validates and dispatches. A turn command (
send_message,trigger_workflow) creates a Task and returns its id immediately; the result flows out on the Event Stream. A control command (create_session,cancel_task,pause_task,approve_action,resume_task) acts on stores synchronously. Queries are not commands — they read the stores directly. - Task actor is an
Asyncfiber with a minimal mailbox (cancel,user_message, pause). Because an LLM turn is almost all waiting on the provider, one process runs many concurrent turns on a fiber scheduler instead of pinning a thread per call. - Event Stream is in-process pub/sub. Every event carries
task_idand a monotonicseq, so one stream multiplexes many turns and replays reliably. Not every event is for the end user — see the edge contract below.
In-process is a real boundary, not a detail: per-session FIFO ordering,
steer, interrupt, pause/cancel and the SSE watch all act on the worker
that holds the session’s actor, and the engine does not promise them across
worker processes. What is cross-process is everything durable — sessions,
tasks, checkpoints, outbox, delegations — behind transactional claims on the
shared store. The deploy-side consequences (and the sticky-routing escape
hatch) are written in Deploy → The process model.
The mirror image of that question — not N workers of one deployment, but N graphs inside one process, which is what happens when you mount Insika into an app you already have — is the embed contract. Short version: a graph owns its store and its LLM credentials, so two graphs no longer swap keys or read each other’s sessions; the process keeps owning signals, the reactor and the Studio, which is why the host installs the drain itself.
What crosses the edge
A turn is not one assistant message. Between the user’s message and the answer the model may narrate the tool loop (“let me look that up”), apologise for a tool that failed, or — when it has no tool to call — reason in prose. All of it arrives as ordinary content chunks, indistinguishable at the token level from the answer.
So the engine publishes rather than relays: :content carries the answer — the
text of the assistant message that ends the turn. It is emitted once, whole,
when that message ends. Everything else the model says rides :intermediate, and
the provider’s own reasoning channel rides :thinking. Both are real events —
the Studio renders them and the trace keeps them, which is how an operator sees
what the model was doing — and /v1/responses deliberately translates neither.
Two consequences worth knowing before you build on it:
- The customer-visible stream is per message, not per token. A consumer that
accumulates deltas gets the same text; one that renders them live gets it in one
piece. Watch
:intermediateif you want the keystrokes. Measured on a real agent, the answer’s frames span 0 ms — which is why a relay costs the customer nothing next to holding an SSE connection for the whole turn. - A turn that dies mid-message publishes nothing. Half a sentence was never an answer; the fragment is still on the stream for whoever is debugging it.
That contract is what makes a channel possible at all. A
channel whose recipient is not on the connection — a WhatsApp
number, your own callback URL — cannot stream anything: it needs one message it
can send. :content is that message, which is why a channel delivers exactly it
and nothing else.
The exception is halt_when: a tool that ends the turn has already answered the
customer, so the model’s lead-in before that call is the turn, and it is
published as the answer.
Neither default is a law. A product with a “thinking” panel wants the reasoning, and a chat UI may want the progress line. An agent opts a channel in:
edge_stream thinking: true, intermediate: false
Two things keep that from re-opening the hole. Nothing crosses unless someone
opted in, per agent. And what crosses gets its own frame type, never the
answer’s — response.reasoning_summary_text.delta for reasoning, and a namespaced
insika.intermediate.delta for the narration, because the Responses protocol has
no honest event for “assistant text that is not the answer” (there, that text is
output_text.delta, told apart only by an output-item index this adapter does not
carry). So a consumer that accumulates output_text deltas into one message —
WhatsApp — is unaffected by the switch, and one that renders reasoning has
something to render.
A turn, end to end
The Executor runs a fixed sequence of stages. Each stage boundary drains the mailbox, so a cancel or pause is honored at a safe point — never mid-write. One of those boundaries sits between the provider’s last word and publishing the answer: a turn cancelled while the model was working publishes nothing, so what the customer read and what the transcript holds never disagree.
The numbers are the engine’s own — the same stage numbers executor.rb uses. The
sequence has no stage 7; the numbering is kept as the code has it rather than
renumbered here.
flowchart TD
s1["1 · Command Bus<br/>send_message → Task on a fiber, task_id returned"]
s2["2 · Context Builder<br/>providers, budget, pinned"]
ck["initial checkpoint<br/>the state at the START of the turn"]
s3["3 · Policy Engine<br/>allowed tools/skills, approval tags"]
s1 --> s2 --> ck --> s3 --> mw
subgraph mw ["4 · Middleware wraps everything below — edge limit → input guardrail"]
s5["5 · assemble chat"]
s6["6 · agent interaction<br/>chat.ask + tool-loop"]
s8["8 · Persistence<br/>checkpoint → session → task"]
s9["9 · Response<br/>task_completed + usage"]
s5 --> s6 --> s8 --> s9
end
s9 --> hook["after-task hook<br/>output guardrail"]
The order is not arbitrary:
- Context before policy. The prompt is assembled first so policy can see what the turn will actually contain (candidate skills come from the catalog, tools from the registry).
- The initial checkpoint is written before the model call. “The checkpoint of turn n holds the state at the start of turn n.” Without it, a crash during the model call would orphan the task with no checkpoint — unrecoverable.
- Middleware wraps the chat stages. For opt-in reranking, the leading edge limiter also wraps context preparation: admission runs once before paid retrieval. Context still precedes policy; the input guardrail and plugin middleware receive the prepared context, including any budget warning, before chat runs. Turns with lexical-only retrieval keep the original order.
- Persistence is a fixed order (checkpoint → session → task) and a pure drain point: the last stage never suspends, so a checkpoint is never left half-written.
- Output validation runs as an after-task hook on the produced content (the output guardrail).
The tool-loop
Stage 6 is the agent interaction. RubyLLM owns the reason→act→observe loop;
ToolEnvelope adds provenance, approval, concurrency, timeout, evidence processing,
fencing, traces and side-effect checkpoints around registered tools.
Calls run serially unless limits[:tool_concurrency] permits parallel execution.
Marked side effects acquire a serial gate before the shared concurrency slot,
so writes in one session cannot overlap and queued writes leave slots for reads.
Turns exposing approval-required tools run serially. See
Tools for limits and timeout behavior.
flowchart TD
ask[chat.ask → model] --> dec{tool call?}
dec -->|no| done[final content]
dec -->|yes| env[ToolEnvelope: skip completed side effects on resume]
env --> prov{declared evidence IDs known?}
prov -->|no| blocked[blocked: provenance]
blocked --> ask
prov -->|yes / undeclared| appr[approval gate]
appr --> gates[side-effect serial gate → shared concurrency slot → timeout]
gates --> kind{tool kind}
kind -->|code| ruby[Ruby / sandbox]
kind -->|HTTP data| http[EgressGuard → HTTP]
kind -->|MCP| mcp[live MCP client]
kind -->|presentation| cards[select turn or session evidence cards → UI event]
ruby --> result[evidence reshape → fencing → checkpoint and trace]
http --> result
mcp --> result
cards --> result
result --> ask
HTTP data tools and presentation tools are both stored definitions. Presentation
runs in-process; MCP remains a live server call (Streamable HTTP or stdio). HTTP egress
is checked before a request. Tool errors reach the model as error results so it
can recover; a provenance refusal reports blocked without contacting the backend
or asking an operator. Completed side effects are recorded for resume.
Ingesting tools: manifest and MCP
Data definitions and MCP instances have separate stores. Both are hot config:
flowchart TD
manifest[POST /v1/tools/manifest] --> substitute[resolve env / secret placeholders]
substitute --> validate[validate HTTP or presentation definition]
validate --> store[(ToolStore)]
store --> data[DataToolRegistry]
config[MCP instance configuration] --> mcpstore[(McpStore)]
mcpstore --> live[McpToolRegistry → live server tools]
data --> catalog[effective registry and catalog]
live --> catalog
catalog --> policy[agent allowlist → tool-loop]
Only manifest ingestion resolves {{env.*}} and {{secret.*}}; other data-tool
write paths require literal values. One malformed manifest tool is reported in
errors[] while valid entries import. MCP tools are never converted to stored
HTTP data tools. See Tools.
Durability: checkpoints and resume
The runtime has no external job queue. Durability is stores plus boot recovery:
at startup the recovery scan finds tasks that were mid-flight and resumes each from
its last valid checkpoint — the same code path a resume_task command uses.
flowchart LR
start(( )) -->|send_message| running[running]
running -->|turn persisted| checkpointed[checkpointed]
checkpointed -->|next turn| running
running -->|approval required| waiting[waiting]
waiting -->|approve_action| running
running -->|process killed| crashed[crashed]
crashed -->|boot recovery / resume_task| running
checkpointed -->|task_completed| completed[completed]
completed --> done((( )))
Resume always replays from the start of the last checkpointed turn. Side-effect calls whose completion was recorded in the checkpoint are skipped when replayed with the same call id. An external effect can succeed before that record is saved; a crash or failed checkpoint write in that window can repeat the effect on resume. Use destination-supported idempotency keys or reconciliation for those operations: checkpointing alone does not provide exactly-once external execution. A resumed turn is also never re-counted against edge-limit ledgers. Cancellation is cooperative — checked at stage boundaries, never in the middle of a store write.
Composition root
The whole graph is wired in one place — Insika::Wiring::Graph — in two phases,
so the two deployment roots (a minimal in-process wiring and the full server
deployment) share the parts that are identical and layer on only what genuinely
differs.
flowchart TD
root1[minimal wiring] --> phase1
root2[server deployment] --> phase1
subgraph phase1 [Phase 1 · spine — infra, identical across roots]
backend["backend<br/>SQLite (INSIKA_DB) or Memory"] --> dstores[domain stores<br/>session · task · checkpoint · memory]
reg[registries<br/>tools · workflows · policies]
caps[capability registry]
es2[event stream]
hk[hooks]
end
phase1 ==>|"Graph.build(spine:)"| phase2
subgraph phase2 [Phase 2 · build — assembled on the spine]
cbz[Context Builder] --> exz[Executor]
pez[Policy Engine] --> exz
mwz["Middleware<br/>edge limiter → input guardrail"] --> exz
exz --> busz[Command Bus<br/>6 core commands]
end
Everything in phase 1 is passed into phase 2: the domain stores, registries,
capability registry, event stream and hooks are all constructor arguments of the
Executor and the Command Bus (spine.* throughout Graph.build). The arrow is one
call, not one wire.
backend_from_env picks the backend: INSIKA_DB set → durable SQLite (the
prerequisite for recovery); unset → ephemeral in-memory (dev/demo). Registering the
operator commands (pause_task, approve_action) in the shared core is what lets
both roots expose the Studio’s controls without a per-root patch. The guardrails
factory contributes the input guardrail as the single middleware and the output
validator as the after-task hook, so both roots enforce content safety identically.
Where the code lives
| Concern | Code |
|---|---|
| Composition root | lib/insika/wiring/graph.rb |
| Command bus + handlers | lib/insika/command_bus.rb, lib/insika/commands/* |
| Turn pipeline | lib/insika/executor.rb |
| Context assembly | lib/insika/context/* |
| Policy | lib/insika/policy/* |
| Stores | lib/insika/stores/*, lib/insika/*_store.rb |
| Recovery | lib/insika/recovery.rb |
| Inbound queue (one turn at a time per session, and what happens to a message that arrives while one is running) | lib/insika/session_actor.rb, lib/insika/queue_policy.rb, lib/insika/steer_injector.rb |
| Channels (a way in and out for people; the reply that travels after the turn ends) | lib/insika/channel_registry.rb, lib/insika/channels/*, lib/insika/channel_delivery.rb, lib/insika/outbox_store.rb, lib/insika/inbound_log.rb |
| Tools (data/manifest/MCP) | lib/insika/tool_definition.rb, lib/insika/tool_manifest.rb, lib/insika/mcp_tool_registry.rb, lib/insika/tools/present.rb |
| Plugin loading (boot) | lib/insika/plugin.rb, lib/insika/plugin/loader.rb |
| Refinement (traffic → report) | lib/insika/refinement/*, lib/insika/refinement_store.rb |
| Post-turn learning (facts, skills, knowledge — extracted from finished conversations) | lib/insika/distill.rb, lib/insika/harvest.rb, lib/insika/knowledge.rb, lib/insika/knowledge_store.rb; the per-turn hook lives in Executor#persist_turn, next to finalize_delegation |
| Evals (cases, judges, gate) | lib/insika/evals/*, lib/insika/golden_store.rb; evals/run.rb is the CLI |
| HTTP/SSE surface | lib/insika/server/* |
See also
- Agents · Tools · Skills · Context · Channels · Plugins · Security — the capability guides.
- Evals — the cases that grade an agent, and the pre-merge gate.
- Refinement — reading a live agent’s own traffic back as a report.
- Deploy — running the engine durably.
- Benchmark — the per-turn engine overhead, reproducible.