البث في JavaScript
يبثّ K-Agent بأحداث Server-Sent Events القياسية، لكنك تقرؤها عبر fetch لا EventSource: فـ EventSource لا يستطيع إرسال ترويسة Authorization ولا جسم طلب POST. الكود أدناه يعمل في المتصفحات وفي Node 18 أو أحدث، بلا أي اعتماديات، ومُختبَر على بثوث تتقطع فيها الأسطر والمحارف العربية والأحداث عبر حزم الشبكة.
1. حلّل البث
رابط القسم «1. حلّل البث»/** 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); } } }}2. أرسل رسالة واعرض الرد
رابط القسم «2. أرسل رسالة واعرض الرد»قواعد التشغيلات والبث تتحول إلى أسطر قليلة:
- احفظ النص لكل
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.
4. تابع جلسة كاملة
رابط القسم «4. تابع جلسة كاملة»بث التشغيل ينتهي بانتهاء تشغيله. ولعرض كل شيء في المحادثة — ردود الوكيل، وردود فريقك أثناء التحويل، وتغيّر الوضع — أبقِ بث الجلسة مفتوحًا:
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. - تظهر أحداث جديدة مع الوقت. تجاهل الأنواع التي لا تتعامل معها، كما في الكود أعلاه.