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 run-to-completion loop receives the next accepted start batch as its next turn; a loop that can absorb input mid-flight may fold new work into the live turn.
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 turn, your
execute_core reads the accepted start batch from state and runs your
loop to completion, streaming through the communicator so the reusable chat component
renders it live. That is the whole integration surface: implement execute_core,
and the event bus does the rest.
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 start batch, one turn. A ported graph or a bespoke loop runs start to finish without watching the lane. It consumes exactly the accepted start batch; a followup that lands mid-turn is not folded. Because it never opens the handler, the shared door accounts for its start batch, releases the reservation for it after the turn, and only then may publish a liveness wake for work queued during the turn. The net contract for your loop: one start batch, one turn, strictly ordered, with queued work eligible to schedule as soon as the current turn releases the lane.
What You Write
For a run-to-completion loop, the integration is one method and one declaration:
class MyAppEntrypoint(BaseEntrypointWithEconomics):
async def execute_core(self, *, state, thread_id, params):
# the triggering event(s) arrive as the external-event batch on state
question = external_events_text(state.get("external_events") or [])
# ... 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, not what is delivered. For a run-to-completion loop (both
false), a message sent mid-turn is queued for the next turn —
the composer shows "Queue for next turn," and the door's finalizer releases the consumer
before its liveness re-wake. The message stays in the lane and is not folded into the running
turn. An agent that genuinely can absorb input mid-flight declares true and owns
the lane handler the way ReAct does — declare what you actually implement, not what you wish
were true.
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 the start batch, 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.