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_queue stores every message exactly as the agent loop emits it.
  • _merged_queue stores 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:

Summary

  • The Kimi CLI Wire protocol relies on an spmc broadcast channel in src/kimi_cli/wire/__init__.py to decouple the agent loop from UI frontends.
  • Two queues—_raw_queue and _merged_queue—let consumers choose between full-fidelity and coalesced message streams.
  • The WireSoulSide automatically buffers and merges MergeableMixin messages, reducing UI chatter without sacrificing debuggability.
  • UI frontends obtain a WireUISide subscription via wire.ui_side(merge=True|False) and consume messages with await ui.receive().
  • Optional JSON-lines recording and graceful shutdown via wire.shutdown() and await 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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →