Перейти к содержанию

Лимиты: сколько пайплайну можно потратить

Пишете свой бэкенд?

Шаг 4 наполняет эти счётчики из настоящих стадий: какие резервируют, какие только начисляют и что из этого показывает отладочная панель.

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

from stageflow import Limits, Policy, Session

plan = Policy(
    stages={"LoadTicket", "ClassifyByRules", "LlmReply"},
    limits=Limits(
        counters={"seconds": 120, "steps": 5_000, "iterations": 2_000,
                  "tokens": 200_000, "llm_calls": 100},
        gauges={"concurrency": 8, "depth": 4, "frame_bytes": 4_000_000},
        max_retries=5,
        max_delay_seconds=30,
    ),
)

session = Session(id="run-1", pipeline=pipeline, context=ctx, policy=plan)
result = await session.run()
result.meters   # {"steps": 412, "seconds": 8.2, "tokens": 5120, "llm_calls": 9}

Один механизм: счётчики

Не три системы для времени, стоимости и размера, а одна, с двумя видами счётчиков:

Вид Поведение Примеры Когда проверяется
counter растёт и не убывает seconds, steps, iterations, tokens при каждом начислении
gauge мгновенное значение concurrency, depth, frame_bytes когда его занимают

steps и iterations начисляет рантайм, пики держит он же; стадии начисляют то, что израсходовали. concurrency — единственный пиковый счётчик, который придерживает, а не роняет: сто элементов на обработку это хороший пайплайн, которому просто не надо держать сто вызовов в воздухе разом, и отказ ему за ширину был бы отказом работе, а не её ограничением. Глубина и размер именно роняют — слишком глубоко это слишком глубоко, ожидание тут не поможет. Имена счётчиков ядро не фиксирует. Хост считает то, что дефицитно у него, а счётчик, который никто не ограничил, просто копится и приезжает в result.meters — из этого и собирается счёт.

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

Что не ограничено, то не измеряется

Любая проверка чего-то стоит, а frame_bytes дороже всех: измерить его — значит сериализовать фрейм. Счётчика нет в лимитах — он не вычисляется вовсе, так что вся эта механика обходится хосту без лимитов в один if.

Что тратит стадия

Стадия начисляет единицы, которые израсходовала: tokens, llm_calls, http_calls, rows. Сколько единица стоит — это прайс-лист; прайс-листы меняются без правки кода и принадлежат хосту. Стадия знает, сколько использовала, а не во что это обошлось.

@register_stage("LlmTriageStage")
class LlmTriageStage(BaseStage):
    """
    description: "Классификация тикета моделью"
    reserve:
      llm_calls: 1
      tokens: "args.max_tokens + size(args.text) / 3"
    """

    timeout = 60

    async def run(self):
        answer = await self.client.chat(self.get_arguments()["text"])
        self.set_outputs({"topic": answer.topic})
        self.charge(tokens=answer.usage.total, llm_calls=1)

Механизма два, и они не пересекаются:

Где Что видит На что отвечает
reserve спека, декларативно args можно ли вообще пытаться
charge() код стадии всё, что стадия знает сколько потрачено на самом деле

reserve декларативен потому, что читается до запуска: редактор покажет «до 4000 токенов за вызов», хост откажет графу не исполняя его. Значения — числа или CEL по args, то есть по аргументам в том виде, в каком их получит стадия.

Порядок — резерв → прогон → расчёт:

  1. резерв вычисляется и удерживается. Не влезает в остаток — стадия не запускается вовсе: дорогой вызов лучше не начинать, чем обрывать, когда деньги уже ушли;
  2. стадия работает;
  3. charge() замещает резерв по названным счётчикам и добавляет неназванные — сумма, и зарезервированная, и начисленная, не считается дважды. Стадия, ничего не начислившая, рассчитывается по резерву, и при падении тоже: вызов, который ушёл и отвалился по таймауту, всё равно потратил то, что потратил.

Стадия, которая ничего не резервирует и не начисляет, бесплатна.

Цикл целиком не резервируется

Резерв берётся по стадиям, в момент запуска каждой: сумма того, во что обойдётся map, перед входом в цикл не считается — резерв итерации зависит от её аргументов, аргументы от элемента, а элемент приходит из данных. Что проверяется заранее, так это число итераций: iterations начисляется до первого прохода, поэтому цикл на тысячу элементов на тарифе с полусотней отвергается, не выполнив ни одного.

То есть цикл останавливается, как только не может оплатить следующую стадию, а сделанные проходы уже оплачены. В режиме parallel перелёт ограничен не одной стадией, а concurrency: вызовы, которые уже в полёте, уже сделаны. Хосту, которому нужен жёсткий потолок на весь цикл, его даёт арифметика, а не предсказание: ограничить iterations так, чтобы их число на худший случай влезало в бюджет.

Что валидация отвергает до запуска

Pipeline.validate(policy) проверяет то, что достоверно выводится из JSON, — чтобы сохраняющий граф узнал, что менять, а не выяснял это на десятом элементе:

Что проверяется Почему это честно
кратчайший путь через граф против steps любой прогон проходит не меньше узлов, значит граф, чей самый дешёвый путь не влезает, не завершится вообще
глубина вложенности объявленных субпайплайнов против depth вложенность записана, а не вычисляется
что просит каждый retry против max_retries и max_delay_seconds числа лежат в JSON

Оценка шагов — нижняя граница, и именно поэтому по ней можно отказывать: граф с длинной дорогой и короткой судится по короткой.

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

Когда бюджет кончился

Потолок жёсткий. BudgetExceeded наследует BaseException и намеренно стоит вне иерархии StageFlowError: блоки try ловят Exception, и машинка retry тоже, — тенант, способный обернуть граф в try и проглотить это, превратил бы любой лимит в украшение.

Ловит его ровно одно место, Session.run, и заканчивает прогон так же, как останавливает stop, а не швыряет исключение сквозь работу, которая действительно была сделана:

result.result   # {"status": "budget_exceeded", "meter": "tokens",
                #  "limit": 200000, "spent": 200512}
result.meters   # всё потраченное — по нему хост и выставляет счёт

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

Ручки тенанта — это просьбы

retry и паузы между попытками приходят из пайплайна, то есть с менее доверенной стороны, поэтому max_retries и max_delay_seconds разбираются с ними дважды.

Узел, объявивший больше, чем разрешает тариф, отвергается на валидации, как и всё остальное, записанное в JSON: тенанту говорят, что менять, а не тихо выдают не то, что он просил. Retry, попавший в сессию другим путём — граф, собранный в Python и не проходивший валидацию, — зажимается: пятьдесят попыток на тарифе с тремя отработают три раза, и в событии будет написано три.

Ожидание человека — не работа

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

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

Чего лимиты не делают

concurrency: 8 ограничивает один прогон. Сто прогонов — это восемьсот задач: ядро не может знать, сколько прогонов существует, поэтому допуск к запуску — дело хоста, очередь и пул на тенанта. Так же и месячная квота: ядро держит потолок одного прогона и сообщает, сколько он потратил, а вычитать это из месячного лимита и переводить единицы в деньги — бухгалтерия платформы.