Как унификация SSE событий LLM удерживает стриминг в порядке
Унификация SSE событий LLM помогает безопасно собрать дельты текста, tool call, ошибки и финальные статусы разных API в один контракт.

Поток от LLM нельзя сводить к циклу for chunk: print(chunk.text). Такой код переживает демо, а в продакшене теряет аргументы инструментов, путает отказ с ошибкой сети и объявляет ответ завершенным только потому, что сокет закрылся.
Нормальный общий слой для стриминга не пытается сделать всех провайдеров похожими. Он фиксирует небольшую последовательность собственных событий, сохраняет порядок блоков и честно отличает корректный финал от аварийного обрыва. Это дает клиентам один контракт, а адаптерам свободу разбирать чужие особенности на границе системы.
SSE задает рамку, но не смысл событий
SSE определяет текстовый способ доставлять сообщения по HTTP с Content-Type: text/event-stream. В HTML Standard каждое сообщение состоит из строк вроде event:, data: и пустой строки, которая завершает запись. Стандарт также допускает несколько строк data: в одном сообщении и идентификатор id, который браузер может вернуть при переподключении.
Из этого не следует, что SSE умеет стримить ответ модели. Он не говорит, что такое дельта, когда открыт вызов инструмента, можно ли переставлять фрагменты местами и что считать успешным окончанием. Это семантика API провайдера.
На практике встречаются как минимум четыре разные модели:
- OpenAI-совместимый Chat Completions поток часто передает JSON в
data:и завершает его строкой[DONE]. - OpenAI Responses API отдает типизированные события вроде
response.created, событий добавления дельт иresponse.completedлибоresponse.failed. - Claude Messages API передает именованные SSE-события и строит ответ через жизненный цикл content block.
- Gemini
streamGenerateContentотдает последовательность объектовGenerateContentResponse, где один объект может содержать очередную часть кандидата или финальные сведения.
Ошибка начинается с названия. Команда говорит «у нас SSE», а потом применяет к Claude обработчик, рассчитанный на choices[0].delta.content, или ожидает [DONE] от Gemini. На проводе у всех может быть SSE, но протокол ответа у них разный.
Общий формат должен описывать блоки, а не токены
Единая дельта текста удобна, пока приложение только печатает чат. Как только модель вызывает функцию, возвращает структурированный объект, выдает отказ или ведет скрытую рассудочную часть, «токен» перестает быть единицей интеграции.
Возьмите за основу блок контента. Блок имеет стабильный block_id, порядковый index и тип. Внутри него приходят дельты. Текстовый блок получает строки, блок вызова инструмента получает имя, идентификатор вызова и фрагменты JSON, а reasoning-блок может получать текст либо отдельную подпись для проверки целостности.
Минимальный контракт, который я использую между шлюзом и приложением, выглядит так:
{
"stream_id": "st_01J...",
"seq": 17,
"type": "block.delta",
"block": {
"id": "b_1",
"index": 1,
"kind": "tool_call"
},
"delta": {
"kind": "json_text",
"text": "{\"city\":\"Alma"
}
}
Для этого контракта нужны шесть обязательных типов:
stream.startсообщает, что шлюз принял поток и назначил ему идентификатор.block.startоткрывает конкретный блок и объявляет его тип.block.deltaдобавляет неизменяемый фрагмент в уже открытый блок.block.stopзакрывает блок и запрещает новые дельты для него.stream.doneподтверждает успешное завершение и несет финальные метаданные.stream.errorзавершает поток ошибкой с кодом, категорией и безопасным сообщением.
Событие stream.start не равно «провайдер начал генерировать». Оно означает только, что ваш публичный контракт начался. Событие stream.done не равно «клиент дочитал HTTP-body». Оно означает, что адаптер получил достаточное подтверждение успешного результата.
Можно добавить stream.heartbeat и stream.warning, но не делайте их обязательными для клиента. Пинг, комментарий SSE и очередная текстовая дельта не должны менять состояние ответа.
Порядок дельт важнее размера фрагмента
Провайдер не обещает удобный размер порций. Одна дельта может содержать слово, несколько предложений, пустую строку, половину Unicode-символа на уровне байтового чтения или кусок JSON без закрывающей скобки. Ваш интерфейс не должен выдавать предположения о границах токенов.
Надежное правило простое: адаптер передает дельты в том порядке, в котором получил их для одного логического блока, и не объединяет разные блоки ради удобства UI. Если блоки имеют индексы, индекс определяет место в финальном массиве контента, а seq определяет порядок доставки событий клиенту.
Рассмотрим ответ, где модель сначала пишет текст, затем вызывает функцию, после результата функции пишет продолжение. Корректная последовательность такая:
stream.start
block.start text index=0
block.delta text="Проверяю расписание."
block.stop text index=0
block.start tool_call index=1
block.delta json_text="{\"date\":\"2026-07-23\""
block.delta json_text="}"
block.stop tool_call index=1
block.start text index=2
block.delta text="На сегодня доступно..."
block.stop text index=2
stream.done
Нельзя показывать пользователю второй текстовый блок раньше, чем приложение выполнило вызов, даже если конкретный провайдер умеет параллельно генерировать части ответа. Нельзя склеивать index=0 и index=2 в один буфер и рассчитывать позже восстановить смысл. Аудит, повторный запуск и инструментальная оркестрация потребуют исходную структуру.
Отдельно решите судьбу reasoning. Если политика продукта запрещает показывать его пользователю, не подменяйте reasoning обычным текстом и не выкидывайте его молча. Передайте блок с visibility: "internal" в защищенный контур или отключите выдачу на уровне запроса. У разных моделей эта часть ответа имеет разные правила доступа и проверки, поэтому универсальный UI-флаг здесь опасен.
Завершение потока требует доказательства от провайдера
Закрытое соединение означает только закрытое соединение. Сервер мог штатно закончить ответ, прокси мог оборвать долгий запрос, пользователь мог потерять сеть, а upstream мог упасть после передачи половины JSON-аргументов. Эти случаи нельзя сводить к одному done.
У каждого адаптера должна быть явная таблица финальных сигналов. Например, для OpenAI-совместимого Chat Completions это может быть [DONE] после валидных JSON-чанков. Для OpenAI Responses API нужно опираться на терминальные response.completed, response.failed, response.incomplete или response.cancelled, а не только на конец HTTP-потока. В документации OpenAI события несут sequence_number, и это полезный внешний материал для диагностики, но ваш публичный seq все равно должен назначать шлюз.
Документация Claude формулирует жизненный цикл особенно ясно: message_start, затем один или несколько блоков через content_block_start, content_block_delta и content_block_stop, после этого message_delta и финальный message_stop. Claude также прямо предупреждает, что новые типы событий могут появляться, а клиент обязан обрабатывать неизвестные типы спокойно. Это хорошее правило для любого адаптера, а не только для Claude.
Gemini в режиме streamGenerateContent передает последовательность объектов ответа, а не обязательный финальный маркер наподобие [DONE]. Адаптеру нужно смотреть на документированные поля завершения кандидата и финальные данные, которые пришли в потоке. Само закрытие после последнего объекта нельзя считать достаточным подтверждением, если перед ним не было данных, по которым можно определить статус.
Внутренне держите три терминальных состояния:
completedозначает, что провайдер подтвердил штатный результат.failedозначает, что провайдер прислал ошибку либо адаптер получил однозначную ошибку протокола.interruptedозначает, что транспорт оборвался или клиент отменил запрос до подтвержденного финала.
Для interrupted можно сохранить уже полученные текстовые дельты как черновик, но нельзя записывать такой ответ как готовый результат агента. Особенно если последний открытый блок был tool_call или JSON-ответ по схеме.
Аргументы инструментов нужно собирать до закрытия блока
Частичный JSON не является JSON. Это кажется банальностью, но именно здесь многие оркестраторы вызывают функцию с обрезанным аргументом или начинают «чинить» модельный вывод регулярными выражениями.
Claude в input_json_delta присылает строковые фрагменты partial_json и рекомендует собирать их до разбора. Документация Gemini для структурированного вывода тоже говорит о валидных частичных JSON-строках, которые надо конкатенировать до получения полного объекта. Это один и тот же инженерный вывод: поток можно отображать инкрементально, но исполняемый объект появляется только после завершения блока.
У адаптера должен быть буфер на каждый block_id, а не один буфер на весь ответ. Вызовы инструментов могут идти несколькими блоками, а некоторые API разрешают параллельные вызовы.
type ToolBuffer = {
name?: string;
callId?: string;
rawArguments: string;
closed: boolean;
};
function appendToolDelta(buf: ToolBuffer, part: string) {
if (buf.closed) throw new Error("delta after block.stop");
buf.rawArguments += part;
}
function closeToolBlock(buf: ToolBuffer) {
buf.closed = true;
const args = JSON.parse(buf.rawArguments);
return { name: buf.name, callId: buf.callId, arguments: args };
}
Этот код намеренно не пытается разбирать rawArguments на каждой дельте. Если UI хочет показать ход формирования вызова, он может показать сырой текст в техническом режиме. Исполнитель инструмента должен дождаться block.stop, затем проверить JSON и схему аргументов.
Не маскируйте ошибку JSON.parse как ошибку инструмента. Это ошибка протокола ответа или несовместимости модели с режимом structured output. В журнале должны остаться ID провайдера, ID модели, исходные события и позиция блока, но не секреты пользователя и не полный промпт по умолчанию.
Ошибка внутри SSE должна дойти до клиента как ошибка
HTTP-статус 200 в начале ответа не гарантирует успешную генерацию. После отправки заголовков сервер уже не может заменить ответ на HTTP 429, 500 или 529. Поэтому провайдеры могут передавать ошибку отдельным событием в середине уже открытого потока.
Claude прямо описывает event: error с объектом ошибки, например overloaded_error. OpenAI Responses API имеет отдельные терминальные события неуспеха. В других API ошибка может проявиться как разрыв транспорта с диагностикой в SDK. Общий слой обязан превратить эти варианты в одно stream.error.
Полезная структура ошибки не должна притворяться, что все ошибки одинаковы:
{
"type": "stream.error",
"stream_id": "st_01J...",
"seq": 24,
"error": {
"category": "upstream_overloaded",
"retryable": true,
"provider_code": "overloaded_error",
"message": "Провайдер временно перегружен"
}
}
После stream.error закрывайте все незакрытые блоки только во внутреннем состоянии, но не отправляйте клиенту фиктивные block.stop. Иначе клиент решит, что аргументы инструмента целы, а текстовый ответ закончен осмысленно.
Повторять запрос автоматически можно только до выдачи внешне наблюдаемого эффекта. Для обычного чата это значит до первой дельты, которую увидел пользователь. Для агента с инструментами граница строже: нельзя без идемпотентного ключа повторять запуск после того, как инструмент мог списать деньги, отправить письмо или изменить запись.
Неизвестные события должны переживать обновления API
Провайдеры расширяют стриминговые протоколы. Новое событие о reasoning, цитате, аудио или серверном инструменте не должно валить обработчик switch с исключением и обрывать пользователю уже готовый текст.
Разделите обработку на два слоя. Первый слой декодирует SSE и сохраняет исходное событие в диагностический след. Второй слой сопоставляет известные типы с вашим контрактом. Неизвестный тип он помечает как ignored, увеличивает метрику и продолжает поток, если сам провайдер не назвал событие терминальной ошибкой.
Это не призыв молча игнорировать изменения навсегда. Неизвестный тип нужно видеть в алертах и тестовых транскриптах. Но падение рабочего потока из-за нового необязательного поля или события ping хуже, чем контролируемое пропускание с наблюдаемостью.
События heartbeat также не должны продлевать бизнесовый таймаут бесконечно. Разделяйте транспортный таймаут, который показывает живое соединение, и таймаут прогресса, который требует появления содержательной дельты или документированного события выполнения инструмента. Иначе upstream может держать соединение пингами, а ваш пользователь будет ждать без конца.
Адаптер должен быть конечным автоматом
Условий в стриминге быстро становится слишком много для набора if. Небольшой конечный автомат делает запрещенные переходы заметными и дает нормальные тесты.
Состояние потока можно описать так:
idle -> open -> receiving -> terminal
| |
v v
interrupted completed | failed
Внутри receiving храните состояние каждого блока: new, open, closed. Разрешайте block.delta только для open, а stream.done только когда все блоки закрыты. Если провайдер прислал терминальное событие с открытым блоком, адаптер должен завершить поток ошибкой совместимости, а не додумывать недостающий JSON.
Ниже пример правил, которые стоит проверять для каждого записанного потока:
- Первый публичный элемент всегда
stream.start. - У каждого
block.deltaесть ранее открытый и еще не закрытый блок. - Номер
seqстрого возрастает в рамках одногоstream_id. - После
stream.doneилиstream.errorнет новых публичных событий. stream.doneне выходит, пока есть открытый блок.
Эти правила ловят ошибки, которые не видно глазами: повторную дельту после закрытия, перепутанный индекс при параллельном tool call, финал до прихода usage и ошибочное повторное применение чанка после reconnect.
Тестировать нужно транскрипты, а не только живые запросы
Живой тест против модели нужен, но он плохо воспроизводит редкие сбои. Настоящая защита появляется, когда адаптер проходит сохраненные транскрипты реальных протокольных событий.
Соберите для каждого провайдера хотя бы пять обезличенных последовательностей: простой текст, несколько блоков, вызов инструмента с фрагментированным JSON, терминальная ошибка после части текста и обрыв без терминального сигнала. Добавьте случай с неизвестным событием. Для Claude добавьте ping, поскольку он может появиться в любом месте. Для OpenAI Responses добавьте неуспешный терминальный статус. Для Gemini добавьте поток, где текст приходит несколькими объектами ответа.
Проверяйте не только ожидаемый финальный текст. Проверяйте последовательность нормализованных событий целиком. Вот форма такого теста:
expect(normalize(transcript)).toEqual([
{ type: "stream.start", seq: 1 },
{ type: "block.start", seq: 2, block: { index: 0, kind: "text" } },
{ type: "block.delta", seq: 3, delta: { text: "Привет" } },
{ type: "block.stop", seq: 4 },
{ type: "stream.done", seq: 5, status: "completed" }
]);
Добавьте property tests поверх транскриптов. Генератор может разрезать одну текстовую дельту на случайные фрагменты, вставлять heartbeat и повторять неизвестные необязательные события. Финальный собранный текст обязан остаться тем же, а автомат не должен нарушать инварианты.
AI Router дает командам OpenAI-совместимый вход для нескольких поставщиков, но совместимость endpoint не отменяет различий в семантике стриминга. Если вы строите шлюз или клиент поверх такого входа, проверьте, какие события сохраняет и нормализует именно ваш маршрут, особенно для инструментов и финальных статусов.
Клиенту нужен простой контракт, а шлюзу нужна полная правда
Фронтенду обычно достаточно текста, статуса и прогресса. Оркестратору нужны индексы блоков, идентификаторы tool call, причины остановки, usage и исходный статус провайдера. Не заставляйте браузер держать эту сложность, но и не выбрасывайте ее в шлюзе.
Хороший раздел обязанностей выглядит так: адаптер читает чужой поток и строит строгие внутренние события, оркестратор принимает решения по блокам и инструментам, а клиент получает безопасную проекцию того, что можно показать. Один и тот же нормализованный журнал при этом должен позволять объяснить, почему ответ закончился, где прервался JSON и что реально пришло от модели.
Если ваш текущий контракт состоит из delta, [DONE] и catch, начните не с переписывания всех интеграций. Возьмите один поток с tool call, зафиксируйте его как транскрипт, добавьте block.start и block.stop, затем запретите done без подтвержденного терминального события. После этого большинство скрытых различий между LLM API перестанет всплывать прямо в пользовательском интерфейсе.
Часто задаваемые вопросы
Чем SSE отличается от формата событий конкретного LLM API?
SSE задает способ передавать события по HTTP: поля event, data, id и границы сообщений. Он не определяет, что означает конкретный фрагмент данных, поэтому один провайдер присылает текстовую дельту, другой открывает и закрывает блоки контента, а третий отдает очередные снимки ответа.
Обязателен ли маркер [DONE] в LLM-стриминге?
Нет. [DONE] стал привычным маркером в ряде OpenAI-совместимых потоков, но это соглашение протокола верхнего уровня, а не часть SSE. Если ваш адаптер считает его единственным доказательством завершения, он неверно обработает потоки с явным финальным событием или с обычным закрытием соединения.
Когда шлюз должен отправлять клиенту событие done?
Событие done стоит отправлять только после того, как адаптер получил подтверждение корректного завершения от провайдера и собрал финальные метаданные. Не привязывайте его к закрытию TCP-соединения: сеть может оборваться после уже показанного пользователю текста.
Можно ли парсить JSON аргументы tool call на каждой дельте?
Нет. Частичный JSON внутри аргументов инструмента может обрываться в середине строки, escape-последовательности или вложенного объекта. Накопите байты для конкретного блока вызова, разберите объект после его закрытия и только затем передавайте инструменту.
Какие поля нужны в общем событии стриминга?
Нужны как минимум идентификатор потока, порядковый номер, тип события, индекс блока, тип содержимого и полезная нагрузка. Для финального события добавьте причину остановки, usage при наличии и статус завершения. Полезно хранить исходный тип события провайдера в диагностическом поле, но не заставлять клиента от него зависеть.
Что делать, если SSE-соединение оборвалось после нескольких токенов?
Показывайте уже подтвержденный текст, но относитесь к нему как к необратимой части текущего ответа. При повторном подключении не просите провайдера "продолжить с последней дельты", если у него нет документированного механизма возобновления. Надежнее создать новый запрос с явной политикой повтора и пометить первый запуск как прерванный.
Может ли один LLM-ответ содержать несколько блоков контента?
Это зависит от модели и API, но общий интерфейс должен допускать несколько независимых блоков: текст, reasoning, вызовы инструментов, отказ и медиа. Склеивать все в одну строку удобно только для простого чата, а потом ломает порядок инструментов и аудит результата.
Можно ли получить usage до завершения ответа?
Не всегда. Некоторые API передают usage в финальном событии, некоторые отдают накопительные счетчики по ходу потока, а некоторые дают их только в нестриминговом ответе. Сохраняйте статус known, estimated или absent, иначе финансовая отчетность начнет выдавать догадки за факт.
Подходит ли браузерный EventSource для прямого вызова LLM API?
Обычный EventSource удобен для браузерного GET с автоматическим переподключением, но многие LLM API требуют POST, заголовки авторизации и тело запроса. На сервере используйте потоковый HTTP-клиент, а в браузере обычно держите свой backend как посредника или используйте fetch с чтением ReadableStream.
Как тестировать адаптеры стриминга для разных LLM?
Сначала зафиксируйте контракт событий и напишите транскрипты реальных потоков для каждого провайдера: обычный текст, tool call, отказ, лимит токенов и обрыв сети. Затем прогоните один и тот же набор инвариантов через каждый адаптер. Ручная проверка в чат-интерфейсе почти ничего не ловит.