Skip to content

Серверы и кластеры

Тот же брокер масштабируется на несколько процессов и хостов: укажите ему общую БД вместо собственного хранилища — вызывающий код не меняется.

Общая БД

Первый аргумент брокера задаёт сразу пул моделей, ключи и журнал:

llmbroker.Broker()                          # собственный каталог llmbroker + ключи из окружения
llmbroker.Broker("broker.db")               # sqlite
llmbroker.Broker("postgresql://host/db")    # postgres
llmbroker.Broker("mongodb://host/db")       # mongodb

Каждому варианту нужен свой extra — см. Установка. Любую часть можно переопределить явно через registry= / secrets= / store=.

Свой реестр вместо нашего

Чтобы задать пул самому, передайте объект, реализующий протокол реестра, и скажите, чему он следует:

broker = llmbroker.Broker(registry=MyRegistry(), sync=None)        # только ваши записи
broker = llmbroker.Broker(registry=MyRegistry(), sync="freetier")  # ваши плюс наши
что передали sync= не указан записи, которые вы внесли сами
ничего или URL базы следует "freetier" обновление их не трогает
объект реестра ошибка — скажите, чему обновление их не трогает

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

from llmbroker.postgres import Secrets, Store

broker = llmbroker.AsyncBroker(
    registry=MyRegistry(),
    secrets=Secrets(pool),            # ключи в вашей базе, а не в окружении
    store=Store(pool),                # журнал туда же
    sync=None,
)

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

Своя запись — в пуле

Свою запись кладут через протокол реестра, а не записью строк: раскладка таблиц принадлежит llmbroker и может меняться от релиза к релизу. Прочитайте, что там есть, добавьте своё и запишите всё обратно — mirror зеркалит целиком, так что всё не переданное удаляется:

from llmbroker import LLMConfig
from llmbroker.postgres import Registry

registry = Registry(pool)
mine = LLMConfig(
    name="my-gateway",
    base_url="https://gw.internal/v1",
    model="m",
    api_key_ref="MY_GATEWAY_KEY",
)
await registry.mirror([*await registry.load(), mine])

Ничто не помечает её как нашу, поэтому никакая синхронизация её не удалит и не перепишет. Это участник пула, и роутер делает на неё failover; endpoint, который вы хотите вызывать по имени, — это объявленная модель, и она нигде не хранится.

Наполнение БД: задача деплоя, а не старта

БД стартует пустой. Синхронизируйте её из своего кода, в том же шаге деплоя, где выполняется alembic upgrade, — фабрикой, которой уже пользуется приложение, чтобы DSN и его секреты жили ровно в одном месте:

broker = build_broker()                   # собственная фабрика приложения
try:
    print(await broker.sync("freetier"))  # курируемый пресет — единственный источник
finally:
    await broker.aclose()

Обратите внимание: это не async with. Вход в брокер инициализирует пул, а задача деплоя должна сделать ровно одно — наполнить базу. А там, где брокеру запрещено ходить в сеть (ниже), контекстный менеджер и вовсе поднимет EmptyRegistryError прямо на входе, до того как вы вызовете sync().

Запускайте это как разовую задачу (release phase, Kubernetes Job, init-контейнер). Поддерживать список свежим дальше не нужно никакой задачей: обслуживающие процессы сами перепроверяют курируемый список примерно раз в сутки, на вызове, который они и так делали. N узлов проверяют безопасно: все они считают одно и то же слияние из одного апстрима и одних ключей, поэтому первая запись всё решает, а проверка любого другого узла не находит работы. Схема избегает другого — узла, сверяющего реестр со своей локальной копией.

sync принимает имя курируемого пресета и больше ничего — ни пути к файлу, ни второго реестра. Строка подключения продолжает следовать курируемому пресету; а если брокеру передан объект реестра, придётся сказать, чему он следует: sync="freetier" или sync=None. В обоих случаях обновление переписывает только записи, сделанные самой синхронизацией, — записи, которые ваша установка задаёт через свой собственный реестр, оно не трогает.

Развёртывание, которому нельзя ходить в сеть во время обслуживания

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

llmbroker.AsyncBroker("postgresql://host/db", sync_interval=None)   # в вашей фабрике
broker = build_broker()
try:
    report = await broker.sync()      # без аргумента: то, чему следует эта установка
    if report is not None:            # один платный каталог ничего не сливает
        print(llmbroker.format_report(report))
finally:
    await broker.aclose()

sync_interval=None останавливает в процессе все часы, которые ходят в сеть, — и курируемый список моделей, и платный каталог, через который разрешаются алиасы из direct=. Вместе с ними останавливается и загрузка, наполняющая пустой реестр на старте: такой брокер поднимет EmptyRegistryError с указанием на эту задачу, а не пойдёт в сеть, чтобы обслужить первый запрос. Это и есть задуманное поведение: поставили переключатель и забыли про задачу — брокер обслуживать не будет.

