Skip to content

Оркестратор пайплайна

PipelineOrchestrator

Основной класс для построения и управления пайплайном из последовательных (и параллельных) этапов.

Конструктор

js
new PipelineOrchestrator({
  config,       // PipelineConfig — этапы и опциональный middleware
  httpConfig?,  // HttpConfig — настройки HTTP-клиента
  sharedData?,  // Record<string, any> — общий пул для всех этапов
  options?,     // { autoReset?: boolean }
})

Методы

МетодОписание
run(onStepPause?, externalSignal?)Выполнить все этапы. Возвращает { stageResults, success }
rerunStep(stepKey, options?)Повторно выполнить один этап (с учётом condition, before, after, middleware)
abort()Прервать выполнение пайплайна (отменяет текущий HTTP-запрос через AbortSignal)
isAborted()Проверить, был ли пайплайн прерван
pause()Поставить на паузу после завершения текущего этапа
resume()Продолжить пайплайн, стоящий на паузе
isPaused()Проверить, стоит ли пайплайн на паузе
exportState()Сериализовать stageResults и логи в обычный объект
importState(state)Восстановить stageResults и логи из снапшота
getStageResults()Синхронный снапшот всех результатов этапов
getRunId()ID текущего/последнего run() или rerunStep() — см. Корреляция запуска
destroy()Выполнить коллбэки очистки всех установленных плагинов
subscribeProgress(listener)Подписаться на обновления прогресса
subscribeStageResults(listener)Подписаться на изменения stageResults
subscribeStepProgress(stepKey, listener)Подписаться на прогресс конкретного этапа
on(eventName, handler)Подписаться на любое событие (step:<key>:start|success|error|skipped|progress, log)
onStepStart/Finish/Error(handler)Подписаться на события жизненного цикла этапа
getProgress()Получить текущий снапшот прогресса
getLogs()Получить все логи пайплайна (ограничены options.maxLogs записями, если задано — см. ниже)
clearStageResults()Сбросить результаты и прогресс

Параметры этапа (PipelineStageConfig)

Каждый хук ниже также получает signal: AbortSignal в объекте параметров — тот же сигнал, что использует orchestrator.abort(). Передавайте его дальше в fetch/axios/и т. д. внутри request/before/after, чтобы отмена действительно останавливала работу «на лету», а не только учёт состояния пайплайна.

ПараметрОписание
keyУникальный идентификатор этапа
request({ prev, allResults, sharedData, signal })Основная функция этапа — возвращаемое значение становится результатом этапа
condition({ prev, allResults, sharedData, signal })Если возвращает false, этап пропускается со статусом "skipped"
before({ prev, allResults, sharedData, signal })Хук предобработки — возвращённое значение заменяет prev, передаваемый в request
after({ result, allResults, sharedData, signal })Хук постобработки — возвращённое значение заменяет результат этапа
errorHandler({ error, key, sharedData, signal })Обработчик ошибок этапа — см. Восстановление после ошибок ниже
retryCountПереопределить количество повторов для этого этапа
timeoutMsПереопределить таймаут для этого этапа
pauseBeforeЗадержка в мс перед выполнением request
pauseAfterЗадержка в мс после выполнения request

Восстановление после ошибок (errorHandler + recoverStep)

По умолчанию всё, что вернёт errorHandler, оборачивается в ApiError, а этап остаётся в статусе "error" — он может преобразовать/дополнить ошибку, но не превратить провал в успех. Верните recoverStep(data), чтобы вместо этого восстановить этап: он фиксируется точно так же, как успешный (статус "success", middleware afterEach, метрики, сохранение состояния), и пайплайн продолжается как обычно.

js
import { recoverStep } from "rest-pipeline-js";

{
  key: "fetchPrice",
  request: async () => fetchPriceFromApi(),
  errorHandler: ({ error }) => {
    if (isNetworkError(error)) return recoverStep(0); // используем значение по умолчанию и продолжаем
    return error; // всё остальное: продолжаем считать провалом, как раньше
  },
}

Поток выполнения этапа

condition? → false → [status: skipped] → следующий этап
           ↓ true
middleware.beforeEach

pauseBefore

хук before()

request()

хук after()

pauseAfter

middleware.afterEach

[status: success] → следующий этап

При ошибке на любом шаге:
  └─► stage.errorHandler (если задан)
        ├─► возвращает recoverStep(data) → [status: success] → следующий этап
        └─► иначе → middleware.onError → [status: error] → остановка

Полный пример

js
import { PipelineOrchestrator } from 'rest-pipeline-js'

const orchestrator = new PipelineOrchestrator({
  config: {
    stages: [
      {
        key: 'fetchUser',
        request: async ({ sharedData }) => {
          const res = await fetch(`/api/users/${sharedData.userId}`)
          return res.json()
        },
      },
      {
        key: 'processData',
        condition: ({ prev }) => prev !== null,
        before: ({ prev }) => ({ ...prev, processed: true }),
        request: async ({ prev }) => prev,
        after: ({ result }) => ({ ...result, finishedAt: Date.now() }),
      },
    ],
    middleware: {
      beforeEach: ({ stage }) => console.log('Starting:', stage.key),
      afterEach: ({ stage, result }) => console.log('Done:', stage.key, result.data),
      onError: ({ stage, error }) => console.error('Error in', stage.key, error),
    },
  },
  httpConfig: {
    baseURL: 'https://api.example.com',
    retry: { attempts: 2, delayMs: 1000, backoffMultiplier: 2 },
    cache: { enabled: true, ttlMs: 60000 },
  },
  sharedData: { userId: 42 },
  options: { autoReset: true },
})

orchestrator.subscribeProgress((progress) => {
  console.log('Stage:', progress.currentStage, 'Statuses:', progress.stageStatuses)
})

orchestrator.on('step:fetchUser:success', (payload) => {
  console.log('fetchUser done:', payload.data)
})

const result = await orchestrator.run()
console.log('Pipeline finished:', result.success)
console.log('Stage results:', result.stageResults)

Параллельные этапы

Группируйте этапы для параллельного выполнения через parallel:

js
const orchestrator = new PipelineOrchestrator({
  config: {
    stages: [
      // Последовательный этап
      { key: 'auth', request: async () => getToken() },

      // Параллельная группа — все выполняются одновременно
      {
        key: 'load-data',
        parallel: [
          { key: 'loadUsers', request: async () => fetchUsers() },
          { key: 'loadProducts', request: async () => fetchProducts() },
          { key: 'loadSettings', request: async () => fetchSettings() },
        ],
      },

      // Последовательный этап после группы
      { key: 'render', request: async ({ allResults }) => render(allResults) },
    ],
  },
})
  • Все этапы в группе parallel выполняются одновременно через Promise.all — если только не задан concurrency (см. ниже).
  • Если любой этап группы проваливается, пайплайн останавливается и помечается как success: false.
  • У каждого параллельного этапа свой ключ и свой результат в stageResults.
  • rerunStep(key) работает и для этапов внутри параллельных групп.

Ограничение конкурентности

Для fan-out по множеству элементов (например, постраничных запросов) задайте concurrency на группе, чтобы ограничить, сколько этапов выполняется одновременно, вместо немедленного запуска их всех:

js
{
  key: "fetch-all-pages",
  parallel: pageNumbers.map((n) => ({
    key: `page-${n}`,
    request: async () => fetchPage(n),
  })),
  concurrency: 5, // не более 5 запросов одновременно
}

Результаты попадают в stageResults под своим ключом независимо от concurrency, в той же форме, что и у неограниченной группы. С билдером pipe(): .parallel(stages, { concurrency: 5 }).

Глобальные middleware

Применяйте хуки ко всем этапам без изменения конфигурации каждого из них:

