PipelineOrchestrator
Основной класс для построения и управления пайплайном из последовательных (и параллельных) этапов.
Конструктор
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
string
Уникальный идентификатор этапа.
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
number
Переопределить количество повторов для этого этапа.
timeoutMs
number
Переопределить таймаут для этого этапа.
pauseBefore
number
Задержка в мс перед выполнением request.
pauseAfter
number
Задержка в мс после выполнения request.
Восстановление после ошибок (errorHandler + recoverStep)
По умолчанию всё, что вернёт errorHandler, оборачивается в ApiError, а этап остаётся в статусе "error" — он может преобразовать/дополнить ошибку, но не превратить провал в успех. Верните recoverStep(data), чтобы вместо этого восстановить этап: он фиксируется точно так же, как успешный (статус "success", middleware afterEach, метрики, сохранение состояния), и пайплайн продолжается как обычно.
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] → остановкаПолный пример
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:
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 на группе, чтобы ограничить, сколько этапов выполняется одновременно, вместо немедленного запуска их всех:
{
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
Применяйте хуки ко всем этапам без изменения конфигурации каждого из них:
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 конкретного этапа.
Пауза / Продолжение
Ставьте пайплайн на паузу после этапа и возобновляйте позже:
const orchestrator = new PipelineOrchestrator({ config })
// Пауза после завершения step1
orchestrator.on('step:step1:success', () => orchestrator.pause())
const runPromise = orchestrator.run()
// В какой-то момент позже (например, после подтверждения пользователем):
await showConfirmDialog()
orchestrator.resume()
await runPromisepause()— пайплайн ждёт после завершения текущего этапа (включая события).resume()— продолжает со следующего этапа.abort()во время паузы разблокирует пайплайн и завершает его.
Экспорт / импорт состояния
Сохраняйте и восстанавливайте состояние пайплайна между перезагрузками страницы или сессиями:
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 самых свежих записей (старые вытесняются первыми):
const orchestrator = new PipelineOrchestrator({
config: {
stages: [/* ... */],
options: { maxLogs: 500 },
},
})Без maxLogs поведение не отличается от предыдущих версий.
Метрики пайплайна
Наблюдайте за выполнением пайплайна без изменения логики этапов:
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'е логирования/трассировки:
orchestrator.on('log', (entry) => sendToLogBackend({ ...entry, runId: orchestrator.getRunId() }))Билдеры пайплайна (createPipeline / pipe)
createPipeline() — короткая фабрика
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-билдер
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(), а не из результатов параллельной группы:
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; она должна бросать исключение при некорректных данных:
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)
Отлавливайте ошибки конфигурации ещё до выполнения:
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.
Система плагинов
Упаковывайте переиспользуемое поведение оркестратора в плагины:
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().
Адаптер сохранения состояния
Автоматически сохраняйте и восстанавливайте состояние пайплайна между перезагрузками страницы:
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()Интерфейс адаптера:
type PipelineStateAdapter = {
save(state: PipelineExportedState): void | Promise<void>
load(): PipelineExportedState | null | Promise<PipelineExportedState | null>
}Оба метода могут быть асинхронными (полезно для IndexedDB или удалённого хранилища).