Skip to main content

GigaChat client with certificate pool, shared limits and stop events

Project description

gigamux

Асинхронный клиент GigaChat: пул сертификатов-каналов, общие in-flight лимиты через Redis (coredis) с fallback на локальные счётчики и drop-in адаптеры LangChain. Сервисы переходят на библиотеку заменой импорта, retry-логика и троттлинг живут в ядре.

Quick start

Ядро напрямую:

from gigamux import GigaCoreClient, config_from_env, ChatRequest, Message

config = config_from_env()
async with GigaCoreClient(config) as client:
    result = await client.chat(
        ChatRequest(
            model="GigaChat-2-Max",
            messages=(Message(role="user", content="привет"),),
        )
    )
    print(result.content, result.meta.channel)

LangChain-сервисы (адаптеры импортируются напрямую, требуют extra langchain):

from gigamux.adapters.langchain_chat import PooledGigaChat

llm = PooledGigaChat(core=client, model="GigaChat-2-Max", temperature=0.1)
message = await llm.ainvoke("привет")

Индексация (один батч = один HTTP-запрос = один слот, нарезка по 20 в адаптере):

from gigamux.adapters.embeddings import PooledGigaEmbeddings

embedder = PooledGigaEmbeddings(core=client, model="EmbeddingsGigaR")
vectors = await embedder.aembed_documents(texts)

Подсчёт токенов без генерации (POST tokens/count):

counts = await client.count_tokens("GigaChat-2-Max", "сколько здесь токенов?")
print(counts[0].tokens, counts[0].characters)

count_tokens принимает строку или список строк и всегда возвращает list[TokenCount].

Проверка контента на провокационность (POST filter/check, модель Gigafilter):

from gigamux import Message

verdict = await client.check_filter(
    "Gigafilter",
    (Message(role="user", content="проверяемый текст"),),
    settings={"neuro": True, "blacklist": True, "whitelist": True},
)
print(verdict.is_profane, verdict.filter_tokens)

is_profane — единый вердикт на весь батч сообщений (не по каждому сообщению). settings опционален. Эндпоинт только v1: на v2-канале запрос автоматически уходит на соседний …/v1.

Function calling и structured output

PooledGigaChat поддерживает инструменты и структурированный вывод — паритет с langchain-gigachat, включая react-агентов LangGraph:

from langchain_core.tools import tool
from langgraph.prebuilt import create_react_agent

@tool
def get_weather(city: str) -> str:
    """Return current weather for a city."""
    ...

agent = create_react_agent(llm, [get_weather])
state = await agent.ainvoke({"messages": [...]})

Extra gigamux[langchain] ставит только langchain-core; для примера выше нужен отдельно установленный langgraph (pip install langgraph).

Структурированный вывод доступен двумя методами эталонной библиотеки: function_calling (по умолчанию) и json_mode (через response_format).

chain = llm.with_structured_output(MyModel)
chain = llm.with_structured_output(MyModel, method="json_mode")
result = await chain.ainvoke("...")

Ограничения и семантика:

  • Стриминг с инструментами на уровне API невозможен, поэтому astream/astream_events с привязанными функциями прозрачно выполняют обычный вызов и отдают ответ одним чанком: токенного стриминга нет, но react-агенты под astream_events/stream_mode="messages" работают.
  • Параллельные tool_calls невозможны на стороне GigaChat: больше одного tool_call в сообщении — ValueError.
  • В with_structured_output поддерживаются pydantic-модели обоих поколений (v2 и pydantic.v1) в обоих методах; v1-класс возвращается экземпляром, как и v2.
  • finish_reason доступен в response_metadata каждого ответа адаптера.
  • with_structured_output методом function_calling возвращает None, если модель ответила текстом вместо вызова функции (поведение эталона); с include_raw=True сырое сообщение доступно в output["raw"].
  • json_mode использует бета-фичу API response_format — при её отсутствии на контуре ошибка идентична поведению langchain-gigachat.
  • few_shot_examples пробрасываются из metadata инструмента.

Стриминг

Ядро отдаёт дельты как StreamChunk; слот канала занят весь стрим, метаданные приходят в последнем чанке:

request = ChatRequest(
    model="GigaChat-2-Max",
    messages=(Message(role="user", content="расскажи анекдот"),),
)
async for chunk in client.stream(request):
    print(chunk.delta, end="")
    if chunk.meta is not None:
        print(chunk.meta.channel, chunk.meta.duration_ms, chunk.usage)

Если API прислал блок usage в стриме, он доступен на чанке, который его принёс, и продублирован на финальном чанке (chunk.usage); у адаптера финальный AIMessageChunk несёт usage_metadata.

Ошибки до первого чанка ретраятся как обычные запросы; после начала стрима — StreamInterruptedError без молчаливого ретрая. У адаптера работает llm.astream(...); без инструментов это настоящий токенный стриминг, а с привязанными функциями ответ приходит одним финальным чанком (см. раздел про function calling).

Потребитель, который может бросить стрим до конца, должен оборачивать его в contextlib.aclosing(...), иначе слот освобождается только сборщиком мусора. Keepalive слота ограничен slot_max_lifetime (дефолт 900 с): после этого лимита lease перестаёт продлеваться и истекает по TTL даже без финализации генератора.

Конфигурация в коде

config_from_env() — удобный путь, но конфиг можно собрать руками:

from gigamux import ChannelConfig, ClientConfig, GigaCoreClient

config = ClientConfig(
    channels=[
        ChannelConfig(
            name="cert-a",
            base_url="https://giga.internal:10501/v1",
            cert_file="/etc/certs/a/tls.pem",
            key_file="/etc/certs/a/key.pem",
            ca_bundle="/etc/certs/ca.pem",
            limits={"GigaChat-2-Max": 2, "EmbeddingsGigaR": 4},
        ),
        ChannelConfig(
            name="cert-b",
            base_url="https://giga.internal:10502/v1",
            limits={"*": 3},
        ),
    ],
    redis_url="rediss://redis.internal:6380/0",
    redis_username="svc",
    redis_password="...",
    acquire_timeout=300.0,
    max_retries=None,
)

async with GigaCoreClient(config) as client:
    ...

Ключевые поля ClientConfig (все имеют дефолты):

Поле Дефолт Назначение
acquire_timeout 300.0 дедлайн на всю операцию: ожидание слота + ретраи 5xx (429 — reroute без ожидания); None — без дедлайна (индексация)
lease_ttl / heartbeat_interval 120 / ttl/3 время жизни слота в Redis и период продления
slot_max_lifetime 900.0 максимум продления слота keepalive-ом; дальше lease истекает по TTL
max_retries / backoff_factor / backoff_max None / 1.0 / 30.0 потолок ретраев 5xx (None — без счётчика, ограничено только acquire_timeout) + экспоненциальная пауза с джиттером
connect_timeout / read_timeout 10 / 120 таймауты HTTP на один аттемпт (не на всю операцию)
redis_failure_threshold / redis_probe_interval 3 / 15.0 circuit breaker: после N ошибок Redis лимиты локальные, проба каждые 15 с

limits канала — единственный источник правды о том, какие модели канал обслуживает: ключ — имя модели, значение — число одновременных запросов, "*" — wildcard. cert_file/key_file задаются только вместе; без них канал работает без mTLS. acquire_timeout можно переопределить и на один вызов: client.chat(request, acquire_timeout=None).

Личность канала: физический сертификат или cn. Локально канал несёт mTLS-пару (cert_file/key_file). За шлюзом, который сам хранит сертификаты и переключает их по CN, вместо путей задаётся cn — тогда каждый запрос канала уходит с заголовком ai-custom-client-cert: <cn>. Каналы, смотрящие в один base_url, обязаны каждый нести личность (пару или cn) — иначе шлюз их не различит и конфиг падает на старте. Одиночный канал без личности валиден (шлюз возьмёт первый сертификат неймспейса).

ChannelConfig(
    name="cert-a",
    base_url="https://giga.internal:10500/v1",
    cn="my-namespace-cert-a",
    limits={"GigaChat-2-Max": 2},
)

Модель таймаутов и ретраев. Три независимых уровня: acquire_timeout — дедлайн на всю операцию (ожидание слота + все ретраи), read_timeout — таймаут одного HTTP-аттемпта в гигу, max_retries — опциональный потолок числа ретраев. По умолчанию max_retries=None, поэтому ретраи 5xx ограничены только временем: запрос терпит серверный шторм пока есть бюджет acquire_timeout, и лишь затем падает с RequestTimeoutError (с логом giving up ... reason=deadline — тихих смертей нет). Жёсткий потолок попыток включается заданием max_retries=N: тогда сдача по тому, что наступит раньше — счётчик или дедлайн. 429 не ретраится: канал под rate-limit сразу перебирается на другой (reroute без backoff), а когда 429 отдали все каналы — немедленный RateLimitError. acquire_timeout=None вместе с max_retries=None означает «ретраить 5xx до победы» (бесконечно).