js
const orchestrator = new PipelineOrchestrator({
  config: {
    stages: [/* ... */],
    middleware: {
      beforeEach: async ({ stage, index, sharedData }) => {
        console.log(`[${index}] Starting: ${stage.key}`)
        sharedData.startedAt = Date.now()
      },
      afterEach: async ({ stage, index, result, sharedData }) => {
        const ms = Date.now() - sharedData.startedAt
        console.log(`[${index}] Done: ${stage.key} in ${ms}ms`, result.data)
      },
      onError: async ({ stage, error, sharedData }) => {
        await reportError({ stage: stage.key, error, context: sharedData })
      },
    },
  },
})

Middleware выполняется в дополнение к (а не вместо) errorHandler конкретного этапа.

Пауза / Продолжение

Ставьте пайплайн на паузу после этапа и возобновляйте позже:

js
const orchestrator = new PipelineOrchestrator({ config })

// Пауза после завершения step1
orchestrator.on('step:step1:success', () => orchestrator.pause())

const runPromise = orchestrator.run()

// В какой-то момент позже (например, после подтверждения пользователем):
await showConfirmDialog()
orchestrator.resume()

await runPromise
  • pause() — пайплайн ждёт после завершения текущего этапа (включая события).
  • resume() — продолжает со следующего этапа.
  • abort() во время паузы разблокирует пайплайн и завершает его.

Экспорт / импорт состояния

Сохраняйте и восстанавливайте состояние пайплайна между перезагрузками страницы или сессиями:

js
const orchestrator = new PipelineOrchestrator({ config })
await orchestrator.run()

// Сохраняем состояние
const snapshot = orchestrator.exportState()
localStorage.setItem('pipelineState', JSON.stringify(snapshot))

// Позже — восстанавливаем и просматриваем без повторного запуска
const saved = JSON.parse(localStorage.getItem('pipelineState'))
const orchestrator2 = new PipelineOrchestrator({ config })
orchestrator2.importState(saved)

console.log(orchestrator2.getProgress()) // восстановленный прогресс
console.log(orchestrator2.getLogs()) // восстановленные логи (метки времени как объекты Date)

exportState() возвращает { stageResults, logs } — обычный JSON-сериализуемый объект. Метки времени в логах хранятся как ISO-строки и восстанавливаются как объекты Date при importState.

Ограничение роста логов (maxLogs)

logs растёт на одну запись за каждое событие шага и никогда не обрезается автоматически — нормально для одного run(), но экземпляр orchestrator, переиспользуемый на множестве запусков без autoReset (например, долгоживущий синглтон в SPA), накапливает логи бесконечно. Задайте options.maxLogs, чтобы хранить только N самых свежих записей (старые вытесняются первыми):

js
const orchestrator = new PipelineOrchestrator({
  config: {
    stages: [/* ... */],
    options: { maxLogs: 500 },
  },
})

Без maxLogs поведение не отличается от предыдущих версий.

Метрики пайплайна

Наблюдайте за выполнением пайплайна без изменения логики этапов:

js
const orchestrator = new PipelineOrchestrator({
  config: {
    stages: [/* ... */],
    metrics: {
      onPipelineStart: ({ timestamp, runId }) => {
        console.log(`[${runId}] Pipeline started at`, new Date(timestamp).toISOString())
      },
      onPipelineEnd: ({ durationMs, success, stageResults, runId }) => {
        analytics.track('pipeline_complete', { durationMs, success, runId })
      },
      onStepDuration: ({ stepKey, durationMs, status, runId }) => {
        console.log(`[${runId}] [${stepKey}] ${status} in ${durationMs}ms`)
      },
    },
  },
})
КоллбэкПолучаетОписание
onPipelineStart{ timestamp, runId }Срабатывает в начале run()
onPipelineEnd{ durationMs, success, stageResults, runId }Срабатывает по завершении run()
onStepDuration{ stepKey, durationMs, status, runId }Срабатывает после каждого выполненного шага

Корреляция запуска (runId)

