انتقل إلى المحتوى

البث في JavaScript

عرض بصيغة Markdown

يبثّ K-Agent بأحداث Server-Sent Events القياسية، لكنك تقرؤها عبر fetch لا EventSource: فـ EventSource لا يستطيع إرسال ترويسة Authorization ولا جسم طلب POST. الكود أدناه يعمل في المتصفحات وفي Node 18 أو أحدث، بلا أي اعتماديات، ومُختبَر على بثوث تتقطع فيها الأسطر والمحارف العربية والأحداث عبر حزم الشبكة.

/** Parses a text/event-stream response into { id, event, data } objects. */
export async function* readEvents(response) {
const reader = response.body.pipeThrough(new TextDecoderStream()).getReader();
let buffer = '';
let ev = { id: undefined, event: 'message', data: [] };
for (;;) {
const { value, done } = await reader.read();
if (done) return;
buffer += value;
const lines = buffer.split('\n');
buffer = lines.pop(); // keep the unfinished last line for the next chunk
for (const raw of lines) {
const line = raw.endsWith('\r') ? raw.slice(0, -1) : raw;
if (line === '') {
// A blank line ends an event.
if (ev.data.length) yield { id: ev.id, event: ev.event, data: JSON.parse(ev.data.join('\n')) };
ev = { id: undefined, event: 'message', data: [] };
} else if (!line.startsWith(':')) {
// Lines starting with ":" are heartbeats; everything else is "field: value".
const i = line.indexOf(':');
const field = i === -1 ? line : line.slice(0, i);
const value = i === -1 ? '' : line.slice(i + 1).replace(/^ /, '');
if (field === 'id') ev.id = value;
else if (field === 'event') ev.event = value;
else if (field === 'data') ev.data.push(value);
}
}
}
}

قواعد التشغيلات والبث تتحول إلى أسطر قليلة:

  • احفظ النص لكل message_id — فالتشغيل الواحد قد ينتج أكثر من رسالة مساعد؛
  • أضف نص message.delta معاينةً؛
  • عند message.completed استبدل المعاينة بالنص المعتمد — أو احذفها حين يكون discarded: true؛
  • احفظ آخر id للأحداث، لتستطيع الاستئناف؛
  • توقف عند حدث run.* النهائي.
const BASE = 'https://api.k-agent.kerneltics.com/v1';
const TERMINAL = new Set(['run.completed', 'run.failed', 'run.cancelled', 'run.superseded', 'run.requires_action']);
const textOf = (message) =>
message.content.filter((part) => part.type === 'text').map((part) => part.text).join('');
/**
* Sends a message and streams the reply. onText receives the reply text so far
* (all assistant messages of the run). Resumes with Last-Event-ID if the connection drops.
* Returns the terminal event: run.completed, run.failed, run.requires_action, …
*/
export async function streamReply({ session, input, token, onText, signal }) {
const headers = { Authorization: `Bearer ${token}`, Accept: 'text/event-stream' };
const texts = new Map(); // message_id → text
let runId;
let lastId;
async function consume(response) {
if (!response.ok) {
const { error } = await response.json();
throw Object.assign(new Error(error.message), { code: error.code });
}
for await (const { id, event, data } of readEvents(response)) {
if (id) lastId = id;
if (event === 'run.created') runId = data.run.id;
else if (event === 'message.delta') {
texts.set(data.message_id, (texts.get(data.message_id) ?? '') + data.delta);
} else if (event === 'message.completed') {
if (data.discarded) texts.delete(data.message_id ?? data.message?.id);
else texts.set(data.message.id, textOf(data.message)); // authoritative text
} else if (event === 'error') {
throw Object.assign(new Error(data.code), { code: data.code });
} else if (TERMINAL.has(event)) {
onText([...texts.values()].join('\n\n'));
return { event, data };
} else continue; // other events: ignore what you don't use
onText([...texts.values()].join('\n\n'));
}
return undefined; // the connection ended early
}
let response = await fetch(`${BASE}/sessions/${encodeURIComponent(session)}/messages`, {
method: 'POST',
headers: { ...headers, 'Content-Type': 'application/json' },
body: JSON.stringify({ input, stream: true, client_message_id: crypto.randomUUID() }),
signal,
});
for (let attempt = 0; ; attempt++) {
try {
const end = await consume(response);
if (end) return end;
} catch (err) {
if (err.code || signal?.aborted) throw err; // an API error, or you aborted
// otherwise a network error: resume below
}
if (!runId || attempt === 3) throw new Error('The stream was interrupted.');
await new Promise((resolve) => setTimeout(resolve, 500 * 2 ** attempt));
response = await fetch(`${BASE}/runs/${runId}/events`, {
headers: lastId ? { ...headers, 'Last-Event-ID': lastId } : headers,
signal,
});
}
}

