Skip to main content

Turn pipeline

The turn pipeline is the glue that makes Mesh do useful work: a message arrives, and something considered, attributed, and auditable goes back out. Mesh processes conversation as a pipeline rather than as a single in-memory chat loop. That is the difference between a system that can explain itself and one that cannot. Every pass through the pipeline belongs to one RuntimeAgent perspective. The control plane may know more than that agent does; the pipeline is not allowed to turn that extra knowledge into hidden conversational context. See context-resolution.md and ../perspective.md. The pipeline also deliberately stops short of assuming that every accepted turn executes as one linear agent loop. The root Run may stay a loop or, when the work genuinely decomposes, compile into a WorkGraph of bounded jobs. See agent-harness.md and work-graph.md.

Responsibilities

  1. receive a normalized message event for one RuntimeAgent perspective
  2. resolve the actor and conversation context inside that perspective
  3. decide whether a turn should exist at all
  4. create or join the one root Run for this (agent, conversation)
  5. assemble model context from that RuntimeAgent’s eligible perspective state
  6. execute the root Run as a bounded loop or admitted WorkGraph
  7. persist prompt snapshots, run/work state, and resulting output
  8. send the reply back through the message connector

Why it exists

Without the turn pipeline, Mesh is just a set of good ideas:
  • Slack can normalize events
  • OpenRouter can answer prompts
  • Postgres can store state
The turn pipeline ties those pieces together into one observable execution path while preserving the owning RuntimeAgent’s perspective end to end.

Inbound flow

  1. Connector receives a webhook/event.
  2. The connector route resolves the owning RuntimeAgent/control-plane connector.
  3. Raw payload is validated and stored.
  4. Event is normalized into Mesh’s internal schema for that RuntimeAgent’s perspective. Files that rode with the message become attachments on it — metadata (name, type, size, the connector’s fetch address), never bytes — and are persisted alongside the event. The turn that answers the message fetches the images among them through the connector and shows them to the model; everything else is named to the model, and the prompt says per file which it is (see ../connectors/slack.md).
  5. Connector identity is resolved inside that perspective. An actor is linked if one exists; seeing an identity does not require materializing an actor.
  6. Conversation context is resolved inside that perspective, including its shape.
  7. The engagement gate runs. The event resolves to observe, engage, or ignore (see addressing.md). The verdict and the signal that produced it are persisted on this perspective’s event.
  8. observe stops here. The event is durable conversation context and no turn is scheduled. In a busy channel this is the overwhelmingly common path, and it costs no inference.
  9. engage enters the conversation’s mailbox. If a turn is already in flight for this (agent, conversation), it does not start a second one — the interrupt policy resolves it (see interrupt-model.md).
  10. Otherwise a turn and root Run are scheduled, subject to per-agent and per-install admission control. A rejected schedule is a recorded state, not a dropped message.
  11. The root Run resolves context through the current RuntimeAgent/perspective.
  12. The harness chooses execution shape: the ordinary case stays one bounded loop; a genuinely decomposable objective may propose a WorkGraph.
  13. If graph mode is proposed, Mesh deterministically admits the topology before any graph node can run. Only admitted ready nodes are scheduled.
  14. The root Run produces the final response from its loop or graph result.
  15. The response is stored as an outbound event in the same perspective.
  16. The outbound connector delivers it. Today SendReplyCommitting renders the reply and sends its parts, returning a plain error on failure.
  17. Delivery and failure states are persisted.