Версии API (v1 и v2)

Библиотека говорит и на v1, и на v2 GigaChat одним и тем же публичным интерфейсом — версия выбирается по каналу. Версия определяется по base_url (суффикс …/v2 → v2) либо явным полем api_version; дефолт — v1.

ChannelConfig(
    name="v2",
    base_url="https://giga.internal:10501/v2",  # api_version выводится как v2
    limits={"GigaChat-2-Max": 4},
)
ChannelConfig(
    name="v2-explicit",
    base_url="https://giga.internal:10501/gateway",
    api_version="v2",
    limits={"GigaChat-2-Max": 4},
)

Один канал = одна версия; разные версии — разные каналы пула. Код вызова (chat/stream/embed/count_tokens) и адаптеры LangChain не меняются: кодек версии собирает тело запроса и разбирает ответ. v2 добавляет tools, content-как-массив, конверт messages[] и событийный SSE (response.message.done); состояние функций (functions_state_idtools_state_id) прокидывается прозрачно. Эмбеддинги и tokens/count существуют только в v1 и на v2-канале автоматически маршрутизируются на соседний …/v1-путь.

Переменные окружения

Переменная Назначение
GIGA_CHANNELS JSON-список каналов (name, base_url, cert_file, key_file, ca_bundle, cn, api_version, limits)
GIGA_EVENTS_INCLUDE_TEXT true/false: класть ли текст промпта/ответа в GigaEvent (дефолт false, PII)
GIGA_REDIS_URL URL Redis для распределённых лимитов; не задан -> локальные лимиты
GIGA_REDIS_USERNAME / GIGA_REDIS_PASSWORD аутентификация Redis
GIGA_REDIS_KEY_FILE / GIGA_REDIS_CERT_FILE / GIGA_REDIS_CA_BUNDLE mTLS для Redis (нужны все три, иначе соединение без TLS с предупреждением)
GIGA_ACQUIRE_TIMEOUT дедлайн на всю операцию (слот + ретраи), дефолт 300; none/null — без дедлайна
GIGA_LEASE_TTL / GIGA_HEARTBEAT_INTERVAL / GIGA_SLOT_MAX_LIFETIME тайминги lease
GIGA_MAX_RETRIES / GIGA_BACKOFF_FACTOR / GIGA_BACKOFF_MAX политика ретраев; GIGA_MAX_RETRIES=none — без счётчика (только по времени)
GIGA_CONNECT_TIMEOUT / GIGA_READ_TIMEOUT / GIGA_WRITE_TIMEOUT / GIGA_POOL_TIMEOUT таймауты HTTP
GIGA_REDIS_CONNECT_TIMEOUT / GIGA_REDIS_STREAM_TIMEOUT таймауты Redis
GIGA_REDIS_FAILURE_THRESHOLD / GIGA_REDIS_PROBE_INTERVAL / GIGA_REPLICA_COUNT circuit breaker и локальный fallback
GIGA_STOPEVENT_MODE stop (дефолт) — при StopEvent летит StopEventError; fallback — переход на резервную модель
GIGA_FALLBACK_MODELS JSON-список резервных моделей для fallback-режима, напр. ["GigaChat-Pro","GigaChat-2"]
GIGA_PREVIEW_ENABLED true/1/yes/on включает preview-маршрутизацию (дефолт выключено)
GIGA_PREVIEW_MODELS JSON-карта основная→preview-модель, напр. {"GigaChat":"GigaChat-preview"}
GIGA_PREVIEW_FRACTION доля трафика на preview, 0..0.05 (потолок 5%), дефолт 0.05
GIGA_PREVIEW_TIMEOUT бюджет одной preview-попытки в секундах (дефолт 5); залипший preview уступает основной модели, не жгя весь acquire_timeout
GIGACHAT_HOST / GIGACHAT_PORT / GIGACHAT_TLS_CERT_FILEPATH / GIGACHAT_KEY_FILEPATH / GIGACHAT_CA_BUNDLE_FILEPATH режим совместимости: один канал из legacy-переменных
GIGACHAT_ENDPOINT legacy-путь base_url в режиме совместимости (по умолчанию /v1)
GIGACHAT_MAX_CONCURRENCY лимит на канал в режиме совместимости (wildcard-модель)

Если задан GIGA_CHANNELS, он имеет приоритет. Иначе из GIGACHAT_* строится пул из одного канала с wildcard-лимитом, что даёт миграцию без изменения конфига.

