Асинхронность и стриминг
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. Для нескольких процессов или хостов наполняйте базу один раз в задаче деплоя — см. Серверы и кластеры.