استخدمها هكذا:

const output = document.querySelector('#reply');
const end = await streamReply({
session: 'order-8812',
input: 'أقدر أغير عنوان التوصيل؟',
token, // a client token in the browser; a secret key only on your server
onText: (text) => {
output.textContent = text; // plain text: never innerHTML for model output
},
});
if (end.event === 'run.completed' && end.data.run.outcome === 'handed_off') {
showBanner('سيرد عليك أحد أعضاء فريقنا هنا.');
}

3. في المتصفح: استخدم رمز عميل

رابط القسم «3. في المتصفح: استخدم رمز عميل»

المفاتيح السرية لا مكان لها في المتصفح أبدًا، والواجهة البرمجية لا ترسل ترويسات CORS لها. في كود الواجهة الأمامية استخدم رمز عميل يصدره خادمك للمستخدم المسجّل (POST /v1/client_tokens). تستطيع رموز العميل استدعاء مسارات الرسائل والأحداث في جلساتها من المتصفح، بما في ذلك الاستئناف مع Last-Event-ID. انظر العملاء والهوية.

حين تنتهي صلاحية رمز العميل أثناء البث يرسل البث event: error مع {"code": "token_expired"} ثم يُغلق. وترمي streamReply خطأً يحمل هذا الرمز code: اجلب رمزًا جديدًا من خادمك واستأنف عبر GET /v1/runs/{run}/events مع آخر id.

بث التشغيل ينتهي بانتهاء تشغيله. ولعرض كل شيء في المحادثة — ردود الوكيل، وردود فريقك أثناء التحويل، وتغيّر الوضع — أبقِ بث الجلسة مفتوحًا:

async function followSession({ session, token, after, onEvent, signal }) {
const url = new URL(`${BASE}/sessions/${encodeURIComponent(session)}/events`);
if (after) url.searchParams.set('after', after); // replay what you missed
const response = await fetch(url, {
headers: { Authorization: `Bearer ${token}`, Accept: 'text/event-stream' },
signal,
});
if (!response.ok) throw new Error((await response.json()).error.code);
for await (const ev of readEvents(response)) {
onEvent(ev); // message.created, message.delta, message.completed, session.updated, …
}
}
  • بث الجلسة لا يُغلق من تلقاء نفسه؛ أوقفه عبر AbortController عند إغلاق الدردشة.
  • دون after يبدأ مباشرًا؛ و?after=0 يعيد كل ما زال محفوظًا (7 أيام). احفظ آخر id ومرّره بعد إعادة تحميل الصفحة.
  • ما دام البث مفتوحًا، أرسل الرسائل مع "background": true واعرض المحتوى من بث الجلسة وحده — فلا تظهر أي رسالة مرتين. هكذا يعمل ودجت K-Agent.
  • session.updated مع mode: "human" إشارة لعرض «أنت تتحدث مع موظف الآن».
  • نبضات الاتصال (: ping) تصل كل 15 ثانية، والمحلّل يتجاوزها. وإذا لم يصلك شيء لمدة أطول بكثير فأعد الاتصال.
  • إيقاف البث لا يوقف التشغيل. إغلاق الاتصال يوقف قراءتك فقط؛ ويكتمل التشغيل ويُحفظ رده. ولإيقاف التشغيل نفسه استدعِ POST /v1/runs/{run}/cancel.
  • الأخطاء قبل البث تعود JSON عاديًا بحالة HTTP ‏(409 session_busy و401 invalid_api_key…)، وترميها consume مع رمز الخطأ code.
  • تظهر أحداث جديدة مع الوقت. تجاهل الأنواع التي لا تتعامل معها، كما في الكود أعلاه.