Начиная с 0.2.1 конфигурация строгая: неизвестное поле в GIGA_CHANNELS или ClientConfig (например, опечатка в имени) — это ConfigError на старте, а не молчаливое игнорирование.

Семантика

Лимит — in-flight слоты на пару канал+модель, общие между репликами. Слот берётся как lease с TTL в Redis и продлевается heartbeat-ом, поэтому упавший под не держит слот навсегда. При занятых слотах вызов ждёт в очереди с джиттером до acquire_timeout (по умолчанию 300 с, индексация — без таймаута), затем SlotWaitTimeoutError. Этот же дедлайн ограничивает ретраи 5xx: когда он истекает под длительным серверным штормом, вызов падает с RequestTimeoutError, а не молча. 5xx на канале прозрачно ретраятся, а 429 перебирается на другой канал (reroute через prefer_not) без backoff; когда 429 отдали все каналы, сразу летит RateLimitError. Ошибка посреди стрима не ретраится молча — StreamInterruptedError.

GigaChat принимает только одно system-сообщение и только первым. Поэтому при сериализации запроса все system-сообщения склеиваются в одно ведущее (содержимое через "\n\n" в порядке появления, пустые/пробельные отбрасываются); относительный порядок остальных сообщений сохраняется. Запрос с единственным ведущим system не меняется. Когда схлопывается больше одного непустого system, пишется одна debug-строка лога — иначе мутация молчит. Это делает библиотеку чуть «прощающей» относительно langchain_gigachat — осознанное улучшение, а не баг.

Ошибки

Все исключения наследуют GigaClientError и импортируются из gigamux:

Исключение Когда летит
ConfigError невалидная конфигурация или переменные окружения
NoChannelForModelError ни один канал не обслуживает запрошенную модель
SlotWaitTimeoutError свободный слот не появился за acquire_timeout
RequestTimeoutError дедлайн acquire_timeout истёк на ретраях 5xx; несёт attempts, last_status, channel
RateLimitError 429 вернули все каналы (reroute исчерпан; при одном канале — сразу)
ServerError 5xx или сетевая ошибка пережили все ретраи
ApiError прочие неретраябельные ответы API (4xx); базовый класс двух предыдущих
HistoryError результат функции в истории не следует за assistant-вызовом с тем же именем
StreamInterruptedError стрим оборвался после уже отданных чанков

У ApiError и наследников доступны status_code, body, channel, request_id; в body кладётся message из тела {status, message}, когда оно есть. Типовая обработка: SlotWaitTimeoutError — перегрузка, имеет смысл отдать 429/503 наверх; ApiError — ошибка запроса, ретраить бесполезно.

Метаданные ответа

Каждый результат (ChatResult, EmbedResult, финальный StreamChunk) несёт meta: ResponseMetarequest_id, headers, channel, model, attempts, duration_ms, status_code. Это сквозной способ узнать, через какой канал ушёл запрос и сколько он занял, например для логов и метрик. На каждый ответ пишется строка лога с request_id (из заголовка ответа), каналом, моделью, статусом и длительностью — для корреляции с логами goprodigy/core; отключается через ClientConfig.log_responses.

Логирование

gigamux пишет через глобальный логгер loguru (from loguru import logger) — без собственных хендлеров и без logger.disable. Из коробки записи видны в stderr; заглушить библиотеку целиком можно через logger.disable("gigamux").

Если сервис строит свою конфигурацию loguru с патчером (проставляет extra-поля, сериализует записи, фильтрует хендлеры по extra), патчер обязан быть глобальным:

from loguru import logger

logger.configure(patcher=my_patcher)   # так — видит записи всех библиотек
patched = logger.patch(my_patcher)     # так НЕЛЬЗЯ полагаться для чужих записей

logger.patch() возвращает новый инстанс, и патчер применяется только к записям, отправленным через него. Записи gigamux (и любой библиотеки на голом глобальном логгере) через инстанс-патчер не проходят: они не получают проставляемых им extra-полей, и дальше два типовых исхода —

  • строгий фильтр хендлера вида record["extra"].get("target") == "log" молча отбрасывает их (логи библиотеки «пропадают»);
  • sink с форматом format="{extra[serialized]}" падает на каждой записи с KeyError (--- Logging error in Loguru Handler ---).

Глобальный патчер видит записи всех библиотек, поэтому он не должен предполагать форму чужих extra (проверяйте типы через isinstance, не обращайтесь к ключам без .get): упавший патчер роняет лог-вызов в точке чужого кода. Хендлер-фильтры при этом можно оставлять строгими — патчер отрабатывает раньше фильтров и сам проставляет недостающие поля.

Проброс заголовков

GigaCoreClient(header_provider=...) принимает колбэк без аргументов, чьи заголовки добавляются в каждый исходящий запрос (unary и stream). Это точка для сквозного трейс-идентификатора: агент кладёт X-Trace-Id/traceparent в contextvar, провайдер читает его на каждую попытку — значение переживает reroute/retry. Ключи/значения приводятся к str; любой сбой провайдера не роняет запрос (заголовки пропускаются с логом ошибки). Генерация и валидация идентификатора — ответственность агента.

from contextvars import ContextVar

trace_id: ContextVar[str] = ContextVar("trace_id")
client = GigaCoreClient(config, header_provider=lambda: {"X-Trace-Id": trace_id.get()})

Fallback и preview-модели

ModelRouter превращает запрошенную модель в упорядоченную цепочку кандидатов, по которой клиент идёт при переключениях. Настраивается через конфиг/env; по умолчанию (режим stop, preview выключен) цепочка равна одной запрошенной модели и поведение не отличается от прежнего.

Fallback на StopEvent. Когда модель отвечает временной недоступностью (HTTP 423 или 403 с сигнальным сообщением → StopEventError), в режиме fallback запрос уходит на следующую модель из плоского списка. Переключение происходит только на StopEvent — 429/5xx/фатальные ошибки обрабатываются в пределах модели и не маскируются. Общий acquire_timeout-дедлайн делится на всех кандидатов.

GIGA_STOPEVENT_MODE=fallback
GIGA_FALLBACK_MODELS=["GigaChat-Pro","GigaChat-2"]

Preview-маршрутизация. Мастер-тумблер (дефолт выключен) уводит долю трафика (≤5%, потолок enforced) на явно указанную preview-модель; на любой ошибке preview запрос откатывается на основную модель. В стриме переход возможен только до первого чанка — после отдачи данных ошибка пробрасывается без рестарта.

GIGA_PREVIEW_ENABLED=true
GIGA_PREVIEW_MODELS={"GigaChat":"GigaChat-preview"}
GIGA_PREVIEW_FRACTION=0.05

Решение «делать ли preview» и мониторинг доли — ответственность агента; библиотека даёт механизм и enforcement потолка. Резервные/preview-модели должны обслуживаться каналами, иначе переключение упрётся в NoChannelForModelError.

Метрики насыщения (saturation)

GigaCoreClient(on_slot_event=...) принимает колбэк, вызываемый на события слота конкурентности: SlotEvent(kind, channel, model, wait_seconds), где kind"acquired" (слот получен, wait_seconds — время ожидания), "released" (освобождён) или "wait_timeout" (не дождались слота за acquire_timeout; channel=""). Это единственный «золотой сигнал», которого нет в трейсинге исходящих вызовов. Пары acquired/released сбалансированы, поэтому занятость считается инкрементами без обращений к лимит-стору. Колбэк best-effort: его ошибка логируется и не роняет запрос.

from prometheus_client import Counter, Gauge, Histogram

inflight = Gauge("giga_inflight", "occupied slots", ["channel", "model"])
wait = Histogram("giga_slot_wait_seconds", "slot wait", ["model"])
rejected = Counter("giga_slot_wait_timeout_total", "acquire timeouts", ["model"])

def on_slot(event):
    if event.kind == "acquired":
        inflight.labels(event.channel, event.model).inc()
        wait.labels(event.model).observe(event.wait_seconds)
    elif event.kind == "released":
        inflight.labels(event.channel, event.model).dec()
    elif event.kind == "wait_timeout":
        rejected.labels(event.model).inc()

client = GigaCoreClient(config, on_slot_event=on_slot)

Latency/Errors/Traffic снимаются из on_response(meta) (duration_ms, статус, канал, модель, attempts) — см. «Метаданные ответа».

События наблюдаемости (GigaEvent)

GigaCoreClient(on_event=...) — один структурный GigaEvent на терминальный исход каждого вызова (chat / stream / embeddings / count_tokens / filter): успех или доменная ошибка. В отличие от on_response, событие видит и то, где HTTP-ответа не было: транспортные сбои (status_code=0), slot-таймауты, decode-ошибки. Внутри — токены из Usage, request_kind, error_class, attempts, duration_ms, ttft_ms (стрим), queue_wait_ms; проекции to_pg_row() (метрики без тяжёлого текста) и to_os_doc() (Q&A + фасеты) — стык двух стоков по request_id.

Что приходит в колбэк

on_event получает один frozen-датакласс GigaEvent. Поля (все, кроме Q&A-блока, заполнены всегда):

Группа Поля
время / трассировка ts (ISO-8601 UTC), request_id (X-Request-ID; None при транспортном сбое), session_id, trace_id
маршрутизация model, channel (None, если слот так и не взяли), request_kind (chat/stream/embeddings/count_tokens/filter), stream
исход status_code (0 = HTTP-ответа не было), is_error, error_class (stop_event/rate_limit/server/timeout/…), error_message, finish_reason, attempts
тайминги duration_ms (весь вызов с ретраями), ttft_ms (стрим), queue_wait_ms (ожидание слотов)
токены prompt_tokens, completion_tokens, total_tokens, cached_tokens, precached_prompt_tokens
инструменты tool_names (объявленные, в порядке объявления), tool_choice (auto/none/имя инструмента), called_tool_name (что вызвала модель) — не PII, заполняются всегда
Q&A (PII, off by default) prompt_text, prompt_messages, response_text, function_call (имя + аргументы вызова), tools_spec (полные схемы инструментов), n_messages
сервис app, env, lib_version

Готовых проекции две — они и есть подсказка, как раскладывать событие на два стока:

  • to_pg_row() → узкая строка метрик: тяжёлый Q&A (prompt_messages, response_text, function_call, tools_spec) выброшен, prompt_text обрезан до 256 символов как превью. tool_names/tool_choice/called_tool_name остаются — по ним считается, сколько запросов шло с инструментами и чем кончилось. Ключи строки = имена колонок.
  • to_os_doc() → документ поиска: @timestamp + request_id + фасеты (model/request_kind/is_error/session_id) + полный Q&A-текст + полные схемы инструментов (tools_spec) и tool-ходы истории.

prompt_messages — не только role/content: у сообщений истории сохраняются name, function_call, functions_state_id, tools_state_id, когда они заданы. Иначе из истории пропадали обе половины tool-раунда: и вызов ассистента, и результат функции. Оба диалекта читаются одинаково — v1 functions/function_call и v2 tools/tool_choice сводятся в одни и те же tool_names/tool_choice.

Три тонкости, о которых лучше знать заранее:

  • Конфликт tool_choice и function_call разрешается ровно так же, как на проводе. tool_choice побеждает, когда это словарь с именем/режимом или строка none/auto; любая другая строка кодеком не распознаётся, и в дело идёт function_call — событие называет то же, что реально ушло.
  • tool_choice не пропускает произвольные строки. Поле заполняется всегда, мимо PII-выключателя, поэтому значение, которое не является ни известным режимом (auto/none/forced/one_function), ни именем объявленного инструмента, схлопывается в "other".
  • Фасеты описывают запрос, а не закодированное тело. v1-канал сериализует только functions (tools молча теряются), но в событии они всё равно будут — это рассогласование стоит видеть, а не прятать. Пустой (не None) tool_names значит «инструменты объявлены, имён извлечь не удалось».

Две проекции соединяются по request_id: числа живут в метрик-сторе, текст — в индексе, JOIN по общему ключу.

Механизм доставки

Колбэк синхронный и вызывается в горячем пути — он не должен блокироваться на I/O. Механизм доставки в комплекте: EventPump (bounded-очередь drop-oldest + фоновый воркер с батчингом; при переполнении теряется старейшее событие, счётчик pump.dropped) и EventSink (Protocol с async def write(events)). Сам write крутится в фоне воркера, не в пути запроса, поэтому в нём можно делать настоящий батч-I/O; исключение из write пампу безопасно — он логирует и роняет батч, а не падает в обработку запроса.

Слив в Postgres (метрики)

from gigamux import EventPump, GigaCoreClient

PG_COLS = (  # порядок = порядок плейсхолдеров ниже
    "ts", "request_id", "session_id", "trace_id", "model", "channel",
    "request_kind", "stream", "status_code", "is_error", "error_class",
    "error_message", "finish_reason", "attempts", "duration_ms", "ttft_ms",
    "queue_wait_ms", "prompt_tokens", "completion_tokens", "total_tokens",
    "cached_tokens", "precached_prompt_tokens", "prompt_text", "n_messages",
    "tool_names", "tool_choice", "called_tool_name",
    "app", "env", "lib_version",
)
INSERT_SQL = (
    f"INSERT INTO giga_events ({', '.join(PG_COLS)}) "
    f"VALUES ({', '.join(f'${i}' for i in range(1, len(PG_COLS) + 1))}) "
    "ON CONFLICT (request_id) DO NOTHING"  # идемпотентность, см. ниже
)

class PgSink:
    def __init__(self, pool):
        self._pool = pool

    async def write(self, events):
        rows = [tuple(e.to_pg_row()[c] for c in PG_COLS) for e in events]
        async with self._pool.acquire() as conn:
            await conn.executemany(INSERT_SQL, rows)

async with EventPump(PgSink(pool)) as pump:
    client = GigaCoreClient(config, on_event=pump.push)
    ...
  • Тип колонок: ts/request_id/строки → text, *_tokens/attempts/status_codebigint/int, *_msdouble precision, is_error/streamboolean. request_idPRIMARY KEY (или уникальный индекс) для ON CONFLICT. tool_namestext[] (в событии это кортеж строк), tool_choice и called_tool_nametext.
  • Не заводите колонку под каждое поле руками, если не нужно: можно один jsonb payload = to_pg_row(e) целиком, а «горячие» фильтруемые поля (ts, model, is_error, request_kind) вынести отдельными колонками для индексов.

Слив в OpenSearch (Q&A + поиск)

class OsSink:
    def __init__(self, os_client, index="giga-events"):
        self._os, self._index = os_client, index

    async def write(self, events):
        actions = []
        for e in events:
            doc = e.to_os_doc()
            _id = e.request_id or f"{e.ts}:{e.model}"   # см. про None ниже
            actions.append({"index": {"_index": self._index, "_id": _id}})
            actions.append(doc)
        await self._os.bulk(body=actions)   # один bulk на батч

async with EventPump(OsSink(os_client)) as pump:
    client = GigaCoreClient(config, on_event=pump.push)
    ...
  • _id = request_id — ключевой приём. to_os_doc() намеренно не кладёт _id; задавайте его сами в bulk-заголовке, тогда повтор события просто перезапишет документ, а не создаст дубль (см. идемпотентность ниже).
  • Маппинг индекса: @timestampdate, request_id/session_id/model/ request_kind/tool_names/tool_choice/called_tool_namekeyword (для фасетов/агрегаций), prompt_text/response_texttext, prompt_messages/function_call/tools_specobject/flattened. Не отдавайте их на dynamic mapping — keyword vs text под фасеты важен, а схемы инструментов лучше flattened, чтобы произвольные parameters не взрывали маппинг.
  • Один to_os_doc() — плоская пара Q&A, поэтому это классический pattern «write-once, search-many»; ротацию делайте по @timestamp (data stream / ILM).

Идемпотентность (важно)

Доставка при отмене/закрытии — at-least-once: если cancel прервал write, памп пере-отправит весь батч целиком на flush-on-cancel. Поэтому оба стока должны быть идемпотентны по request_id (ON CONFLICT DO NOTHING / bulk _id). Тонкость: у транспортных сбоев request_id is None — для них дедуп по ключу не работает; давайте им синтетический ключ (ts+model, как выше) либо миритесь с редким дублем именно на этих событиях.

Сквозной trace_id

trace_id — ваш сквозной идентификатор запроса (например, порождённый фронтом). Передаётся keyword-аргументом в любой вызов и возвращается тождественным:

res = await client.chat(request, trace_id="T-123")
res.meta.trace_id            # "T-123"

try:
    await client.chat(request, trace_id="T-123")
except GigaClientError as e:
    e.trace_id               # "T-123" — id доступен и на ошибке

Принимают trace_id все методы, но возвращается он по-разному — потому что не у всех результатов есть meta:

Метод На результате На исключении В логе и GigaEvent
chat, embed res.meta.trace_id да да
stream meta финального чанка да да
count_tokens, check_filter — (TokenCount/FilterResult без meta) да да
probe — (ProbeResult без meta) — (probe не бросает) да

Для трёх нижних значение тождественно переданному, поэтому вызывающий его и так держит — отдельное echo-поле в эти типы не вводилось намеренно. Учти, что probe гасит ошибки и возвращает ProbeResult(ok=False, error=…), так что при его провале trace_id виден только в логе и в событии.

⚠️ В стриминге meta несёт только финальный чанк (у остальных meta is None). Если потребитель прерывает итерацию раньше — например, по finish_reason == "stop" из чанка кодека — он до meta.trace_id не доберётся. Дочитывайте генератор до конца либо держите trace_id на своей стороне.

В GigaChat не отправляется. Дока провайдера не содержит слота под сквозной trace: есть только X-Request-ID / X-Session-ID / X-Client-ID, где X-Request-ID уже соответствует ResponseMeta.request_id, назначаемому шлюзом. Библиотека trace_id не генерирует, не изменяет и не читает из ответа — только переносит, поэтому вернувшееся значение всегда равно переданному.

Куда попадает: ResponseMeta.trace_id (успех), атрибут trace_id любого GigaClientError (ошибка), строка лога ответа рядом с request_id, и GigaEvent.trace_id — то есть автоматически в ваши Postgres/OpenSearch-стоки, без обёрток над on_event.

Смысл в том, что на обратном пути RAG читает id с объекта ответа или ошибки, а не достаёт из ambient-контекста: это переживает границы контекста (очереди, хендоффы задач, буферы стрима), где contextvar рвётся.

PII

По умолчанию событие несёт только метрики и идентификаторы. Q&A-текст (prompt_text, prompt_messages, response_text, function_call) и полные схемы инструментов (tools_spec — в них человекописанные description и примеры) заполняются только при events_include_text=True (GIGA_EVENTS_INCLUDE_TEXT=true) — слать ли его, хэшировать или семплировать, решает сток, не либа. Поэтому текст и живёт в OpenSearch-проекции (её включают осознанно), а метрик-строка Postgres несёт лишь 256-символьное превью. app/env заполняет потребитель (в обёртке EventSink или враппере события); trace_id с 0.5.3 проставляет сама библиотека из аргумента вызова. Под выключатель не попадают tool_names/tool_choice/called_tool_name: имя инструмента — это код, а не пользовательский текст, и без него нельзя посчитать tool-использование, не включая PII.

Диагностика доступности (probe)

await client.probe(model) выполняет дешёвый реальный вызов (count_tokens, без генерации) и возвращает ProbeResult(ok, model, duration_ms, error), не бросая на ожидаемых сбоях — сеть, ошибка API, stop-event, таймаут слота дают ok=False с текстом в error. По умолчанию acquire_timeout=5.0 — проба fail-fast, не виснет на занятом пуле.

result = await client.probe("GigaChat")
if not result.ok:
    logger.warning(f"gigachat probe failed: {result.error}")

Это строительный блок для «контрольной корзины» приложения, а не Kubernetes readiness-проба: доступность GigaChat — внешняя зависимость и не должна гейтить readinessProbe сервиса. Проба идёт обычным путём, поэтому при включённом fallback отражает эффективную обслуживаемость (в т.ч. через резервную модель).

Тесты

python -m venv .venv
.venv\Scripts\pip install -e ".[dev]"
.venv\Scripts\python -m pytest

Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

gigamux-0.5.4.tar.gz (126.8 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

gigamux-0.5.4-py3-none-any.whl (73.0 kB view details)

Uploaded Python 3

File details

Details for the file gigamux-0.5.4.tar.gz.

File metadata

  • Download URL: gigamux-0.5.4.tar.gz
  • Upload date:
  • Size: 126.8 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: uv/0.11.16 {"installer":{"name":"uv","version":"0.11.16","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":null,"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":null}

File hashes

Hashes for gigamux-0.5.4.tar.gz
Algorithm Hash digest
SHA256 0c665d24e8d4acec382a035441f4684a56851d2b86ceff4694ef1c8603cfaa2c
MD5 d4861a44ba85542dc624e90a4369c03a
BLAKE2b-256 17e82bd719bed780ae0a607c523573cd238ff31b89c15a12b64b23cab5ae4257

See more details on using hashes here.

File details

Details for the file gigamux-0.5.4-py3-none-any.whl.

File metadata

  • Download URL: gigamux-0.5.4-py3-none-any.whl
  • Upload date:
  • Size: 73.0 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? No
  • Uploaded via: uv/0.11.16 {"installer":{"name":"uv","version":"0.11.16","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":null,"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":null}

File hashes

Hashes for gigamux-0.5.4-py3-none-any.whl
Algorithm Hash digest
SHA256 11022edf2d2feefc4ee6dd484db2f8131b6ab2ff8cde59cd8a9913ade1105329
MD5 ff5a25bce4d1716056872749200325d1
BLAKE2b-256 bd631ed1768cf73b398d5ad3ab061bbc58da42a19ecc9a6e49bd101c32e4d298

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page