Каждый вызов run() генерирует новый runId (UUID либо запасной вариант на основе timestamp в окружениях без crypto.randomUUID), общий для всех коллбэков метрик, записей логов (getLogs()) и событий шагов (PipelineStepEvent.runId), произведённых во время этого запуска — включая все попытки pipelineRetry. rerunStep() генерирует свой отдельный runId. Используйте orchestrator.getRunId(), чтобы прочитать текущий/последний, либо читайте runId из любого события/лога/коллбэка метрик, чтобы связать всё, что произошло за одно выполнение, в вашем backend'е логирования/трассировки:

js
orchestrator.on('log', (entry) => sendToLogBackend({ ...entry, runId: orchestrator.getRunId() }))

createPipeline() + pipe() builder

createPipeline() — короткая фабрика

js
import { createPipeline } from 'rest-pipeline-js'

const orchestrator = createPipeline(
  [
    { key: 'fetchUser', request: async () => fetchUser() },
    { key: 'process', request: async ({ prev }) => process(prev) },
  ],
  {
    httpConfig: { baseURL: 'https://api.example.com' },
    sharedData: { userId: 42 },
    pipelineOptions: { continueOnError: false },
    metrics: {
      onStepDuration: ({ stepKey, durationMs }) => console.log(stepKey, durationMs),
    },
  },
)

pipe() — fluent-билдер

js
import { pipe } from 'rest-pipeline-js'

const orchestrator = pipe()
  .step({ key: 'auth', request: async () => getToken() })
  .step({ key: 'fetchUser', request: async ({ prev }) => fetchUser(prev) })
  .parallel([
    { key: 'loadPosts', request: async () => fetchPosts() },
    { key: 'loadNotifs', request: async () => fetchNotifications() },
  ])
  .stream({
    key: 'liveUpdates',
    stream: async function* () {
      yield* subscribe('/events')
    },
    onChunk: (chunk) => updateUI(chunk),
  })
  .build({ httpConfig: { baseURL: 'https://api.example.com' } })
Метод билдераОписание
.step(stage)Добавить последовательный этап
.parallel(stages, options?)Добавить параллельную группу (key/concurrency опциональны, см. Ограничение конкурентности)
.subPipeline(item)Встроить суб-пайплайн как этап
.stream(stage)Добавить потоковый этап (AsyncIterable)
.build(options?)Создать и вернуть PipelineOrchestrator
.toConfig(options?)Вернуть PipelineConfig без создания оркестратора

Типизированная цепочка вызовов (TypeScript)

В TypeScript pipe().step(...) отслеживает тип prev по всей цепочке: prev каждого .step() типизирован как возвращаемое значение предыдущего шага (undefined для самого первого шага, что соответствует реальному поведению оркестратора во время выполнения). .parallel() / .subPipeline() / .stream() его не меняют — точно как во время выполнения, где prev для следующего шага по-прежнему берётся из последнего обычного .step(), а не из результатов параллельной группы:

ts
const orchestrator = pipe()
  .step({ key: 'auth', request: async (): Promise<string> => getToken() })
  .step({ key: 'fetchUser', request: async ({ prev }) => fetchUser(prev) }) // prev: string — выведено, автодополнение
  .step({ key: 'oops', request: async ({ prev }) => prev.totallyNotAMethod() }) // ✗ ошибка компиляции: неверный тип для prev
  .build()

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

Валидация схемы (validateInput / validateOutput)

Каждый этап может валидировать (и, поскольку возвращаемое значение заменяет данные, опционально приводить к нужному типу) свои входные и выходные данные через validateInput/validateOutput. Ни один из них не зависит от конкретной библиотеки схем — передайте любую функцию (data) => T; она должна бросать исключение при некорректных данных:

ts
import { createRestClient, PipelineOrchestrator } from 'rest-pipeline-js'
import { z } from 'zod' // не зависимость этого пакета — подключите свою

const userSchema = z.object({ id: z.number(), name: z.string() })