Step 16 is the intended home of the egress contract, and it is where Mesh will refuse to leak — but that contract is target behaviour, not current behaviour. This PR documents the target; it changes no runtime code. The plan: transport failures against the model (OpenRouter 5xx/429, timeout, empty completion) will be absorbed by bounded retry-and-fallback before a turn produces its response, at zero happy-path cost; an unrecoverable failure will become a structured error event that the user-facing send path cannot accept by type, and the user will see a human-shaped silence or holding line instead of “LLM request failed.” The current SendReplyCommitting path does none of this yet — it returns a plain error and has no EgressError type or model-call recovery. The full target contract — failure taxonomy, the typed egress boundary, the degradation policy, and the heuristic-gated quality judge — is egress-reliability.md. Steps 7 through 10 are the load-bearing ones for ingress stability, and they are the ones a single-user harness does not have. Without them, the gate effectively reads “a turn is scheduled” unconditionally, which in a channel with fifteen people talking means one concurrent tool loop per message until the host runs out of CPU. That is not a hypothetical; it is why this section is numbered the way it is. Steps 12 and 13 are load-bearing for complex execution. Graph fan-out is allowed only after the turn has already passed addressing/admission and after the graph itself has passed a second, deterministic admission gate. A planner drawing ten parallel nodes does not get to manufacture ten unconstrained workers. The perspective boundary is equally load-bearing for correctness. If step 11 or a WorkGraph child node can read deployment-global actor history or another RuntimeAgent’s memory, Mesh has reintroduced a hidden shared-user/shared-context bucket at retrieval time even if the write model looked isolated. Ingress is deliberately lightweight and durable; the expensive work happens in workers. That split is the answer to the “single-threaded event loop chokes” problem, and it is what buys backpressure, retry control, observability, per-turn isolation, and horizontal scaling. A queue is necessary and not sufficient. A queue with unbounded concurrent consumption is a fan-out with extra steps: the work still all runs at once, it just arrives there via a durable table. The properties above require explicit concurrency bounds — one in-flight turn per (agent, conversation), plus per-agent and per-install ceilings — and admission control at the point where a worker picks work up. WorkGraph parallelism is constrained by those same ceilings plus the graph’s own width/budget limits. “Backpressure” with no ceiling to push back against is a word, not a mechanism.

Context resolution

When a worker handles a root Run, it must be able to answer:
  • which RuntimeAgent perspective owns this turn?
  • who is speaking now inside that perspective?
  • who else is in the thread — and how many of them are there?
  • was this actually addressed to me?
  • who is the reply directed to?
  • what has this Actor said recently to this RuntimeAgent?
  • what connectors are involved in this perspective?
  • what policy or persona applies to this RuntimeAgent?
  • which memory and prior conversations are eligible for this audience?
A WorkGraph child worker adds a narrower set of questions:
  • which WorkNode objective am I executing?
  • which admitted incoming edges authorize which WorkArtifacts as inputs?
  • which subset of the root Run’s tools, budget, and perspective context am I delegated?
A pipeline that cannot answer these has collapsed into the flat model ../philosophy.md exists to reject. The normative retrieval contract is context-resolution.md. In particular:
  • there is no context-resolution entry point without a RuntimeAgent/perspective key
  • another RuntimeAgent’s actors, conversations, identity links, or memory are not retrieval sources
  • control-plane knowledge that another participant is also hosted by Mesh does not enter the prompt
  • WorkGraph child contexts remain inside the same RuntimeAgent perspective
  • sibling worker transcripts are not ambient shared context; typed WorkArtifacts cross node boundaries through admitted WorkEdges
  • prompt snapshots eventually record the exact perspective-scoped sources and WorkArtifacts used

WorkGraph execution

Graph execution is an inner scheduling mode of one root Run, not a new turn-routing system. The root Run owns:
  • the user-visible objective
  • attribution to the triggering conversation/actors
  • aggregate budget
  • interrupt/commit state
  • graph revisions
  • final answer
WorkGraph nodes own bounded pieces of that objective. Model-backed agent_loop nodes execute as child Runs or equivalent durable child execution units so they reuse the harness’s RunStep, budget, cancellation, and workspace semantics. The scheduler, not the model, determines which admitted nodes are ready to run. Independent nodes may execute in parallel; sequential dependencies stay sequential. Completed node outputs cross edges as typed WorkArtifacts. A join owns fan-in. A verifier owns independent checking. A human gate owns attributed approval. The root Run remains the only execution that directly consumes new conversation mailbox events. Child nodes do not independently listen to Slack. New user input may cause the root to amend, enqueue, cancel, or propose a WorkGraph revision through interrupt-model.md.

