Оркестратор пайплайна
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 | Уникальный идентификатор этапа |
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, метрики, сохранение состояния), и пайплайн продолжается как обычно.
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() builder
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 или удалённого хранилища).