sync() без аргумента синхронизирует то, чему следует установка: пресет, названный в sync=, а если она не следует ни одному — только платный каталог, который обновляет объявленные алиасы, ничего не сливает в реестр и потому не возвращает отчёта.

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

Свежесть списка теперь ваша забота. Провайдеры закрывают бесплатные endpoint'ы без предупреждения, поэтому список, который никто не обновляет, вырождается в пул, который не может обслуживать. Запускайте задачу рядом с миграциями на каждом деплое и по собственному расписанию между деплоями — встроенные часы проверяют примерно раз в сутки, и повторить это будет безопасным выбором.

Переезд установки между бэкендами

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

from llmbroker.mongodb import Registry as MongoRegistry
from llmbroker.postgres import Registry as PostgresRegistry

old = PostgresRegistry(old_pool)
new = MongoRegistry(new_db)
await new.mirror(await old.load())

Секреты и журнал переезжают так же через свои бэкенды, если это нужно; обычно ключи выдаются заново, а журнал остаётся на старом месте.

Платная модель, которую вы зовёте по имени, объявляется там же, где фабрика создаёт брокер — AsyncBroker(dsn, direct=["opus"]), — а не пишется в реестр. Одна строка в фабрике, которая у вас уже есть, покрывает весь кластер, и каждый процесс сам разрешает алиас заново по своим часам обновления, так что долгоживущий деплой не сидит вечно на том id модели, с которым его развернули. См. Прямые вызовы модели.

sync приводит записи, которые сделал он сам, в соответствие с курируемым списком, которому следует установка: запись, всё ещё присутствующая в списке, обновляется, запись, которой в списке больше нет, удаляется, новая добавляется. Ничто не взвешивает, могла бы выброшенная запись ещё работать здесь, а запись, которую ваша установка сама положила в реестр, не трогается никогда. Удаление ограничено там, где список курируется: запись покидает его лишь тогда, когда её уже нельзя вызвать. Возвращаемый SyncReport говорит, что именно произошло, на каждом запуске, включая запуск без изменений, и называет ключ, который стал не нужен. Ненулевой код выхода задачи и её лог — тот же канал для админа, которым уже пользуется упавшая миграция; кому нужно переслать это дальше, читает broker.last_sync_report.

Наблюдение за пулом с админского экрана

snapshot() — один вызов, который наполняет весь экран: строки по моделям плюс вердикт по пулу целиком, а тот же вердикт llmbroker пишет в лог, так что алерт обходится без опроса. См. Наблюдение и журнал.

SQLite: общая база и WAL

Штатный режим — одна база, общая для llmbroker и вашего приложения: брокер держит свои таблицы llmbroker_* рядом с вашими и больше ничего не трогает (хук Alembic убирает их из автогенерации миграций).

В том числе PRAGMA user_version — слот в заголовке файла, который используют многие миграционные инструменты, остаётся вашим. Свою версию схемы брокер хранит в таблице llmbroker_schema_version, поэтому удаление таблиц llmbroker_* сбрасывает всё, что llmbroker держит в файле.

llmbroker никогда не устанавливает и не меняет journal_mode SQLite — WAL это персистентное свойство уровня файла, принадлежащее владельцу файла БД, поэтому включать его — ваша задача, не брокера. На общем файле владелец — ваше приложение: включите WAL там, если нужна конкурентность чтения и записи.

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

sqlite3 broker.db 'PRAGMA journal_mode=WAL'

Это касается только SQLite. У Postgres и MongoDB такого файлового лока нет — общая база с приложением нормальна, а отдельная схема или база — опциональная аккуратность, не потребность конкурентности.

Ошибки запуска

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

  • SyncRefusedErrorsync() отказался применять результат, после которого в рабочем реестре не осталось бы ни одной записи. Ничего не записано; в report лежит то, что слияние собиралось сделать.
  • EmptyRegistryError — в реестр ещё ничего не синхронизировано. Безобидно: установка не настроена, а не сломана.
  • SchemaVersionError — в хранилище версия схемы, с которой этот релиз работать не может. Фатально и требует действий оператора: удалить таблицы llmbroker_* и перезапустить (сначала выгрузите реестр/секреты/журнал, если они нужны). Обе версии лежат в found и expected.

Все три лежат на llmbroker (llmbroker.SchemaVersionError и так далее) и наследуют LLMBrokerError, который сам является RuntimeError, — ловите на нужной вам гранулярности:

try:
    models = broker.snapshot()
except llmbroker.EmptyRegistryError:
    models = {}   # ещё ничего не настроено — пустой экран, а не 500

