Connect Your Agentic Loop To The Event Bus
Bring your own agent — your graph, your framework, your control loop. KDCube hands it user events through the event bus: ordered and one turn at a time per conversation. A wake names one accepted occurrence; a run-to-completion adapter folds it and every other still-pending occurrence into one start snapshot. A live loop may fold new work while the turn runs.
You have a working agentic loop. On KDCube it becomes an app, and users talk to it. So the first question is not "how do I run my graph" — it is "how does a user's message reach my loop, and what happens when a second message arrives while the first is still running?" The answer is the same for every agent on KDCube, whether it is the built-in ReAct engine or a ported LangGraph graph: the event bus orders the messages, and current lane state decides whether an open live handler folds the event or the shared door starts the next turn.
One Door For Every Turn
Nothing starts a new turn except a wakeup. An accepted reactive batch is appended to the conversation's event-bus lane — an ordered log — and its primary wakeup is enqueued in the same atomic operation. When the processor claims that wakeup, it consults current lane state. An open live handler may already own the event; otherwise the processor schedules the shared door:
turn-starting event (prompt / queued followup)
-> appended to the conversation event-bus lane (ordered log)
-> primary wakeup enqueued [same atomic operation]
-> current open handler? yes -> fold into that live turn
no -> run() (@on_reactive_event)
-> execute_core(...) runs the next turn
run() is the reactive-event door. For a newly scheduled foreign turn, your
execute_core performs one read-only fold of the lane: the wake occurrence plus
every other still-pending occurrence, ordered by lane sequence. It maps all of those messages
into the framework and runs the loop to completion, streaming through the communicator so
the reusable chat component renders it live. The wake identifies work; it is not the complete
turn input.
Ordered, One At A Time
Turns of a single conversation are serialized by a per-conversation lock.
A second message that arrives while a turn is running does not start a
second execute_core — it waits its turn, in order:
Event 1 -> wakeup -> conversation lock acquired -> run()/execute_core (turn 1 running)
Event 2 arrives now -> wakeup enqueued -> cannot acquire the lock -> requeues (waits)
turn 1 ends -> lock released -> Event 2's wakeup claimed -> run()/execute_core (turn 2)
- Same conversation: one turn at a time, in arrival order. The lock holds across processor workers, so two workers cannot run two turns of one conversation at once.
- Different conversations: run in parallel, on their own locks.
You do not implement any of this. What a turn is responsible for is releasing the piece of the event bus it was handed — so the message waiting behind it can run. That piece is the lane reservation.
The Lane Reservation
When the processor dispatches a wakeup it reserves the event-bus lane consumer for that conversation before the turn runs. The reservation is how the platform knows a turn is responsible for this lane. Whoever holds it must release it when the turn ends. A fresh duplicate wake is correctly deferred while a real starter is loading. Finalization releases the reservation before publishing any additional liveness wake, so a queued event can schedule the next turn.
Releasing the reservation is the one place the two kinds of agent differ.
Two Ways To Consume
ReAct — a live consumer that folds mid-turn. A ReAct turn opens the lane handler and reads the event bus while it runs. A followup that lands mid-turn is folded into the running turn at a decision boundary. Once the close gate closes the handler, it stops the live reader immediately. After artifacts persist, finalization releases the consumer and only then may publish a duplicate liveness wake for anything still unconsumed. ReAct owns its lane lifecycle inside its own workflow.
Run-to-completion — one pending fold, one turn. A ported graph or a bespoke loop first snapshots the wake occurrence and every other still-pending lane occurrence, then runs start to finish without folding later content. The snapshot may cross several ingress batches; each batch id remains attribution, not the turn boundary. A supported adapter may watch control read-only, but ordinary arrivals remain pending. Because this path never opens the ReAct handler, the shared door accounts for the snapshot's exact ids, releases the reservation for it after the turn, and only then may publish a liveness wake. The net contract is one pending start fold, one turn.
What You Write
For a run-to-completion loop, the integration is one method and one declaration:
from kdcube_ai_app.apps.chat.sdk.protocol import external_events_texts
from kdcube_ai_app.apps.chat.sdk.solutions.foreign_runtime import (
fold_turn_external_events,
)
class MyAppEntrypoint(BaseEntrypointWithEconomics):
async def execute_core(self, *, state, thread_id, params):
# The wake names one occurrence, not the complete turn input.
state["external_events"] = await fold_turn_external_events(self, state)
messages = external_events_texts(state.get("external_events") or [])
question = "\n".join(
f"{index}. {text}" for index, text in enumerate(messages, start=1)
)
# ... run YOUR loop / graph to completion ...
# stream tokens + steps through the current communicator (comm_ctx)
# per agent, in the app descriptor
conversation:
accepts_followup: false # this loop does not fold a new message mid-turn
accepts_steer: false # this loop does not cancel + finalize mid-turn
That is the entire wiring. The door serializes turns, releases the lane reservation for you, and lets queued work schedule in order. You never touch the event bus.
Followup Without Folding
accepts_followup / accepts_steer change what the composer
offers and must match what the adapter implements; they do not alter lane ordering.
With both false, a message sent mid-turn is queued for the next
turn. The door releases the consumer before its liveness re-wake, and that next turn
folds all eligible pending messages together.
A foreign adapter may instead declare accepts_steer: true and wrap its native
run with the shared read-only control watcher. The ported LangGraph adapter uses that path:
steer cancels its streaming task, while ordinary messages are not folded into the active
graph. At handoff a bare steer is spent as the stop boundary. A steer carrying text stays
pending for a later fold but does not wake a turn by itself. An agent that genuinely absorbs
content mid-flight needs a live handler like ReAct's; a control watcher is not content
folding.
The Failures This Prevents
Before the door finalized the lane for run-to-completion turns, such a turn left its
reservation held. The next turn's wakeup, arriving inside the freshness window, was dropped
as scheduled_consumer_fresh: the turn "completed" in the UI, but the
next message never reached execute_core — nothing in the processor log
— and only recovered after the window went stale. It read as an intermittent "second or third
turn hangs." The finalizer now accounts for every exact id in the turn's pending snapshot,
releases the reservation, and then publishes the optional liveness wake.
A second failure lived one layer lower. Cancellation could arrive after Redis accepted the
lane-state lock but before the client observed the acknowledgement. The lock then survived
until its TTL, a wake timed out during normalization, and valid work could be acknowledged as
invalid. Lock acquisition and release now complete to a known outcome, uncertain ownership
is cleaned by exact token, and a pre-execution lock timeout requeues the same wake. A fresh
active consumer suppresses a wake only while its handler remains
open; closed + active is recoverable finalization residue, not proof
that someone can still fold the event.
Read More
- The Conversation Is a Lane — the deeper wake, handler, and ownership model.
- Reactive Turn Delivery — the framework-neutral event-bus delivery contract.
- Connect Your Agentic Loop To Ordered Message Delivery — the builder recipe.
- Conversation Event Lane State — the event-bus state primitive.
- Settle Your Solution In A KDCube App — the end-to-end settlement this delivery slice fits into.