Глава 10 · Поток
Streaming — ответ модели как лента событий
Модель не отдаёт ответ одним JSON-объектом — она шлёт его как поток SSE-событий, которые приходят с интервалом в десятки миллисекунд. Между событиями может прилететь ошибка, которую обычный try/catch не поймает, — и harness обязан красиво её обработать.
Зачем streaming
Без streaming жизненный цикл вызова такой: harness отправляет запрос на /v1/messages и ждёт. Модель целиком генерирует ответ, складывает его в один JSON и возвращает. При длинных ответах это легко превращается в десятки секунд тишины: UI стоит, harness не может реагировать. Хуже того — SDK Anthropic в не-стримящем режиме валидирует запрос против 10-минутного таймаута, и долгие генерации без stream: true просто рвутся.
Со streaming картина другая: модель шлёт токены по мере генерации, и текст появляется на экране почти сразу после первого байта. UX оживает, harness может реагировать раньше — например, начать выполнять tool_use, как только модель его сформулировала. Та же модель и тот же ответ, но разрезанный на куски и упакованный в Server-Sent Events — текстовый стрим, где каждое сообщение начинается с event:, тело идёт в data:, разделитель — пустая строка. Wire — это WHATWG EventSource: многострочный data: склеивается через \n, строки с : в начале — комментарии-keepalive.
Анатомия SSE-потока
Один стримящийся ответ — это шесть типов событий, идущих в фиксированном порядке. Сначала message_start: в data лежит каркас будущего Message с пустым content: [], моделью, начальным usage (только input_tokens + output_tokens: 1) и stop_reason: null. На этом событии harness обычно замеряет TTFT — time-to-first-token.
Дальше идут контент-блоки. Каждый открывается content_block_start (с index и пустым телом — {"type":"text","text":""} или {"type":"tool_use","input":{}}), внутри прилетает один или несколько content_block_delta, а закрывается content_block_stop. У дельт есть подтипы: text_delta — кусочек текста; input_json_delta — кусок частичного JSON для tool_use.input (парсить на каждой дельте нельзя, это O(n²); правильно — накапливать строку и парсить один раз на content_block_stop); thinking_delta — фрагмент рассуждений; signature_delta — подпись thinking-блока перед его закрытием; citations_delta — цитата, которая дописывается в текущий text-блок.
После всех контент-блоков прилетает один или несколько message_delta с финальным stop_reason (end_turn, tool_use, max_tokens, stop_sequence, refusal…) и обновлённым usage. И в конце — message_stop, пустое событие-маркер «поток закончен». Между любыми из них могут вкраплятся ping-события — keepalive, чтобы прокси не убил соединение. Последнее правило: спецификация явно разрешает добавлять новые типы событий в будущем — клиент обязан игнорировать неизвестные, а не падать на них.
Cumulative vs additive — главная ловушка usage
Самая тихая и дорогая ошибка в streaming-клиентах живёт в одной строке документации: «The token counts shown in the usage field of the message_delta event are cumulative». То есть когда вы видите message_delta с "usage":{"output_tokens":15}, это не дельта на 15 токенов — это «всего на данный момент сгенерировано 15 токенов». Если клиент будет наивно прибавлять каждое значение к счётчику, он насчитает в разы больше реальной стоимости.
В claude-cli это разруливается одной функцией updateUsage со специальным правилом: для input_tokens, cache_creation_input_tokens, cache_read_input_tokens поле перезаписывается только если новое значение > 0 — иначе поздний message_delta с явным 0 затрёт корректное число из message_start. Для output_tokens — наоборот, перезапись безусловная: каждое message_delta приносит обновлённый итог. Глава «Cost» разбирает, как из этого usage считается итоговая стоимость; здесь важно одно — все цифры в message_delta.usage кумулятивные, а не приращения.
Ошибки в середине потока
HTTP-ответ уже пришёл со статусом 200 OK. Заголовки прочитаны, поток открыт, первые события прилетели. И посередине, между content_block_delta и следующим content_block_delta, приходит:
event: error
data: {"type": "error", "error": {"type": "overloaded_error", "message": "Overloaded"}}
Это не сетевая ошибка, не HTTP 5xx, не SDK-исключение. Это событие в SSE-стриме с типом error, и его документация Anthropic называет своим именем: «It's possible that an error can occur after returning a 200 response, in which case error handling wouldn't follow these standard mechanisms». Обычные retry-механизмы поверх HTTP его не ловят, потому что для них запрос завершился успешно. Поймать event: error может только тот код, который парсит сам поток — то есть SSE-клиент в harness'е.
Типичные in-stream-ошибки: overloaded_error (соответствует HTTP 529 в не-стримящем мире — Anthropic временно перегружен), сетевые обрывы, разрыв пайпа на стороне прокси. К ним добавляется отдельный сценарий «проксированный 200, но без событий»: HTTP-ответ пришёл, заголовки в порядке, но message_start либо не прилетел, либо за ним нет ни одного content_block_stop. Claude-cli явно ловит это после for-await как «stream ended without receiving any events» и трактует как corrupt-proxy — это тоже мид-стримная ошибка, просто молчаливая.
Recovery patterns: три стратегии
Что harness может сделать с partial-content, когда прилетел event: error? Зависит от того, что уже накоплено — есть три рабочие стратегии.
Rollback. Самая безопасная: откатить накопленный текст, выбросить partial-message, повторить запрос с теми же messages. Поток начнётся с чистого message_start, пользователь увидит ответ заново. Тонкость: блоки tool_use и thinking невосстановимы частично — у них атомарные структуры с подписями. Если ошибка случилась внутри них, rollback — единственный вариант.
Retry с продолжением. Если накоплен полезный кусок текста, его можно сохранить и попросить модель продолжить. Для Claude 4.5 и ранее канонический способ — добавить partial-ответ в messages как assistant-сообщение и снова запросить стрим. Для Claude 4.6+ схема меняется: partial добавляется в user-сообщение с явной инструкцией «ваш предыдущий ответ оборвался на [текст]. Продолжите с того места». Resume работает только по границе последнего text-блока — внутри tool или thinking продолжить нельзя.
Fallback на меньшую модель. Если overloaded_error прилетает на главной модели, имеет смысл попробовать ту же задачу на меньшей (Opus → Sonnet → Haiku). Худший по качеству, но самый быстрый путь дать пользователю хоть какой-то ответ. В claude-cli живёт отдельная ветка fallback'а — tengu_streaming_fallback_to_non_streaming: при corrupt-proxy и idle-timeout harness переключается на не-стримящий запрос с теми же параметрами. Эта ветка отключаема через CLAUDE_CODE_DISABLE_NONSTREAMING_FALLBACK, потому что однажды fallback под StreamingToolExecutor привёл к двойному выполнению тулзы. Recovery — всегда trade-off, и оператору нужен переключатель.
Watchdog: тишина — тоже отказ
Не все обрывы громкие. Бывает так: стрим открыт, HTTP-соединение живо, события приходили — а потом замолчали. Никакого event: error не было, провайдер просто «завис». Снаружи это выглядит как «модель долго думает», а внутри — мёртвое соединение, которое будет держать ресурсы, пока кто-то его не оборвёт.
Claude-cli защищается двумя независимыми детекторами. Stall-detector (30 с): отслеживает lastEventTime и эмитит информационный tengu_streaming_stall в дашборд — «канал медленный», без действия. Watchdog (90 с по умолчанию, настраивается CLAUDE_STREAM_IDLE_TIMEOUT_MS): на половине таймаута пишет warning, на полном — вызывает releaseStreamResources() и бросает исключение, чтобы триггернуть fallback. Разделение важно: пик медленности — нормально; полторы минуты тишины — broken-proxy. Параллельный контракт: request-id приходит в HTTP-заголовке ответа, и его нужно сохранить до чтения тела стрима — на ошибке внутри потока его уже не достанешь, а без него инцидент в support не разобрать.
Special errors: prompt_too_long, max_tokens, model_context_window_exceeded
Не все ошибки прилетают через event: error. Часть приходит «в норме потока», но через message_delta.stop_reason — и трактуется как recoverable. max_tokens — модель упёрлась в лимит выходных токенов и оборвалась посреди мысли; recovery — поднять лимит и попросить продолжить (claude-cli делает это через escalating-retry с системным «продолжи без извинений»). model_context_window_exceeded — весь массив messages с генерируемым ответом не помещается в окно; recovery — сжать контекст (reactive-compact, удаление старых сообщений; механики разобраны в главе «Context Window»). prompt_too_long может прилететь до открытия потока (HTTP 400) или в самом начале — лечится тем же compact'ом. В claude-cli все три превращаются в один внутренний createAssistantAPIErrorMessage с apiError: 'max_output_tokens', который верхний loop ловит как «не падать, а пойти на recovery-итерацию». Для оркестратора цикла из главы «Цикл» это просто ещё одно значение stop_reason в его switch.
It's possible that an error can occur after returning a 200 response, in which case error handling wouldn't follow these standard mechanisms.
— Anthropic · «Errors» / «Streaming messages» docsИсточники главы
- [ant] Streaming messages — Anthropic Docs (event-таксономия, cumulative usage, in-stream errors, error recovery 4.5/4.6+)
- [cli] Конспект: Streaming и SSE (watchdog 90 с, stall 30 с, corrupt-proxy detection, fallback к non-streaming, cumulative-vs-additive дисциплина в
updateUsage) - [cli] Конспект: Anthropic API distill (раздел «Streaming events and partial message behavior» — каноничный event flow, WHATWG
EventSource, «handle unknown event types gracefully», resume contract) - [ant] Errors — Anthropic Docs («200 then error» формулировка,
overloaded_error= HTTP 529 equivalent, request-id в headers)