Назад к разделу Technology

Стриминг ИИ-агента: это задача про протокол

Ответ в чате выглядит как сокет, из которого сыплются слова. В продакшене это протокол: аддитивные типы событий, идентификатор, приходящий последним, правило ретрая, которое знает, успел ли выполниться инструмент, и бюджет таймаутов, где каждый слой обязан пережить самый медленный ход. Здесь описано, что на самом деле отправляет наш чат-стрим и зачем нужна каждая его часть.

Категория
Общее
Обновлено
Автор
Stan Kharlap

Демо-версия чата с ИИ собирается за полдня. Открываешь потоковый ответ, пробрасываешь токены, текст появляется буква за буквой, все в восторге. Потом ты ставишь это перед реальными пользователями с реальной работой и узнаёшь, что стриминг никогда и не был фичей. Фича, это протокол, и почти всё интересное в нём про то, что происходит, когда счастливый путь не случается.

Встроенный ассистент Norman, это агент, который пользуется инструментами. Он может посмотреть твои транзакции, подготовить счёт, собрать платёж по счёту поставщика, объяснить, что такое НДС-номер. За время работы в продакшене мы записали десятки тысяч ходов чата и почти сотню тысяч трассированных запусков агентов по всем нашим ИИ-сценариям. Ходы чата, самые медленные и самые заметные из них, и у них форма, с которой наивный поток токенов справляется плохо:

  • Медианный ход завершается примерно за две секунды.
  • 90-й процентиль, около пятнадцати секунд.
  • 99-й процентиль, около тридцати секунд, а самые тяжёлые из записанных ходов длятся больше минуты.

Примерно каждый пятый ответ вызывает хотя бы один инструмент. Именно они и есть медленные, и именно в них пользователь сильнее всего ждёт результата, потому что он попросил ассистента что-то сделать, а не что-то объяснить.

Это распределение и есть всё техзадание. Протокол, которому комфортно на двух секундах и который разваливается на тридцати, не годится, потому что важны как раз тридцатисекундные ходы.

Транспорт нарочно скучный

Наш бэкенд, это синхронное Django-приложение. Цикл агента асинхронный. Эти два факта не хотят дружить, и есть известное искушение переписать половину стека, чтобы они подружились.

Мы не стали. Стриминговая вьюха запускает асинхронный цикл агента в daemon-потоке, а этот цикл складывает готовые строки в обычную потокобезопасную очередь. Генератор ответа работает в обычном синхронном воркере, разгребает очередь и отдаёт каждую строку в HTTP-ответ, пока не увидит sentinel:

def generate_streaming_content():
    q: sync_queue.Queue = sync_queue.Queue()

    async def run_async() -> None:
        try:
            async for chunk in wrapper.process_assistant_by_prompt_streamed(...):
                q.put(chunk + "\n")
        except Exception:
            logger.exception("Streaming view error")
            q.put(json.dumps({"type": "error", "message": "..."}) + "\n")
        finally:
            q.put(None)  # sentinel

    threading.Thread(target=lambda: asyncio.run(run_async()), daemon=True).start()

Полезная нагрузка, это JSON, разделённый переводами строк. Одно событие на строку, никаких трюков с фреймингом, никаких недособранных объектов. Клиент держит последнюю неполную строку в буфере и парсит только целые. Это четыре строки клиентского кода и единственное правило фрейминга во всём протоколе.

Мне нравится этот слой именно потому, что про него нечего рассказать. Он ни разу не был тем, что сломалось.

Один потоковый ход: асинхронный цикл агента наполняет очередь, синхронный генератор разгребает её в HTTP-ответ из построчного JSON, а браузер буферизует неполные строки. По проводу идут событие tool_start с именем работающего инструмента, дельты текста с ответом по токенам, необязательная карточка действия, событие completion, завершающее видимый ответ, и приходящее последним событие message_saved с идентификатором сохранённой записи.
Транспорт скучный намеренно. Все интересные решения, это какие события существуют и в каком порядке им разрешено приходить.

Каждый тип события аддитивен

Поток несёт небольшой набор типизированных событий: дельты текста, старт инструмента, изображение, карточку действия, завершение, ошибку и финальное подтверждение сохранения. Клиент разбирает их по полю type обычной цепочкой if и молча игнорирует всё, чего не знает.

Именно это свойство, незнакомые события пропускаются, а не роняют разбор, позволяет нам развивать протокол без согласованного релиза. Веб, iOS и Android читают один и тот же поток и обновляются по совершенно разным графикам. Если бы новый тип события означал ошибку парсинга, каждое улучшение ассистента превращалось бы в релизный поезд.

Оба события, которые мы выкатили на этой неделе, это доказательство, что схема работает. Ни одному из них не требовалось изменение клиента, чтобы деплой был безопасным, изменение нужно было только чтобы событие стало полезным.

Проблема в тишине, а не в медленности

Вот о чём не предупреждают, когда ты добавляешь агенту инструменты: вызов инструмента не стримит ничего. Модель решает вызвать инструмент, и ход замолкает целиком на всё время его работы. Ни токенов, ни прогресса, ничего в проводе.

Все эти десятки секунд наш интерфейс показывал индикатор набора текста. В такой момент индикатор набора, это враньё. Модель не печатает, она читает три месяца транзакций, а пользователь не может отличить «усердно работает» от «повисло».

Лечится это не улучшенной анимацией, а событием. Цикл стриминга и так знал имя инструмента, потому что использует его, чтобы раскладывать вывод инструмента по карточкам действий, так что мы начали его отправлять:

if event.item.type == "tool_call_item":
    self.current_tool_id = _tool_name_of(event.item.raw_item)
    # Ход может провести десятки секунд внутри вызовов инструментов,
    # не отправив ни токена, и без этого события интерфейс показывает
    # только индикатор набора, а пользователь не понимает, идёт ли работа.
    if self.current_tool_id:
        yield json.dumps({"type": "tool_start", "tool_name": self.current_tool_id})

Теперь клиент может сказать «Создаю счёт» вместо трёх анимированных точек. Та же задержка, совершенно другой опыт. Медленная операция, которая говорит, чем занята, терпима. Быстрая, которая молчит, нет.

Идентификатор приходит после текста, и это намеренно

У наших ответов в чате есть кнопки «нравится» и «не нравится». Реакции адресуют сообщение по его сохранённому идентификатору.

У потокового ответа его нет. Строка в базе появляется только после того, как поток закончился, потому что до этого нет финального текста. Значит, клиент рисовал сообщение без идентификатора, и единственные ответы, на которые реально можно было поставить реакцию, это те, что пришли при перезагрузке истории. Сигнал обратной связи, который мы хотели больше всего, на практике был почти недостижим. Продакшен-данные говорят ровно это: реакций за всю историю продукта почти ноль, и это очень вежливый способ сказать, что кнопка была декоративной.

Лечится финальным событием. После completion мы сохраняем сообщение и отправляем его идентификатор:

saved_message_id = await self._save_assistant_message_async(accumulated_message)
# Отправляется после "completion" намеренно: строка появляется только
# когда поток закончен, и клиенты не должны ждать её, чтобы дорисовать ответ.
yield json.dumps({"type": "message_saved", "message_id": saved_message_id})

Комментарий про порядок здесь и есть главное. Аккуратнее было бы отправить идентификатор первым и иметь одно авторитетное событие со смыслом «ответ готов, вот всё про него». Но это означало бы, что неудачная запись строки может помешать показать совершенно нормальный ответ. Отрисовка никогда не должна зависеть от сохранения. Поэтому ответ завершается, а идентификатор догоняет мгновением позже и включает вторичную возможность.

Ретрай безопасен только до первого инструмента

Потоки рвутся. Соединения с провайдером отваливаются, провайдеры отдают временные серверные ошибки, сети ведут себя как сети. Очевидная реакция, повторить ход.

Очевидная реакция опасна. Наш ассистент создаёт счета, транзакции и платежи. Повтор хода, который уже выполнил create_transaction, даёт не лучший ответ, а две транзакции.

Поэтому правило ретрая опирается на два состояния, которые цикл стриминга и так ведёт:

is_transient = isinstance(e, APIError) and getattr(e, "type", None) == "server_error"
if is_transient and attempt < MAX_STREAM_RETRIES and not tool_executed:
    await asyncio.sleep(attempt + 1)
    continue

tool_executed становится истиной в момент, когда любой вызов инструмента вернул результат. После этого ход перестаёт быть чистой функцией от входа, и притворяться иначе нам нельзя. То же правило действует для оборванного соединения, с одной добавкой: мы не повторяем и тогда, когда часть текста уже ушла в поток, потому что пользователь её видел, а начинать ответ заново, это отдельный вид поломки.

Это тот же инстинкт, что проходит через всю остальную систему. Если что-то уже коснулось внешнего мира, единственный безопасный ход, это вперёд.

Частичный ответ лучше красной ошибки

Если повторить нельзя, мы спасаем. Когда поток обрывается после того, как пришёл настоящий контент, мы не выбрасываем текст. Мы сохраняем то, что есть, отправляем событие completion с ним и следом идентификатор сообщения, ровно как при успешном ходе. Пользователь видит короткий ответ вместо ошибки, и он остаётся в истории как любое другое сообщение.

Видимой ошибкой становится только обрыв, при котором не накопилось ничего, и даже тогда это простая фраза с просьбой попробовать ещё раз, а не стек-трейс.

Принцип такой: ход, который что-то произвёл, заканчивается сообщением, а ход, который не произвёл ничего, заканчивается извинением. Третьего случая, где пользователю приходится гадать, нет.

Соединение, о котором ты забыл

Мой любимый баг в этой системе не имел к модели никакого отношения.

Сохранять сообщения из асинхронного цикла означает, что ORM работает в потоках экзекьютора. Django перерабатывает соединения с базой по сигналам начала и конца запроса, а эти сигналы никогда не срабатывают для потоков, которые ты запустил сам. Поэтому длинный ход может оставить соединение из пула простаивать всё время вызова модели, достаточно долго, чтобы сервер на той стороне его закрыл. Следующий запрос падает на том, что соединения уже нет, внутри в остальном совершенно здорового реквеста.

Починка маленькая и конкретная:

def _retry_on_stale_connection(func):
    """Восстановление после соединения с Postgres, потерянного в простое.

    Каждый обёрнутый хелпер делает ровно один create или save без
    закоммиченных до этого побочных эффектов, поэтому выбросить мёртвое
    соединение и повторить ровно один раз, это идемпотентно.
    """
    @wraps(func)
    def wrapper(*args, **kwargs):
        try:
            return func(*args, **kwargs)
        except (OperationalError, InterfaceError):
            close_old_connections()
            return func(*args, **kwargs)
    return wrapper

Обрати внимание на обоснование в докстринге, а не только на механизм. Ретрай разрешён потому, что обёрнутая операция, это одна запись без закоммиченных изменений до неё. В этом разница между безопасным ретраем и ретраем на удачу, и это стоит записать там, где следующий человек это прочитает.

Каждый таймаут в стеке обязан пережить самый медленный ход

Вернёмся к распределению задержек: p99 около тридцати секунд. А теперь представь таймаут запроса в воркере в двадцать пять секунд, вполне разумное значение по умолчанию для транзакционного API.

Эта комбинация, машина по производству загадочных сбоев. Воркер убивают посреди хода, асинхронный цикл умирает без уборки, осиротевшие подпроцессы инструментов остаются висеть, а пользователь получает общую ошибку ровно на тех запросах, где работы было больше всего. В логах приложения не будет слова «таймаут», потому что процесса, который бы его записал, уже нет.

Урок обобщается в правило, которое мы теперь применяем осознанно:

  1. Таймаут запроса должен превышать самый медленный ход, который ты готов обслужить, а не средний.
  2. Каждый внутренний таймаут, особенно на вызовы инструментов, должен быть короче таймаута запроса, чтобы медленный инструмент падал ловимым исключением, а не убитым процессом.
  3. Стриминговым эндпоинтам не место в том же бюджете воркеров, что и транзакционным. Запрос, который держит воркер тридцать секунд, и запрос, который отвечает за восемьдесят миллисекунд, не должны конкурировать за одни и те же слоты.

Есть и четвёртое правило, которое мы сейчас выкатываем, и оно наименее очевидное: отправь байт пораньше. Прокси перед API обычно повторяют запрос, который ещё не выдал ни одного байта ответа. Для обычного эндпоинта это полезная страховка. Для стримингового, который может двадцать секунд думать до первого токена, это риск двойного выполнения. Тривиальный первый чанк, отправленный до старта агента, закрывает это окно.

Модель, это простая часть

Ничего из этого не про промптинг. Каждое сложное решение здесь, это старое решение из распределённых систем в новой одежде: что идемпотентно, что можно повторять, что можно переставлять местами, что никогда не должно блокировать отрисовку и какой таймаут в цепочке самый короткий.

Если ты строишь такое, самая короткая полезная выжимка, которую я могу дать: запиши контракт событий до того, как напишешь цикл стриминга, реши для каждого события, обязательное оно или аддитивное, и реши для каждого сбоя, во что он превращается, в сообщение или в извинение. Потоковая передача токенов заработает с первого раза. Всё остальное и есть настоящий продукт.

Norman берет операционную финансовую работу на себя

От invoicing до bookkeeping: Norman организует повторяющиеся финансовые процессы так, чтобы вы успевали к дедлайнам с меньшим объемом ручной работы.