rest-pipeline-js 2.1.0: живое демо вместо документации и вычищенные гонки

В 2.0.0 были закрыты вопросы для бэкенда — распределённый rate limiter, circuit breaker, кэш, трассировка, идемпотентность. Фичи получились честные, но проверить их вживую было особо негде: демо-приложение — это пять страниц про покупку авиабилетов, бьющих в реальный внешний сервис, который мог просто лечь посреди показа.
2.1.0 — релиз на трёх ногах. Во-первых, восемь новых возможностей: более умный джиттер в retry, проактивный троттлинг по заголовкам rate limit, валидация схем на границах стадий, утилита пагинации, mock-адаптер для тестов, WebSocket как полноценный тип стадии, офлайн-очередь и, наконец, UMD/CDN-сборка. Во-вторых — новый демо-стенд, который наконец показывает всё это в работе, а не в README. В-третьих — то, что не видно на скриншотах: пакет прогнали через плотное ревью, нашли и закрыли десяток гонок и утечек в конкурентном коде, и переложили src/ из одной плоской кучи файлов в структуру по доменам.
Обратной несовместимости в этом релизе нет — все изменения либо аддитивные, либо чисто внутренние.
Новые возможности
Настраиваемый джиттер в retry
RequestExecutor и раньше умел ретраить с backoff, но добавлял к задержке фиксированный джиттер — плюс-минус 10% сверху номинального значения, без права выбора. Проблема этого подхода не в самой формуле, а в том, что происходит на масштабе: если у бэкенда лежит сотня клиентов, и все ретраят с одинаково узким разбросом задержек, второй залп повторных запросов приходит почти синхронно — начинается thundering herd прямо на восстанавливающийся сервис.
Теперь retry.jitterStrategy принимает три значения:
ts
const client = createRestClient({
baseURL: "https://api.example.com",
retry: { attempts: 5, delayMs: 200, backoffMultiplier: 2, jitterStrategy: "full" },
});
"fixed"(по умолчанию) — прежнее поведение, для обратной совместимости."full"и"decorrelated"— реализация алгоритмов из статьи AWS Exponential Backoff and Jitter, которые куда лучше размазывают повторные попытки множества параллельных клиентов по времени, чем узкий фиксированный разброс.
На Retry-After это не влияет — значение из заголовка всегда берётся как есть, сервер лучше знает, когда можно повторить. Заодно исправлена мелкая несостыковка: fallback-путь (backoff, когда Retry-After присутствует, но не парсится) раньше вообще не получал джиттер — теперь получает, как и обычный backoff.
Проактивный троттлинг по заголовкам rate limit
Раньше rate limiter в пакете реагировал только постфактум — получил 429 с Retry-After, подождал, повторил. Но многие API честно сообщают, сколько запросов у вас осталось, в каждом ответе: X-RateLimit-Remaining, черновой IETF-стандарт RateLimit-Remaining или что-то своё у конкретного вендора. Дожидаться реального 429, чтобы отреагировать на эту информацию, — терять данные, которые уже были в руках.
ts
const client = createRestClient({
baseURL: "https://api.example.com",
rateLimit: {
onRateLimitHeaders: (headers, control) => {
const remaining = Number(headers["x-ratelimit-remaining"]);
const resetSec = Number(headers["x-ratelimit-reset"]);
if (remaining === 0 && Number.isFinite(resetSec)) {
control.throttleFor(resetSec * 1000);
}
},
},
});
Коллбэк вызывается после каждого ответа — и успешного, и ошибочного (429 обычно несёт те же заголовки, что и обычный ответ). Пакет намеренно не парсит конкретный формат заголовков сам — единого стандарта нет, а значит, разбор — забота потребителя (в examples/proactive-rate-limit-throttling.ts есть заготовки под три частые схемы). control.throttleFor(ms) не заменяет maxConcurrent/maxRequestsPerInterval, а комбинируется с ними — сдвигает следующий acquire() минимум на ms вперёд, причём более поздний короткий вызов не сокращает уже запланированное более длинное ожидание. Работает и поверх распределённого rateLimit.store — троттлинг применяется локально, ещё до похода в общий стор.
Валидация схем на входе и выходе стадии (validateInput/validateOutput)
Раньше несовпадение формата ответа с ожиданиями всплывало где-то в глубине пайплайна — стадия N+3 падала с невнятным Cannot read property of undefined, потому что стадия N вернула не то, что от неё ждали. Теперь стадия может проверить (и, поскольку возвращаемое значение подменяет данные, — скорректировать) свой вход и выход прямо на границе:
ts
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}`),
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; validateOutput — после after, прямо перед тем, как результат зафиксируется как успешный. Оба разделяют ту же сигнатуру (data, { allResults, sharedData, signal }), что и остальные хуки стадии, и работают внутри ParallelStageGroup без дополнительного кода — просто потому что используют тот же путь исполнения, что и обычные стадии верхнего уровня. Библиотека не тянет zod или другую схемную библиотеку как зависимость — подходит любая функция (data) => T, которая бросает исключение на невалидных данных; брошенная ошибка идёт по тому же пути, что и обычная ошибка стадии, значит errorHandler может её поймать и вызвать recoverStep(fallback), как при упавшем request.
Пагинация: paginate() / paginateAll() / flattenPages()
Курсорная и офсетная пагинация — два разных по форме, но одинаковых по сути API, которые раньше каждый обрабатывал вручную. Теперь это одна утилита, скрывающая разницу:
ts
import { paginate, paginateAll, flattenPages } from "rest-pipeline-js";
// Cursor-based
for await (const page of paginate({
fetchPage: (cursor) => client.get("/items", { params: { cursor } }).then((r) => r.data),
})) {
console.log(page.length, "items");
}
// Offset/limit-based
for await (const page of paginate({
strategy: "offset",
limit: 50,
fetchPage: (offset, limit) =>
client.get("/items", { params: { offset, limit } }).then((r) => r.data),
})) {
console.log(page.length, "items");
}
fetchPage возвращает { items, nextCursor } (курсорная стратегия — стоп, когда nextCursor — null/undefined) или { items, total? } (офсетная — стоп, когда страница короче limit, или offset достиг total, если API его сообщает). paginateAll() разворачивает всё в один плоский массив — самый простой вариант, если датасет заведомо помещается в память. flattenPages() превращает поток страниц в поток отдельных элементов — то, что нужно, если onChunk должен стрелять на каждый элемент, а не на каждую страницу; естественный источник для StreamStageConfig.stream. И paginate(), и сам fetchPage принимают signal — abort() останавливает обход между страницами, а не только между запросами.
createMockAdapter() — тестирование без реального бэкенда
Новая, отдельная точка входа rest-pipeline-js/testing — по аналогии с /vue и /react, чтобы тестовый код не протекал в продакшен-бандл (сама точка входа — 603 B brotli). Заменяет сеть набором маршрутов:
ts
import { createRestClient } from "rest-pipeline-js";
import { createMockAdapter } from "rest-pipeline-js/testing";
const adapter = createMockAdapter([
{ method: "GET", url: "/users/1", respond: { data: { id: 1, name: "Ada" } } },
// Динамический ответ — читает сам запрос
{
method: "POST",
url: "/orders",
respond: (info) => ({ data: { id: 42, ...(info.data as object) }, status: 201 }),
},
// Последовательность ответов — первые две попытки падают, третья успешна
{
method: "GET",
url: "/flaky",
respond: [{ error: true, status: 503 }, { error: true, status: 503 }, { data: { ok: true } }],
},
]);
const client = createRestClient({ baseURL: "https://api.example.com", adapter });
url матчится по подстроке (string) или через .test() (RegExp) относительно config.url, не полного URL. Ответ со status >= 400 по умолчанию отклоняется — как это делают axios/fetch — с .status/.response.status/.response.data/.response.headers, этого достаточно, чтобы retry.retriableStatus, circuitBreaker и error-интерсепторы отработали как на настоящей сети, без адаптации теста под мок. adapter.calls хранит историю всех запросов по порядку — удобно для ассертов. Запрос без подходящего маршрута сразу бросает понятную ошибку вместо того, чтобы тест молча завис.
WebSocket как тип стадии пайплайна
До этой версии единственным «живым» транспортом в пайплайне были SSE-подобные стримы (StreamStageConfig). Чаты, presence-фиды, живые ордербуки, совместное редактирование — всё, что требует полнодуплексного соединения, — не вписывалось. Теперь WebSocketStageConfig встаёт рядом со стрим-стадией на равных:
js
const orchestrator = pipe()
.step({ key: "auth", request: async () => getToken() })
.websocket({
key: "chatFeed",
url: ({ prev }) => `wss://chat.example.com/rooms/general?token=${prev}`,
onOpen: () => console.log("connected"),
onMessage: (data) => JSON.parse(data),
onChunk: (message, sharedData) => updateUI(message),
closeOn: (message) => message.text === "__end__",
onClose: ({ wasClean }) => console.log("closed, clean:", wasClean),
onError: (error) => console.error(error),
timeoutMs: 5 * 60_000,
})
.build();
Сообщения, вернувшиеся из onMessage (может быть async), собираются в массив результата стадии — тот же паттерн, что чанки у стрим-стадии; onChunk стреляет по каждому сообщению в реальном времени. Успех/ошибка решается по событию close, а не error напрямую — большинство реализаций WebSocket шлют error прямо перед close, так что onError сам по себе стадию не роняет: чистое закрытие (wasClean: true) — успех со всем, что успело собраться, нечистое — ошибка по обычному пути continueOnError. closeOn(data) позволяет завершить стадию по содержимому сообщения, не дожидаясь, пока сервер сам закроет соединение. createWebSocket по умолчанию берёт globalThis.WebSocket (браузер, Deno, Node ≥22) — для более старого Node достаточно передать фабрику поверх пакета ws, тот же паттерн инъекции, что у HttpAdapter.
Офлайн-очередь
Мутирующий запрос, отправленный без связи, раньше просто падал с сетевой ошибкой — обработка «повторить, когда вернётся интернет» была целиком заботой вызывающего кода. Теперь это встроено:
js
import { createRestClient, OfflineQueuedError } from "rest-pipeline-js";
const client = createRestClient({
baseURL: "https://api.example.com",
offlineQueue: {
enabled: true,
persistAdapter: {
save: (queue) => localStorage.setItem("offline-queue", JSON.stringify(queue)),
load: () => JSON.parse(localStorage.getItem("offline-queue") ?? "null"),
},
onFlushSuccess: (request, response) => console.log("synced", request.url, response.data),
onFlushError: (request, error) => console.error("failed permanently", request.url, error),
},
});
try {
await client.post("/orders", cart);
} catch (err) {
if (err instanceof OfflineQueuedError) {
console.log("Order queued, will sync automatically:", err.queueId);
} else {
throw err;
}
}
По умолчанию в очередь уходят мутирующие методы (POST/PUT/PATCH/DELETE) — shouldQueue можно переопределить, GET по умолчанию никогда не ставится в очередь (устаревшее чтение бессмысленно «переигрывать»). persistAdapter — тот же интерфейс PipelineStateAdapter, который уже использовался для сохранения состояния пайплайна, а не новая сущность. isOnline/onOnlineChange по умолчанию завязаны на navigator.onLine и браузерное событие "online"; для Node/React Native подставляется своя реализация (например, NetInfo). Каждый запрос в очереди получает Idempotency-Key, переиспользуемый на каждой попытке реплея — тот же механизм, что и в ручной идемпотентности, так что бэкенд, который его поддерживает, не продублирует мутацию, реально прошедшую прямо перед обрывом связи. client.getQueuedRequests()/client.flushQueue() — для бейджа «N действий ждут синхронизации» и ручного форс-флаша.
UMD/CDN сборка
Раньше пакет требовал бандлер или Node — собрать HTML-страничку с одним <script> было невозможно. Теперь npm run build последним шагом выпускает самодостаточный IIFE:
html
<script src="https://unpkg.com/rest-pipeline-js@2.1.0/dist/umd/rest-pipeline.umd.min.js"></script>
<script>
const { createRestClient } = window.RestPipeline;
const client = createRestClient({ baseURL: "https://api.example.com" });
</script>
axios забандлен внутрь — он и так обычная (не peer) зависимость пакета, и это даёт настоящий опыт «один тег <script>», ничего больше подключать не нужно. Сборка собирается esbuild'ом уже из скомпилированного dist/esm/index.js, а не заново из src/, так что ESM/CJS/UMD гарантированно расходятся из одного и того же кода, а не трёх параллельных сборок с шансом разойтись поведением. Бюджет в .size-limit.json — 28 KB brotli, факт — ~25.3 KB, почти вровень с основным ESM-бандлом (разница только в IIFE-обвязке вместо tree-shaking).
Новое демо: CI/CD-пайплайн и торговый терминал вместо авиабилетов
Старое демо было живой иллюстрацией собственной проблемы: пять страниц («рейс», «параллельные запросы», «retry», «кэш», «трассировка») хором дёргали macrulez-api.ru, и если этот сервис моргал — моргало и демо, никак не связанное с качеством самого пакета. Заодно оно показывало едва ли треть того, что пакет реально умеет — WebSocket-стадии, DAG-переходы, sub-pipeline, плагины, auth provider, пагинация, схемная валидация, проактивный троттлинг, офлайн-очередь нигде не были видны вживую.
Теперь у демо есть свой сервер. Vite dev-плагин (demo/server/) поднимает HTTP-роуты и WebSocket поверх уже существующего dev-сервера — npm run demo:vue по-прежнему один процесс, просто теперь он же и бэкенд. Слой управляемых сбоев (demo/server/flaky.ts) даёт то, чего не даст внешний сервис: задержки, детерминированное «зафейлить N попыток подряд, потом починиться», малформленные ответы по флагу — retry, circuit breaker и schema-валидацию теперь можно увидеть надёжно и по требованию, а не ждать, пока настоящий бэкенд решит икнуть сам.
Два флагманских сценария:
CI/CD Pipeline — параллельная сборка сервисов (.parallel()), вложенный subPipeline для frontend (build → test), нестабильный test-шаг с видимым retry/backoff и выбором jitterStrategy прямо из UI, ручной approval-гейт на настоящем pause()/resume() перед деплоем, next()-переход на hotfix-ветку в обход тестов, живая консоль сборки на .stream() поверх SSE (автоскроллится к последней строке по мере поступления новых), PipelinePlugin, зеркалящий внутренние события оркестратора в консоль, отдельная песочница circuit breaker с настоящим getCircuitBreakerState(), и exportState()/importState() через localStorage — можно поставить сборку на паузу, перезагрузить страницу и продолжить с того же места.
Trading Terminal — логин через AuthProvider, параллельная загрузка позиций и котировок (кэшируемых), живой тикер цен на .websocket()-стадии со спарклайнами (hand-rolled SVG, без графической библиотеки) и flash-анимацией на изменение цены, размещение ордера с autoIdempotencyKey (двойной клик безопасен — плюс кнопка «повторить тот же ордер» для наглядного дедупа), настоящие убывающие заголовки X-RateLimit-*, включающие проактивный троттлинг через onRateLimitHeaders/throttleFor(), переключатель «офлайн», уводящий ордера в offlineQueue до восстановления связи, и circuit breaker на эндпоинте брокера для имитации отказа.
Ещё четыре страницы — не флагманы, а витрины отдельных фич, которым не нашлось естественного места в первых двух: изолированный цикл AuthProviderDemo (login → 401 → refresh → retry в отдельном логе, без шума остального терминала), PaginationDemo (paginate() по истории сборок с переключателем cursor/offset на одном датасете), SchemaValidationDemo (validateOutput на ответе размещения ордера, три ветки — успех / спасено recoverStep() / провал) и OrchestrationDemo (компактный проход по next(), subPipeline, PipelinePlugin, exportState()/importState(), по-настоящему восстановленному в свежий инстанс оркестратора, без визуального шума дашбордов).
Дизайн — тёмный «пульт управления»: живая консоль/лог-панель, спарклайны, тикер с flash-анимацией. Консоль сборки и лог событий терминала автоскроллятся к последней записи при поступлении новых строк — не нужно тянуться к низу вручную во время демонстрации.
Найденные и исправленные гонки
Прогнали пакет через плотное ревью — нашли и закрыли десяток конкурентных багов, часть из которых требовалась специально сконструированным тестом, чтобы вообще проявиться.
Rate limiter мог разбудить лишний запрос на один освободившийся слот. drainQueue() инкрементировал счётчик активных запросов только внутри продолжения самого acquire() — если очередь освобождающихся слотов разбирал синхронный цикл, он успевал разбудить больше waiter'ов, чем реально было свободных мест. Инкремент перенесён внутрь drainQueue(), синхронно и до пробуждения каждого waiter'а. Отдельно — waitForWindow() резервировал timestamp отдельным шагом после проверки вместимости (классический check-then-act); теперь резервирование и проверка — одна атомарная операция внутри цикла.
Распределённый rate limiter мог уйти в бесконечный retry-луп. _acquireViaStore() при постоянной перегрузке backend-стора не имел верхней границы по времени. Добавлен дедлайн — по истечении лимитер fail-open'ит вместо того, чтобы держать запрос вечно.
Офлайн-очередь могла отправить один и тот же запрос дважды или удалить не тот элемент. Два независимых flush(), вызванных почти одновременно (например, дважды сработавшее событие online), гонялись за одной и той же очередью — добавлен guard, запрещающий параллельный flush(). Отдельно: удаление отправленного элемента шло через queue.shift() — если за время await в enqueue() сработал maxQueueSize-трим и сдвинул очередь, shift() удалял не тот запрос. Заменено на удаление по id.
WebSocket-стадия могла обработать сообщения не по порядку — и зависнуть навсегда при ошибке в onClose. Если новое message-событие приходило раньше, чем завершался асинхронный onMessage предыдущего, обработка могла перекрыться — collected/onChunk/прогресс-события шли не в том порядке, в котором пришли данные. Обработчик теперь сериализован через цепочку промисов: каждое сообщение ждёт завершения предыдущего перед тем, как начать своё. Отдельно — исключение внутри пользовательского onClose раньше оставляло промис стадии висеть навечно (resolve/reject просто не вызывались); теперь обёрнуто в try/catch.
Слушатель событий утекал при отмене пайплайна через два сигнала одновременно. mergeSignals() подписывался на оба входных сигнала с { once: true } — сработавший сам себя отписывал, а второй, так и не сработавший, оставался подписан навсегда. Теперь при срабатывании любого из сигналов слушатели снимаются с обоих.
Stream- и WebSocket-стадии иногда исполнялись как обычные HTTP-стадии. findStageByKey() (используется, в частности, rerunStep()) исключал из обычного поиска только subPipeline-стадии — стадии с .stream()/.websocket() могли попасть в код, ожидающий обычную request-функцию. Добавлено то же исключение для них.
Два клиента с разными auth-провайдерами могли получить один и тот же закэшированный инстанс. Ключ кэша getRestClient() для присутствия auth собирался как !!config.auth — булев флаг «есть хоть какой-то провайдер», без учёта того, какой именно. Два конфига с разными провайдерами, но иначе идентичные, схлопывались в один клиент. Теперь ключ — стабильный id из WeakMap<AuthProvider, string>, присваиваемый каждому провайдеру при первом использовании.
Снимок прогресса пайплайна мутировал уже отданные наружу снимки. getProgress() копировал верхний объект ({ ...this.progress }), но массив stageStatuses внутри оставался общим по ссылке между всеми снимками — следующее изменение прогресса незаметно меняло уже возвращённые ранее данные. Добавлено глубокое клонирование массива в снимке.
React-хук логов не обновлялся сразу при смене оркестратора. usePipelineLogsReact инициализировал состояние через lazy useState-инициализатор, который запускается один раз при монтировании — при передаче нового orchestrator в тот же компонент старые логи продолжали висеть на экране, пока новый оркестратор сам не выстрелит первым событием. Добавлен немедленный setLogs() внутри useEffect.
createMockAdapter() неверно матчил повторно используемый global/sticky RegExp. RegExp.test() на паттерне с флагом g/y двигает lastIndex при каждом вызове — один и тот же маршрут с таким паттерном чередовал true/false на последовательных запросах вместо стабильного матчинга. lastIndex теперь сбрасывается перед каждой проверкой.
Наведение порядка в src/
Без изменений публичного API или поведения — только структура каталогов.
pipeline-orchestrator.ts (1677 строк) разошёлся на фасад и пять извлечённых модулей (pause-resume, stream-stage, websocket-stage, sub-pipeline, state-persistence) в src/pipeline/orchestrator/ — паттерн «свободная функция + typed context-параметр», сам класс стал тонким делегатом (1677 → 1225 строк). types.ts (1276 строк) разбит по доменам на types/http.ts/types/pipeline.ts/types/plugins.ts, с реэкспортом из types.ts для обратной совместимости импортов.
Дальше — просто расселение по папкам: Vue/React-хуки переехали в src/plugins/vue//src/plugins/react/ (с удалением избыточных суффиксов -vue/-react в именах файлов — папка уже говорит, какой это фреймворк), а сам HTTP-клиент и пайплайн-движок — в src/http/ и src/pipeline/. Публичные пути импорта (rest-pipeline-js, /vue, /react, /testing) не изменились ни на символ. Заодно снесён src/vue-demo/ — мёртвый, никуда не подключённый прекурсор текущего демо.
Все проверки после каждого шага — тесты, coverage, lint, размер бандлов — дают то же самое число, что и до переноса.
Сценарии использования
Ниже — четыре собранных сценария, которые показывают новые фичи не по одной, а так, как они обычно используются вместе.
1. Offline-first мобильное приложение
Пользователь оформляет заказ в метро, где связь то есть, то нет. Раньше это была ошибка сети, которую нужно было ловить и обрабатывать руками. Теперь offlineQueue берёт это на себя, а бейдж синхронизации — обычное реактивное состояние Vue-компонента, обновляемое через onFlushSuccess/onFlushError:
vue
<script setup>
import { ref } from "vue";
import { createRestClient, OfflineQueuedError } from "rest-pipeline-js";
const pendingSync = ref(0);
const client = createRestClient({
baseURL: "https://api.example.com",
offlineQueue: {
enabled: true,
persistAdapter: {
save: (queue) => localStorage.setItem("offline-queue", JSON.stringify(queue)),
load: () => JSON.parse(localStorage.getItem("offline-queue") ?? "null"),
},
// isOnline/onOnlineChange не заданы — используется дефолт: navigator.onLine
// и браузерное событие "online" (для Capacitor/Cordova-обёртки подставьте
// сюда нативный детектор сети соответствующего плагина).
onFlushSuccess: () => pendingSync.value--,
onFlushError: (req) => toast.error(`Не удалось отправить: ${req.url}`),
},
});
async function placeOrder(cart) {
try {
return await client.post("/orders", cart);
} catch (err) {
if (err instanceof OfflineQueuedError) {
pendingSync.value++;
return { queued: true, queueId: err.queueId };
}
throw err;
}
}
</script>
<template>
<button @click="placeOrder(cart)">Оформить заказ</button>
<span v-if="pendingSync > 0" class="sync-badge">{{ pendingSync }} ждут синхронизации</span>
</template>
Каждый заказ в очереди несёт один и тот же Idempotency-Key на всех попытках реплея — если запрос на самом деле дошёл до сервера прямо перед обрывом связи, повторная отправка после реконнекта не создаст вторую запись. pendingSync — обычный ref, реактивность бейджа в шаблоне не требует ничего специфичного для offlineQueue, кроме того, чтобы обновлять его из уже знакомых коллбэков конфига.
2. Интеграция с жёстко лимитированным партнёрским API
Сервис за балансировщиком из нескольких инстансов ходит в API поставщика с лимитом 100 запросов в минуту. Без координации между инстансами каждый считает лимит независимо — три инстанса реально бьют по 300 запросов в минуту. Совмещаем распределённый rateLimit.store из 2.0.0 с новым проактивным троттлингом и «размазанным» джиттером — так, чтобы инстансы не только не превышали общий лимит, но и не ретраили залпом одновременно:
ts
const client = createRestClient({
baseURL: "https://partner-api.example.com",
rateLimit: {
store: redisRateLimiterStore,
key: "partner-api", // общий бакет на все инстансы
onRateLimitHeaders: (headers, control) => {
const remaining = Number(headers["x-ratelimit-remaining"]);
if (remaining <= 5) control.throttleFor(2000);
},
},
retry: {
attempts: 5,
delayMs: 300,
backoffMultiplier: 2,
jitterStrategy: "decorrelated", // не совпадают по времени с ретраями соседних инстансов
},
});
Проактивный троттлинг сглаживает нагрузку заранее, ещё до первого 429; decorrelated-джиттер не даёт залпу ретраев после реального сбоя ударить по партнёрскому API синхронно со всех инстансов сразу.
3. Потоковый импорт большого каталога
Партнёрский API отдаёт миллион товаров постранично. Загружать всё в память ради paginateAll() — не вариант; нужен построчный стриминг прямо в БД, с проверкой формата каждой страницы на лету:
ts
const importStage = {
key: "importCatalog",
stream: () =>
flattenPages(
paginate({
strategy: "offset",
limit: 200,
fetchPage: (offset, limit) =>
client.get("/catalog/items", { params: { offset, limit } }).then((r) => r.data),
}),
),
validateOutput: (item) => catalogItemSchema.parse(item), // ловит порчу схемы на лету
onChunk: (item) => db.upsert("catalog_items", item),
};
const orchestrator = createPipeline([importStage]);
await orchestrator.run();
onChunk вызывается на каждый отдельный товар, а не на страницу — flattenPages() разворачивает поток страниц в поток элементов. Если партнёр однажды пришлёт бракованную запись посреди импорта, validateOutput поймает её раньше, чем она долетит до БД.
4. Юнит-тесты retry-логики без реального бэкенда
Нужно проверить, что клиент действительно повторяет запрос после 503 и открывает circuit breaker после серии отказов — без похода в сеть и без флакующих таймаутов в CI:
ts
import { createRestClient } from "rest-pipeline-js";
import { createMockAdapter } from "rest-pipeline-js/testing";
it("retries on 503 and eventually opens the circuit", async () => {
const adapter = createMockAdapter([
{ method: "GET", url: "/orders", respond: { error: true, status: 503 } }, // всегда 503
]);
const client = createRestClient({
baseURL: "https://api.example.com",
adapter,
retry: { attempts: 2, delayMs: 1 },
circuitBreaker: { failureThreshold: 3, openMs: 60_000 },
});
await expect(client.get("/orders")).rejects.toThrow();
await expect(client.get("/orders")).rejects.toThrow();
await expect(client.get("/orders")).rejects.toThrow();
expect(await client.getCircuitBreakerState()).toBe("open");
expect(adapter.calls.length).toBeGreaterThan(3); // retry реально сработал на каждом вызове
});
adapter.calls даёт честную историю вызовов — можно проверить не только итоговый результат, но и то, что retry действительно попытался повторить запрос нужное число раз, а не просто упал с первой попытки.
Обратная совместимость
Все восемь новых возможностей — аддитивные поля конфига и новые точки входа, ничего не ломающие в существующем коде. Рефакторинг src/ и правки гонок не меняют публичный API и внешне наблюдаемое поведение (кроме, собственно, того, что гонки больше не происходят). Мигрировать с 2.0.0 на 2.1.0 можно без единой правки в вызывающем коде.
Релиз 2.1.0 уже на npm:
bash
npm i rest-pipeline-js@2.1.0
Репозиторий: github.com/macrulezru/pipeline-js
Демо: npm run demo:vue после клонирования — CI/CD Pipeline и Trading Terminal в браузере, без внешних зависимостей
npm: npmjs.com/package/rest-pipeline-js