Outbound flow

Outbound replies are events, not side effects. That matters because it is what gives us audit trails, retries, idempotency, delivery state, and connector-specific formatting. A reply that exists only as a side effect cannot be retried safely or explained afterward. The outbound event belongs to the same RuntimeAgent perspective as the turn. Another RuntimeAgent that later observes that Slack message creates its own observation of the sender through its own connector/Actor graph. It does not get a hidden runtime_agent_id link as conversational context. The final reply is owned by the root Run even when a WorkGraph produced it. Individual graph workers do not independently speak into the conversation unless an explicit node is itself an admitted external side effect, in which case the root Run’s commit boundary and attribution still apply.

Render at send, persist canonical

There is exactly one place a reply is translated into a surface’s wire format: turn.SendReply, immediately before the connector’s Send. Everything upstream of that line is canonical Markdown, and that is what the outbound event stores (turn.ReplyBody, body_format = markdown). The direction of the two data flows is the point:
The branch happens at SendReply, not before it, so no wire-format string ever exists as a value that something else could hand to the graph. Why it is arranged that way rather than rendering earlier and storing the result: the graph is read back as prompt history. Presentation written into the record comes back as context and the model imitates it, which is the compounding bug — the actor prefix that compounded into victor — victor — victor —. Formatting decays the same way, more quietly. Store *bold* and history teaches the model to write single asterisks; the renderer then reads those as emphasis, and bold degrades to italics over successive turns. The corollary is that internal/turn must not know which surface it is replying through. It depends on connector.Outbound — render plus deliver, as one interface — and internal/app wires the Slack implementation (slack.NewOutbound). Before render-at-send the send boundary was typed on a concrete slack.SendRequest, which is precisely why per-surface formatting had nowhere to live: the only code that knew a reply was being sent already knew it was going to Slack. connector.Rendered.Parts is plural because chunking belongs to the renderer — only it knows where a code fence opens and closes. SendReply delivers every part in order and keys the outbound event off the first part’s id. See ../connectors/slack.md for the mrkdwn mapping itself. See ../persistence.md for the current write path and why the outbound event is persisted after delivery.

First slice

The initial implementation kept the control flow small:
  • one normalized Slack message in
  • one perspective-scoped prompt assembled
  • one OpenRouter completion out
  • one Slack reply back
The first useful loop now adds one deliberately narrow capability to that path:
  • load up to the 100 newest live memory entries owned by this RuntimeAgent into its prompt — this is a recency-bounded first slice, not relevance retrieval; memory-retrieval.md defines the target eligibility, hybrid candidate, and reranking contract
  • offer four provider-neutral memory tools: remember (create only), recall (perspective-bounded search), amend (revise, optionally re-attribute), and forget (retire, preserving history)
  • accept a bounded three tool rounds (at most eight validated calls each), write every memory entry and its revision row transactionally, then make one final completion with the tool results
  • attribute every runtime-authored memory to the persisted inbound message that triggered it and to the actor who sent it, both taken from the graph record — never from tool arguments or message text
