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

3. Что разрешено собирать

Вернёмся к фразе из шага 1, теперь уже с последствием: register_stage кладёт класс в реестр, общий на весь процесс, а значит любая импортированная стадия доступна любому пайплайну внутри него.

Пока пайплайны ваши — это удобно. Как только они не ваши, это дыра. А «разработчик задаёт блоки, а собирает из них кто-то менее доверенный» — ровно та схема, ради которой StageFlow и существует. Без дополнительной защиты ChargeCardStage доступна любому JSON, который кто угодно отправил в /api/run.

Такая защита — Policy. Её задаёт хост; из пайплайна её нельзя ни прочитать, ни тем более расширить.

from stageflow import Limits, Policy

BASIC = Policy(
    stages={"LoadTicketStage", "ClassifyByRulesStage", "SearchKnowledgeStage",
            "RenderReplyStage", "SendReplyStage", "EscalateStage",
            "SetValueStage", "ConcatStage", "TemplateStage"},
    node_types={"entry", "stage", "condition", "switch", "terminal"},
    limits=Limits(
        counters={"seconds": 15, "steps": 200, "kb_lookups": 20, "replies_sent": 1},
        gauges={"concurrency": 2, "depth": 2, "frame_bytes": 200_000},
        max_retries=2, max_delay_seconds=2,
    ),
)

Три поля, и None в любом из них означает «ограничений нет» — это не то же самое, что пустое множество: Policy() разрешает всё, а Policy(stages=set()) не разрешает ни одной стадии.

Два рубежа, один и тот же ответ

pipeline.validate(policy)                       # до: весь граф целиком
Session(id=..., pipeline=..., policy=policy)    # во время: каждый узел, когда до него дошли

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

В примере валидация стоит прямо в обработчике запроса, ещё до запуска потока:

def start(self, payload: dict, policy: Policy) -> Run:
    pipeline = Pipeline.from_dict(payload.get("pipeline") or {})
    pipeline.validate(policy)      # отказ придёт ответом на POST, а не через минуту из потока
    ...
    session = Session(..., policy=policy)

Предупредить редактор, чтобы он предупредил автора

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

@router.get("/stages")
async def stages(caller: CallerPlan) -> dict:
    policy = policy_for(caller)
    return {"stages": {name: cls.get_specs()
                       for name, cls in get_stages().items()
                       if policy.allows_stage(name)}}


@router.get("/meta")
async def meta(caller: CallerPlan) -> dict:
    policy = policy_for(caller)
    return {"api": 1, "plan": caller, **capabilities(policy), "limits": _limits_of(policy)}

Отдать спецификацию запрещённой стадии — значит заставить редактор нарисовать её в палитре, а прогон потом её отвергнуть. Ради устранения такого расхождения политику и завели — а здесь оно воспроизводится на один эндпоинт позже. capabilities(policy) делает то же самое для типов узлов.

Вот что это даёт на дешёвом тарифе из примера, если открыть пайплайн, которому нужен map:

Граф, который тариф не запустит, — сказано до запуска

Четыре типа узлов в палитре погашены, на проблемном узле стоит метка, а в строке состояния написано per_ticket: plan 'basic' does not include a 'map' node — и всё это до того, как кто-нибудь нажал «Run».

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

Из limits статическая валидация берёт только то, что достоверно следует из JSON:

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

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

Потолок жёсткий

BudgetExceeded унаследован от BaseException и намеренно стоит вне иерархии StageFlowError. Блоки try ловят Exception, механизм retry — тоже, и если бы отказ по бюджету можно было поймать, любой тариф превращался бы в украшение: достаточно обернуть граф в try. Ловит его ровно одно место, Session.run, и завершает прогон так же, как это делает stop:

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

Дальше: во что обходится прогон — откуда в этих счётчиках берутся числа.