Streaming de sesiuni
Urmărește o sesiune live prin SSE sau WebSocket, reia de la orice seq și scrie idempotent, ca un retry să nu dubleze niciodată.
Ambele canale livrează aceleași frame-uri. Fiecare frame poartă v: 1 și un tip t ∈ event | presence | lease | control. La conectare serverul redă întâi evenimentele stocate cu seq > after (până la 1000), apoi trimite un frame lease, apoi trece pe live. Frame-urile event duplicate sunt suprimate după seq, așa că un mesaj pe care îl ai deja nu sosește niciodată de două ori.
| SSE | WebSocket | |
|---|---|---|
| Cale | GET /v1/sessions/:id/stream?after=<seq> | GET /v1/sessions/:id/ws?after=<seq> (upgrade) |
| Direcție | Doar down | Down + up (controls, events, presence) |
| Auth | Header Authorization | Header Authorization — o cheie în query string este respinsă |
| Dispozitiv | Header x-codai-device, obligatoriu | Header, sau ?device=<uuid>&platform=<p> |
| Link share | x-codai-share-token sau ?share= | La fel |
| Forma frame-ului | event: <t> + data: {"v":1,"seq"?,…data} — data aplatizată | { "v":1, "t", "seq"?, "data": {…} } — data imbricată |
| Keep-alive | Comentariu : ping <ms> la fiecare 15 s de tăcere | — |
| Conexiune maximă | 30 min, apoi serverul închide | — |
SSE
curl -N "https://ai.codai.ro/v1/sessions/5b3e…/stream?after=6" \
-H "Authorization: Bearer $CODAI_API_KEY" \
-H "x-codai-device: 9a0b…" -H "x-codai-device-platform: web"event: event
data: {"v":1,"seq":7,"kind":"tool_call","ts":1757757852512,"sender_device_id":"1c2f…","turn_id":"t-1","client_event_id":"c9f2…","payload":{"name":"open_app"}}
event: lease
data: {"v":1,"holder_device_id":"1c2f…","expires_at":"2026-09-13T10:04:40.000Z"}
event: presence
data: {"v":1,"device_id":"9a0b…","user_id":"…","role":"editor","executor":false,"last_seen":1757757852600,"driving":false,"online":true}
event: control
data: {"v":1,"id":"ctl-8d1a…","kind":"steer","text":"Use the second result.","turn_id":"t-1","ask_id":null,"from_device_id":"9a0b…","seq":8,"applied":false}
event: control
data: {"v":1,"id":"ctl-8d1a…","kind":"steer","seq":8,"applied":true}
: ping 1757757870000Un frame lease cu "holder_device_id": null, "expires_at": null înseamnă că lease-ul este liber. Un frame control sosește de două ori per control: o dată la acceptare (applied: false, body complet) și o dată când executorul îl marchează ca aplicat (applied: true, body scurt).
EventSource din browser nu poate seta Authorization sau x-codai-device, așa că folosește un cititor SSE bazat pe fetch (fetch + ReadableStream) sau WebSocket-ul. Nu trimite niciodată cheia API în URL — protocolul interzice asta și WebSocket-ul o impune.
Un cititor minimal bazat pe fetch, care urmărește seq și se reconectează:
async function follow(sessionId: string, apiKey: string, device: string, onFrame: (t: string, d: any) => void) {
let after = 0;
for (;;) {
const res = await fetch(`https://ai.codai.ro/v1/sessions/${sessionId}/stream?after=${after}`, {
headers: { Authorization: `Bearer ${apiKey}`, 'x-codai-device': device, 'x-codai-device-platform': 'web' },
});
const reader = res.body!.pipeThrough(new TextDecoderStream()).getReader();
let buf = '';
for (;;) {
const { value, done } = await reader.read();
if (done) break; // tăierea de 30 min sau cădere de rețea → reconectare cu ultimul seq
buf += value;
let i: number;
while ((i = buf.indexOf('\n\n')) >= 0) {
const block = buf.slice(0, i);
buf = buf.slice(i + 2);
const t = /^event: (.+)$/m.exec(block)?.[1];
const raw = /^data: (.+)$/m.exec(block)?.[1];
if (!t || !raw) continue; // comentariile ": ping"
const d = JSON.parse(raw);
if (t === 'event') {
if (d.seq <= after) continue;
after = d.seq;
}
onFrame(t, d);
}
}
}
}WebSocket
const ws = new WebSocket(`wss://ai.codai.ro/v1/sessions/${id}/ws?after=${after}&device=${device}&platform=web`, {
headers: { Authorization: `Bearer ${apiKey}` }, // Node `ws`; browserele nu pot seta headere pe WebSocket
});Coduri de închidere: 4401 cheie lipsă/invalidă (sau o cheie în query string) · 4403 nu ești membru · 4404 sesiune negăsită · 4400 altă eroare de client (de ex. dispozitiv lipsă) · 1011 eroare internă.
Frame-urile down sunt obiectele SSE cu data imbricat:
{ "v": 1, "t": "event", "seq": 7, "data": { "kind": "tool_call", "ts": 1757757852512, "sender_device_id": "1c2f…", "turn_id": "t-1", "client_event_id": "c9f2…", "payload": { "name": "open_app" } } }
{ "v": 1, "t": "lease", "data": { "holder_device_id": "1c2f…", "expires_at": "2026-09-13T10:04:40.000Z" } }
{ "v": 1, "t": "presence", "data": { "device_id": "9a0b…", "role": "viewer", "executor": false, "driving": false, "last_seen": 1757757852600, "online": true } }
{ "v": 1, "t": "control", "data": { "id": "ctl-8d1a…", "kind": "send", "text": "…", "from_device_id": "9a0b…", "seq": 9, "applied": false } }Mesajele up (frame-uri text, JSON):
t | Body | Rol minim | Ack |
|---|---|---|---|
control | { "t": "control", "id", "kind", "text"?, "turn_id"?, "ask_id"? } | editor | { "v": 1, "t": "ack", "ref": <id>, "accepted": true, "seq", "duplicate" } |
events | { "t": "events", "events": IncomingEvent[1..200], "expected_last_seq"? } | deținătorul lease-ului | { "v": 1, "t": "ack", "last_seq", "events": [{ "client_event_id", "seq" }] } |
presence | { "t": "presence", "driving": bool } | viewer | Un frame presence către toți |
O eroare la un mesaj up primește răspunsul { "v": 1, "t": "error", "ref"?: <control id>, "error": { "message", "type", "code", "details" } } și socket-ul rămâne deschis. Frame-urile binare sunt ignorate.
Reluarea
Există exact un mecanism de reluare: ?after=<seq>. Frame-urile nu poartă o linie id: și Last-Event-ID nu este citit. Păstrează cel mai mare seq pe care l-ai procesat; la orice deconectare — rețea, tăierea SSE de 30 de minute, un redeploy — reconectează-te cu el și serverul redă ce ai pierdut.
Două lucruri de făcut corect:
- Replay-ul este limitat la 1000 de evenimente. Dacă
last_seq − afterpoate depăși asta (o sesiune lungă pe care nu ai urmărit-o de ceva vreme), paginează întâiGET /v1/sessions/:id/events?after=&limit=1000până ajungi lalast_seq, apoi deschide stream-ul de acolo. Altfel evenimentele dintre sfârșitul replay-ului și primul frame live nu sunt livrate niciodată. - Doar frame-urile
eventauseq. Frame-urilelease,presenceșicontrolsunt stare, nu intrări de jurnal — ia-l pe cel mai recent, nu încerca să le ordonezi față de evenimente. (Frame-ulcontrolpoartă totușiseq-ul evenimentuluicontrolpe care îl face ecou — așa le corelezi pe cele două.)
Scrierea idempotentă
Retry-urile sunt sigure când fiecare scriere poartă propria cheie:
| Scriere | Cheia ta | La retry |
|---|---|---|
POST …/events / WS events | client_event_id per eveniment | Duplicatul nu este stocat; primești seq-ul lui original în events[] și este exclus din accepted. Trimite unul pe fiecare eveniment — un eveniment fără el nu este deduplicat niciodată. Cheia este limitată la dispozitivul tău: același client_event_id de pe alt dispozitiv este un eveniment diferit. |
POST …/control / WS control | id-ul controlului | 200 { accepted: true, seq: <original>, duplicate: true } — niciun eveniment nou, nimic refăcut ecou. |
POST …/dispatch | control_id | 200 duplicate: true, niciun push retrimis. Furnizează mereu propriul control_id când e posibil să faci retry. |
POST …/lease | dispozitivul tău | Re-preluarea reîmprospătează expires_at. |
Pentru a te feri să scrii pe o vedere învechită, adaugă expected_last_seq la un batch de evenimente: dacă sesiunea a avansat, întregul batch este refuzat cu 409 seq_conflict și nu se scrie nimic — reîncarcă, apoi reîncearcă.
Bucla executorului
Executorul de referință (telefonul) se comportă astfel; ceilalți executori ar trebui să facă la fel.
- La pornire sau la o trezire prin dispatch, pentru fiecare sesiune deschisă:
POST …/lease; la succes pornește rularea agentului. - Scrie fiecare linie de trace local și trimite-o în batch la
POST …/events(sau WSevents) — flush la fiecare ≤ 200 ms sau 20 de evenimente, mereu cuclient_event_id. - Consumă
GET …/controls?applied=false(opțional&target=me), apoi abonează-te la frame-urilecontrollive. - Aplică fiecare control; când ai terminat,
POST …/control/:cid/applied. Sari controls al cărortarget_device_ideste setat și nu este propriul tău id. - Heartbeat
PUT …/leasela fiecare 10 s. La409: oprește execuția, devino viewer, afișează „condus de<device>”. - Dacă un turn are nevoie de ecran, emite
{ "kind": "screen_wait", "payload": { "position" } }ca viewerii să înțeleagă de ce sesiunea stă.
Referința protocolului
Shared Sessions v1 — obiecte, headere, fiecare rută HTTP, regulile lease-ului, idempotență, erori și limite.
Partajare și organizații
Acordă un rol pe o sesiune unei persoane, unei organizații sau unui link; listează ce este partajat cu tine; gestionează apartenența la organizații.