Multiple rounds exist for one specific shape: recall → amend → answer. A model cannot correct an entry without first obtaining its stable id, and one round could not express that. The bound is structural rather than advisory — on the last permitted round the tool list is removed from the request, so the model is unable to continue rather than merely asked not to. Every reply that follows a tool round is preceded by an authoritative ledger of what actually committed, built from the store’s return values. A reply claiming a durable mutation that did not commit fails the turn. See memory-retrieval.md §“Truthful confirmation” for the production incident that contract exists to prevent. The synchronous control flow is enclosed by a durable Run. Every bounded round, model request/response, raw and validated memory-tool call, observation, memory edge, budget counter, and terminal outcome is persisted. The final Run is linked to the outbound message only after that message is stored. Recovery frees the active-run slot without replaying uncertain external effects: it is lease-aware , failing only runs whose lease is absent or expired, so a run a live replica still heartbeats is spared and starting a replica does not fail healthy runs. Slack and Telegram now persist intake before acknowledgment and use leased capture/turn workers; the run itself is leased on the run queue and heartbeat for the life of the turn. Generic ingress still uses its in-process handoff. See durable-intake.md (repository) for the implemented scope, the two ownership systems (work-item intake lease vs. run-execution lease), recovery boundaries and required deployment configuration. A resumed run continues from the ledger’s record of its interrupted attempt rather than re-inferring from the trigger alone (see interrupt-model.md, “What a resume does with the record”); awaiting_input is not implemented. Exec-class tools run in sandboxed workspaces; see sandboxed-execution.md. That is enough to prove a memory-backed shared-conversation loop end to end for one agent. Audience appropriateness is now enforced rather than assumed: visibility is a hard SQL eligibility filter, so a memory learned in one conversation is not reusable elsewhere unless it was stored as perspective-wide. Relevance is still coarse — lexical search plus a current-subject boost and recency, with no semantic retrieval or salience scoring yet. The multi-agent isolation proof is to run two RuntimeAgents in the same external room and confirm that each observes the other only through the connector while retaining independent actors, conversation state, and memory. WorkGraph execution is deliberately not required for that MVP. The architecture is specified now so the durable harness, budgets, child Runs, and sandbox boundaries do not accidentally assume complex work will always remain linear.

Async execution

The first slice is synchronous on purpose, but synchronous call-and-response is not the end state. Real agent work needs several model calls and tool calls per request, which means the pipeline above becomes the outer path into a durable root Run:
  • ingress enqueues and returns instead of blocking on the model
  • a loop-mode Run iterates until it has a final reply, a budget is exhausted, or it needs input
  • a graph-mode Run schedules admitted WorkNodes until its terminal condition is satisfied, a budget is exhausted, or it needs input/approval
  • each loop step, graph revision, node attempt, and artifact is persisted, so progress can be reported and a crashed worker can resume
  • a lifecycle reaction on the triggering message tells the channel the root Run is working
See agent-harness.md for durable loop semantics, work-graph.md for graph execution, and sandboxed-execution.md for where exec-class tool calls actually execute. Async execution also means a turn is no longer the only thing happening in its conversation. Messages keep arriving while it runs, so the root pipeline consumes a mailbox rather than a queue slot: see interrupt-model.md. Four bounds have to exist before the pipeline is pointed at a real shared channel with graph execution, and all are cheaper to build in than to retrofit:
  • Not every message enters the pipeline. In a channel, most inbound messages are context, not requests. The engagement gate decides, before any turn is scheduled and before any inference (see addressing.md).
  • One turn per (agent, conversation), in flight, maximum. “Ingress enqueues and returns” without that bound is unbounded fan-out with a durable log attached: ten messages in one channel become ten concurrent root Runs, and the host, not the design, decides how many is too many.
  • Graph topology is admitted. A model proposal does not directly create workers; node count, width, depth, budget, authority, and effect policy are validated before scheduling.
  • One perspective per context read. No root or child worker may assemble a prompt by reading deployment-global graph state and filtering it after the fact. The RuntimeAgent is part of the query boundary from the beginning.

Why Slack first

Slack is the first connector because it exercises the hardest part of the conversation problem: channels, threads, DMs, multiple participants, multiple RuntimeAgents, event routing, and attribution all at once. It also exercises the perspective boundary cleanly: two RuntimeAgents may occupy the same Slack room while remaining conversationally unaware that the other is co-resident in Mesh. WorkGraph execution is connector-independent by design. A Slack turn, GitHub request, Telegram chat, or synthetic internal objective should all reach the same root Run / WorkGraph machinery after addressing and context resolution. If Slack is right, the rest of the connector ecosystem is a portability problem rather than a conceptual one. See ../connectors/slack.md. Document attachments can additionally be read on demand through the scoped list_documents / read_document path (repository). Original attachment rows remain metadata-only; delivered extracted text enters normal tool observations.