SchemaVersionError пробрасывается: его сообщение — инструкция оператору, поэтому проглотив его вы превратите несовпадение схемы в «провайдеры не настроены». LLMBrokerError ловит все три случая, RuntimeError — их и всё остальное.

Ошибка самого запроса приходит из отдельного дерева (LLMRequestError и его наследники) — см. Когда ответить некому.

Закрытие брокера

Закрывайте брокер явно, если долгоживущий процесс создаёт брокеры повторно или подключена внешняя БД:

with llmbroker.Broker("broker.db") as broker:
    reply = broker.ask("...")

AsyncBrokerasync with или await broker.aclose().

Журнал вызовов

Каждая попытка вызова оставляет запись: кто отвечал, чем кончилось, во что обошлось, с каким trace_id её звали и как её потом оценили. Читают его broker.calls(...) и broker.stats(...), и они не инициализируют пул, поэтому работают и на несинхронизированной установке. См. Наблюдение и журнал.

Журнал самоочищается: записи старше retention удаляются, по умолчанию это 90 дней. Глубина хранения — свойство бэкенда журнала, а не брокера, поэтому строкой подключения её не задать: соберите порты сами и задайте её тому, который пишет журнал.

from datetime import timedelta

from llmbroker.postgres import Registry, Secrets, Store

broker = llmbroker.AsyncBroker(
    registry=Registry(pool),
    secrets=Secrets(pool),
    store=Store(pool, retention=timedelta(days=365)),
    sync="freetier",                  # реестр передан объектом — скажите, чему он следует
)

Один брокер, вызывающий на запрос

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

Четыре развёртывания, по нарастанию потребностей:

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

broker = llmbroker.Broker()
print(broker.ask("hi").text)

Долгоживущий процесс с базой. Создавайте брокер там же, где создаёте движок базы — один раз, на старте, — и закрывайте на остановке:

@asynccontextmanager
async def lifespan(app: FastAPI):
    app.state.broker = llmbroker.AsyncBroker("postgresql://host/db")
    try:
        yield
    finally:
        await app.state.broker.aclose()

Кластер на общих ключах. Каждый обработчик берёт собственного вызывающего брокера. Один пул, один набор ключей, один пул соединений на процесс:

def llms(request: Request) -> llmbroker.AsyncLLMs:
    return request.app.state.broker.llms

@app.post("/ask")
async def ask(prompt: str, llms: llmbroker.AsyncLLMs = Depends(llms)):
    return (await llms.ask(prompt)).text

Тот же кластер, но с ключом на пользователя. Меняется только зависимость:

def llms(request: Request) -> llmbroker.AsyncLLMs:
    return request.app.state.broker.for_scope(request.headers["x-user-id"])

Ключ пользователя лежит под ref'ом с его скоупом впереди: <скоуп>/<REF>. Вызывающий со скоупом u-42 сначала спросит у хранилища секретов u-42/GROQ_API_KEY и только потом упадёт на общий GROQ_API_KEY инсталляции; общее значение читается один раз на всех. Значит, чтобы дать пользователю свой ключ, вы кладёте его под этим именем — в переменную окружения, в свою БД, в AWS или Vault, туда же, где лежат общие. Скоуп — это просто строка, которую вы передали в for_scope(...); никакого понятия пользователя внутри llmbroker нет. У Vault одна оговорка про / в имени — см. API-ключи.

Каждая строка журнала, которую пишет такой вызывающий, несёт его скоуп, так что история одного пользователя — это broker.for_scope(user).calls(...). Отдельного параметра scope= у calls() нет: скоуп берётся у вызывающего, через которого вы читаете. Собственные broker.calls() и broker.stats() — это взгляд всей инсталляции, они видят строки всех скоупов сразу.

Пул, накопленное качество и ограничение parallel на модель принадлежат брокеру, а не вызывающему: отдельный счётчик на пользователя — это уже не ограничение. Ключ, который провайдер отверг у одного вызывающего, перестаёт предлагаться именно ему; вызывающий со своим ключом не затронут, а те, кто платил одним и тем же значением, теряют его вместе — потому что это один и тот же ключ.

Что процессы делят, а что нет

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

Правка реестра, сделанная соседом, доедет на следующей перестройке пула. Пул перестраивается на старте, по часам обновления (примерно раз в сутки), на явном sync() и тогда, когда пул только что не смог ответить. Больше ничто порты не перечитывает, поэтому успешный вызов не стоит базе ничего, кроме собственной строки журнала.

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

Alembic

Чтобы автогенерация миграций игнорировала таблицы llmbroker_*:

# alembic/env.py
import llmbroker.integrations.alembic

context.configure(
    connection=connection,
    target_metadata=target_metadata,
    include_object=llmbroker.integrations.alembic.include_object,
)

Свой include_object скомпонуйте с ним через and.