WuKongIM Docs

Agent Streaming Replies

Understand the flow, connect a real model, and update one reply through live SDK events.

One base message + incremental events = one continuously updated chat bubble.

Before you start

Use easyjssdk@2.0.5 and complete the Web Quickstart. The server must include online EVENT delivery; verify support in older released images.

1. Understand the flow

LocationResponsibility
Frontend SDKSend questions, receive events, update bubbles
Product backend / BFFIssue credentials, authorize access, expose history and cancellation
Agent backendHold model API keys, call the model, serialize event writes

2. Run it

Start a WuKongIM cluster supporting this feature, then run the demo:

cd demo/streamdemo
npm ci
npm run build
npm start

Open http://127.0.0.1:5175/streamdemo/ and select 真实模型 (Real model) in settings:

SettingValue
Model URLCompatible Chat Completions /v1 base or full /chat/completions URL
API keyModel service key
Model nameLeave blank to try /models; enter explicitly when discovery is unavailable
WuKongIM URLDefault http://127.0.0.1:5001; editable

The demo uses a local proxy for model CORS. Production keys and WuKongIM HTTP API calls stay on the backend. Run instructions

3. Frontend: receive SDK events

npm install --save-exact easyjssdk@2.0.5

The BFF supplies bootstrap. Implement both UI callbacks using the receiver code as a reference. This example uses group Channel 2 with the user and Agent as members.

import { WKIM, WKIMEvent } from 'easyjssdk';

const im = WKIM.init(bootstrap.websocketUrl, {
  uid: bootstrap.uid, token: bootstrap.token, deviceFlag: 1,
}, { singleton: false, debugLogging: false });

// On exit: use the same callback references.
function stop() {
  im.off(WKIMEvent.Message, upsertMessage);
  im.off(WKIMEvent.CustomEvent, applyStreamEvent);
  im.destroy();
}
window.addEventListener('pagehide', stop, { once: true });

im.on(WKIMEvent.Message, upsertMessage);
im.on(WKIMEvent.CustomEvent, applyStreamEvent);
await im.connect();
await im.send(channelId, 2, { type: 1, content: 'Hello' });
CallbackMerge rule
upsertMessageCreate the base bubble; reuse a bubble created by an earlier event
applyStreamEventFilter Channel, Agent UID, and lane first; update one reply by client_msg_no
Duplicates / gapsDeduplicate event_id; merge text_offset in UTF-8 bytes; wait for a full snapshot when text is missing

4. Backend: model deltas → message events

This is the normal-flow core. agent is a connected backend SDK. Authorize the question and membership before creating an idempotent job. chunks comes from the model SSE adapter; configure the model URL, key, and name on the backend.

const clientMsgNo = job.replyClientMsgNo; // Persist once per reply.
const ack = await agent.send(channelId, 2, { type: 1, content: '' }, {
  clientMsgNo, setting: { stream: true },
});
if (ack.reasonCode !== 1) throw new Error('Base message rejected');

await append('stream.open', { kind: 'text' });
let text = '';
for await (const delta of chunks) {
  await append('stream.delta', { kind: 'text', delta }); // Serialize writes.
  text += delta;
}
await append('stream.finish', { snapshot: { kind: 'text', text } });

Implement append on the backend: call POST /message/event (OpenAPI reference), check HTTP and application results, and save the ID and payload before requesting. Reuse them on retries. Example request:

{
  "channel_id": "room-42", "channel_type": 2, "from_uid": "agent-42",
  "client_msg_no": "reply-42", "event_id": "reply-42:2",
  "event_key": "main", "event_type": "stream.delta", "visibility": "public",
  "payload": { "kind": "text", "delta": "Hello" }
}

5. Finish and recover

ScenarioAgent actionUI / history
Completestream.finish + full snapshotComplete answer
User cancelsAbort model → stream.cancel → stream.finishPartial text and cancelled state
Model failsstream.error → stream.finishPartial text and failed state
Uncertain IM writeStop; reconcile and retry the original eventNo blind new error or completion event

Online reception uses SDK push. History sync is for offline recovery: no polling or /message/eventsync.

KeepReason
Persistent base with stream: trueAppend only after successful SENDACK; no NoPersist / SyncOnce
Serialized writes and stable event IDsPrevent reordering and repeated text
Complete terminal snapshot and independent write timeoutLive events may be lost; cancellation still needs finalization
Bounded buffers and full textBuffer during recovery; replay a full snapshot before finish after Leader cache loss
main and __finish__Finish uses the reserved lane; do not filter it out
Backend authorization and history filteringPrivate events are not broadcast; visibility is not product authentication

6. Full implementation and acceptance

Check normal reply → duplicate event → cancellation → truncated model stream → offline reconnect. The same bubble must match history each time. After fixtures pass, verify your actual model.

API details: Message Events · History Sync · JSON-RPC

On this page