const orchestrator = new PipelineOrchestrator({
  config: {
    stages: [
      {
        key: 'fetchUser',
        request: async ({ sharedData }) => client.get(`/users/${sharedData.userId}`),
        // Валидирует форму ответа перед тем, как он будет сохранён как результат
        // этого этапа и передан следующему этапу как `prev` — отлавливает
        // изменение контракта backend'а на границе пайплайна, а не ниже по потоку.
        validateOutput: (data) => userSchema.parse((data as { data: unknown }).data),
      },
      {
        key: 'greet',
        validateInput: (data) => userSchema.parse(data),
        request: async ({ prev }) => `Hello, ${prev.name}!`,
      },
    ],
  },
})
  • validateInput выполняется после хука before (видит его результат), непосредственно перед request — так что может валидировать/приводить значение, которое request вот-вот получит как prev.
  • validateOutput выполняется после хука after (видит его результат), непосредственно перед тем, как шаг фиксируется как успешный — так что валидирует/приводит фактическое значение, которое становится data этого этапа и prev следующего.
  • Оба получают (data, { allResults, sharedData, signal }), как и другие хуки этапа.
  • Выброшенная ошибка проходит тот же самый путь, что и любая другая ошибка этапа — errorHandler может её осмотреть и вернуть recoverStep(fallbackValue) для восстановления, точно так же, как при провале request.
  • Оба также применяются к этапам внутри ParallelStageGroup (они разделяют тот же путь выполнения, что и этапы верхнего уровня).

Полный вспомогательный адаптер withZodSchema() см. в examples/zod-validation.ts.

validatePipelineConfig()

Отлавливайте ошибки конфигурации ещё до выполнения:

js
import { validatePipelineConfig } from 'rest-pipeline-js'

const { valid, errors } = validatePipelineConfig({
  stages: [
    { key: 'step1', request: async () => data },
    { key: 'step1', request: async () => other }, // дубликат!
    { key: '', request: async () => other }, // пустой ключ!
  ],
})

if (!valid) console.error(errors)
// ["[root] duplicate stage key: "step1"", "[root] stage key must be a non-empty string"]

Валидирует: дублирующиеся ключи, пустые/некорректные ключи, пустой массив stages, некорректные типы полей (request, condition, retryCount, timeoutMs), а также рекурсивно валидирует вложенные конфигурации subPipeline.

Система плагинов

Упаковывайте переиспользуемое поведение оркестратора в плагины:

js
const loggingPlugin = {
  name: 'logging',
  install(orchestrator) {
    const off = orchestrator.on('log', (event) => {
      if (event.type === 'step:success') console.log('✓', event.stepKey)
      if (event.type === 'step:error') console.error('✗', event.stepKey, event.error)
    })
    return () => off() // очистка при orchestrator.destroy()
  },
}

const orchestrator = new PipelineOrchestrator({
  config: {
    stages: [/* ... */],
    options: { plugins: [loggingPlugin, analyticsPlugin] },
  },
})

// Вызывайте, когда оркестратор больше не нужен:
orchestrator.destroy()
  • install(orchestrator) — получает экземпляр оркестратора; может подписываться на события, настраивать middleware и т. д.
  • Если install возвращает функцию, она регистрируется как коллбэк очистки и вызывается через destroy().

Адаптер сохранения состояния

Автоматически сохраняйте и восстанавливайте состояние пайплайна между перезагрузками страницы:

js
const localStorageAdapter = {
  save: (state) => localStorage.setItem('pipeline', JSON.stringify(state)),
  load: () => {
    const raw = localStorage.getItem('pipeline')
    return raw ? JSON.parse(raw) : null
  },
}

const orchestrator = new PipelineOrchestrator({
  config: {
    stages: [/* ... */],
    options: { persistAdapter: localStorageAdapter },
  },
})

// run() загружает сохранённое состояние при старте; сохраняет после каждого завершённого шага
await orchestrator.run()

Интерфейс адаптера:

ts
type PipelineStateAdapter = {
  save(state: PipelineExportedState): void | Promise<void>
  load(): PipelineExportedState | null | Promise<PipelineExportedState | null>
}

Оба метода могут быть асинхронными (полезно для IndexedDB или удалённого хранилища).