How Kimi CLI's Wire Protocol Facilitates Communication Between Agent Loops and UI Frontends
Kimi CLI's Wire protocol uses a lightweight spmc broadcast channel with dual raw and merged queues to shuttle Pydantic-typed messages from the agent loop to any UI frontend, letting consumers choose between high-fidelity or coalesced streams.
The MoonshotAI/kimi-cli repository implements a custom communication layer called the Kimi CLI Wire protocol that decouples the core agent loop from interface components such as the shell, ACP, or print backends. Defined by a small set of Pydantic models and a canonical version constant, this protocol ensures ordered, versioned message delivery across multiple UI consumers without blocking the agent.
Architecture of the Kimi CLI Wire Protocol
The protocol centers on the Wire class in src/kimi_cli/wire/__init__.py, which instantiates two internal WireMessageQueue objects backed by a broadcast mechanism. This architecture provides a single-producer-multiple-consumer (spmc) channel that cleanly separates message emission from consumption.
Protocol Version and Message Types
In src/kimi_cli/wire/protocol.py, the version constant WIRE_PROTOCOL_VERSION = "1.10" declares the current protocol revision. All messages inherit from Pydantic models defined in src/kimi_cli/wire/types.py, including ContentPart, ToolCallPart, and the MergeableMixin marker that enables message coalescing.
Dual-Queue Design: Raw and Merged Streams
Every Wire instance maintains two queues:
_raw_queuestores every message exactly as the agent loop emits it._merged_queuestores a coalesced stream where consecutive mergeable messages are combined to reduce UI chatter.
self._raw_queue = WireMessageQueue() # line 24
self._merged_queue = WireMessageQueue() # line 25
How the Agent Loop Publishes Messages via WireSoulSide
The producer interface, WireSoulSide, exposes a send() method that the agent loop calls to dispatch WireMessage instances. When wire.soul_side.send(msg) is invoked, the message is immediately published to _raw_queue. If the message implements MergeableMixin, the side attempts to merge it into an internal buffer; when merging is no longer possible, the buffer flushes to _merged_queue. Non-mergeable messages trigger an immediate flush and pass straight to the merged queue.
self._raw_queue.publish_nowait(msg) # line 82
match msg: # line 88
case MergeableMixin(): … # line 89-95
case _: self.flush(); self._send_merged(msg) # line 96-98
This design means verbose internal events—such as incremental "thinking" progress updates—are automatically compacted before they reach the UI layer, while the raw archive preserves full fidelity for debugging.
How UI Frontends Consume Messages via WireUISide
On the consumer side, a UI frontend obtains a handle through wire.ui_side(merge=True|False). Setting merge=True subscribes the consumer to _merged_queue for a compacted, low-chatter view, whereas merge=False subscribes to _raw_queue for full-fidelity, low-latency delivery. The UI then simply awaits ui.receive() to retrieve the next message in order.
if merge:
return WireUISide(self._merged_queue.subscribe()) # line 46-48
else:
return WireUISide(self._raw_queue.subscribe()) # line 49-50
Because each queue is a broadcast channel, multiple independent UI components can consume the same stream concurrently without missing messages.
Recording to Disk and Graceful Shutdown
The Wire constructor optionally accepts a WireFile backend. When provided, an internal _WireRecorder task subscribes to _merged_queue and appends every coalesced message to a JSON-lines file. This enables later replay or offline debugging of agent runs.
self._recorder = _WireRecorder(file_backend, self._merged_queue.subscribe()) # line 31-33
When the session ends, wire.shutdown() flushes any pending merge buffer, shuts down both queues, and terminates the recorder background task. A subsequent await wire.join() ensures all recorder data is fully flushed to disk.
self.soul_side.flush() # line 52
self._raw_queue.shutdown() # line 54-55
End-to-End Example: Sending and Receiving Wire Messages
The following pattern demonstrates a complete lifecycle: instantiate a Wire, emit messages from the agent loop, consume them in a UI coroutine, and shut down cleanly.
from kimi_cli.wire import Wire
# 1️⃣ Create a Wire instance – optionally persist to a file
wire = Wire() # no file backend in this example
# 2️⃣ Agent loop (soul) sends messages
wire.soul_side.send(ContentPart(text="Starting…"))
wire.soul_side.send(ToolCallPart(name="search", args={"q": "python"}))
# Mergeable messages (e.g., incremental progress) will be coalesced automatically
# 3️⃣ UI front‑end consumes either the raw or merged stream
ui = wire.ui_side(merge=True) # get the compacted stream
async def ui_loop():
while True:
msg = await ui.receive()
# UI can pattern‑match on the concrete message type
match msg:
case ContentPart(text=t):
print("OUTPUT:", t)
case ToolCallPart(name=n, args=a):
print(f"Tool called: {n} with {a}")
# 4️⃣ Shutdown cleanly when the run ends
wire.shutdown()
await wire.join() # wait for recorder flush if enabled
Key files supporting this behavior include:
src/kimi_cli/wire/protocol.py– definesWIRE_PROTOCOL_VERSION.src/kimi_cli/wire/types.py– provides Pydantic models such asContentPart,ToolCallPart, andMergeableMixin.src/kimi_cli/wire/__init__.py– implementsWire,WireSoulSide,WireUISide, and_WireRecorder.src/kimi_cli/wire/file.py– handles JSON-lines persistence for recorded messages.src/kimi_cli/wire/jsonrpc.py– supplies JSON-RPC definitions for remote UI (ACP) communication.
Summary
- The Kimi CLI Wire protocol relies on an
spmcbroadcast channel insrc/kimi_cli/wire/__init__.pyto decouple the agent loop from UI frontends. - Two queues—
_raw_queueand_merged_queue—let consumers choose between full-fidelity and coalesced message streams. - The
WireSoulSideautomatically buffers and mergesMergeableMixinmessages, reducing UI chatter without sacrificing debuggability. - UI frontends obtain a
WireUISidesubscription viawire.ui_side(merge=True|False)and consume messages withawait ui.receive(). - Optional JSON-lines recording and graceful shutdown via
wire.shutdown()andawait wire.join()enable durable debugging and clean teardown.
Frequently Asked Questions
What is the current Wire protocol version in Kimi CLI?
As defined in src/kimi_cli/wire/protocol.py, the current version constant is WIRE_PROTOCOL_VERSION = "1.10". This identifier allows both the agent loop and UI frontends to agree on a shared message schema.
How does Kimi CLI decide whether to merge consecutive messages?
The WireSoulSide checks if a message implements MergeableMixin using structural pattern matching. Mergeable messages are absorbed into an internal buffer until a non-mergeable message arrives or the buffer is flushed, at which point the coalesced result is published to _merged_queue.
Can multiple UI frontends consume the same Wire stream simultaneously?
Yes. The underlying WireMessageQueue is a broadcast queue, so calling wire.ui_side() multiple times creates independent consumers. Each subscriber receives every message in order, whether they attach to the raw or merged stream.
How do I persist Wire messages for later replay?
Pass a WireFile backend when constructing the Wire instance. The internal _WireRecorder subscribes to _merged_queue and appends each coalesced message as a JSON-lines entry. After shutdown, the resulting file can be replayed or inspected offline.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →