Skip to content

Асинхронность и стриминг

AsyncBroker — основной движок; Broker — его блокирующая обёртка. Методы те же, только с await:

async with llmbroker.AsyncBroker() as broker:
    reply = await broker.ask("Привет")
    print(reply.text)

Асинхронный tool-цикл — await llmbroker.arun_tool_loop(...), см. Инструменты и агенты.

Стриминг

Стриминг — только асинхронный: пул отдаёт дельты по мере поступления, с роутингом и фейловером:

stream = broker.stream("Напиши хокку про брокеров", operation="write")
async for delta in stream:
    print(delta, end="", flush=True)

print(stream.llm_name, stream.usage)   # кто ответил и во что это обошлось
await stream.record_quality(0.9)       # оценить, не называя вызов самому

stream(...) возвращает handle: по нему вы итерируетесь за дельтами, а он же называет ответившую модель и — когда ответ закончился — во что он обошёлся.

Фейловер работает как обычно вплоть до первой дельты: зарейтлимиченная или сломанная модель уходит в кулдаун, её место незаметно занимает следующая. Модель, чей ответ закончился, так и не дав ни одной дельты, сломана в том же смысле: до вас ничего не дошло, поэтому брокер переключается и на ней, а не отдаёт вам пустой стрим. Но как только текст пошёл, переключаться уже некуда — поэтому стрим, оборвавшийся посреди ответа, бросает StreamInterruptedError; полученные дельты остаются при вас.

Бюджет считает весь ответ

wait ограничивает ответ целиком, а не первый токен, и считает только то время, которое библиотека провела в ожидании провайдера:

stream = broker.stream("Напиши длинный ответ", wait=20.0)
try:
    async for delta in stream:
        print(delta, end="", flush=True)
except llmbroker.LLMTimeoutError as exc:
    print(f"\nсдались: {exc}")

Сколько времени вы тратите между дельтами — ваше дело: часы снимаются в момент выдачи дельты и продолжают идти с того же места, когда вы просите следующую, так что медленный читатель бюджет не тратит. А модель, которая мгновенно открыла стрим и потом цедит текст, в бюджет не укладывается — ровно ради этого он и считает ответ, а не его начало.

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

По той же причине до первой дельты handle ничего не говорит: llm_name и call_id до неё равны None, потому что вызов ещё может уехать на другую модель. usage появляется ещё позже — когда ответ закончился.

Оценка ждёт того же момента, что и счётчики: конца ответа, когда вызов попадает в журнал. Попросите раньше — получите ValueError, а не оценку, которая тихо никуда не денется.

Оборвать чтение можно в любой момент, но слот модели возвращает именно закрытие стрима: break сам по себе handle не закрывает — брошенный, он освободит слот только когда до него доберётся сборщик мусора. Закрывайте сами; заодно это единственный способ оценить то, что успело прийти, — закрытие завершает вызов:

stream = broker.stream("Напиши хокку про брокеров")
async for delta in stream:
    if looks_wrong(delta):
        break

await stream.aclose()
await stream.record_quality(0.0)

Либо пусть закрывает контекстный менеджер — там, где оценивать нечего:

async with contextlib.aclosing(broker.stream("...")) as stream:
    async for delta in stream:
        ...

Стримить одну названную модель, без пула и фейловера, умеет и direct — см. Прямые вызовы модели.

Один процесс, один файл, без шага инициализации

sqlite держит модели, ключи и журнал в одном файле — этого хватает однопроцессному сервису, — а курируемый список наполняет его перед первым вызовом:

async with llmbroker.AsyncBroker("broker.db") as broker:
    print((await broker.ask("Привет")).text)

База стартует пустой и наполняется до инициализации пула — отдельного шага инициализации помнить не нужно, а дальше список поддерживается свежим сам. Это best-effort: если каталог недоступен, в лог уйдёт предупреждение, а процесс стартует на том, что уже лежит в файле, — а когда файл пуст, на копии пресета, вшитой в пакет.

Если этот файл общий с вашим приложением, стоит знать про WAL и файловый лок — см. SQLite: общая база и WAL. Для нескольких процессов или хостов наполняйте базу один раз в задаче деплоя — см. Серверы и кластеры.