Потоковые данные
Потоковые этапы (SSE / AsyncIterable)
Этап, функция stream которого возвращает AsyncIterable<T>. Оркестратор собирает все полученные чанки в массив (результат этапа). onChunk вызывается для каждого чанка в реальном времени.
const orchestrator = createPipeline([
{ key: 'auth', request: async () => getToken() },
{
key: 'liveData',
stream: async function* ({ prev }) {
const source = new EventSource(`/api/stream?token=${prev}`)
yield* eventSourceToAsyncIterable(source)
},
onChunk: (chunk, sharedData) => {
sharedData.partial = (sharedData.partial ?? '') + chunk
updateUI(sharedData.partial)
},
},
{
key: 'finalize',
// allResults.liveData.data — полный массив чанков
request: async ({ allResults }) => allResults.liveData.data.join(''),
},
])- Учитывает
abort()— проверяет сигнал прерывания между каждым чанком. - Поддерживает
continueOnError— провалившиеся потоковые этапы можно пропускать, как и любой другой шаг. - Генерирует стандартные события шага:
step:start,step:success,step:error.
Пагинация
paginate() обходит страницы постраничного API как AsyncGenerator<T[]>, скрывая разницу между курсорными и offset/limit-based API:
import { paginate, paginateAll, flattenPages } from 'rest-pipeline-js'
// На основе курсора (стратегия по умолчанию)
for await (const page of paginate({
fetchPage: (cursor) => client.get('/items', { params: { cursor } }).then((r) => r.data),
})) {
console.log(page.length, 'items')
}
// На основе offset/limit
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? } (offset-стратегия — останавливается, когда страница короче limit, либо когда offset достигает total, если API его сообщает).
paginateAll(options)— собирает все страницы в один плоский массив; самый простой вариант, когда весь набор данных достаточно мал, чтобы поместиться в памяти целиком.flattenPages(pages)— превращает поток страниц в поток отдельных элементов; полезно как источникStreamStageConfig.stream, когдаonChunkдолжен срабатывать на каждый элемент, а не на каждую страницу (см.examples/pagination-stream.ts).- И
paginate(), и её коллбэкfetchPageпринимают опциональныйsignalдля поддержкиabort().
Если вы заранее знаете количество страниц и хотите получать их параллельно, а не последовательно, см. вместо этого examples/pagination-fanout.ts.
Этапы WebSocket
Этап, работающий поверх постоянного WebSocket-соединения вместо одиночного запроса/ответа — чаты/фиды присутствия, живые биржевые стаканы, события совместного редактирования и т. д. Сообщения, возвращённые onMessage, собираются в массив data этапа (тот же паттерн, что и чанки потоковых этапов); onChunk срабатывает на каждое сообщение в реальном времени.
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()url— строка или функция от{ prev, allResults, sharedData, signal }, те же параметры, что получаетrequest.createWebSocket— фабрика для базовой реализации. По умолчаниюglobalThis.WebSocket(браузеры, Deno, Node ≥22). Для Node <22 передайте реализацию на базе пакетаws:createWebSocket: (url, protocols) => new WS(url, protocols).onMessage(обязателен) — вызывается на каждое сообщение сevent.data; может бытьasync. Значение, отличное отundefined, собирается в массив результатов этапа и передаётся вonChunk.- Успех/ошибка определяются событием закрытия, а не событием ошибки: большинство реализаций WebSocket генерируют
errorнепосредственно передclose, так что сам по себеonErrorне проваливает этап — чистое закрытие (wasClean: true) успешно завершает этап со всем, что было собрано к тому моменту; нечистое закрытие отклоняет его, проходя черезcontinueOnError, как и любой другой этап. closeOn(data)— вернитеtrue, чтобы закрыть соединение и успешно завершить этап, как только вы увидели нужное, вместо ожидания закрытия сервером.timeoutMs— общий таймаут соединения (не сбрасывается сообщениями); закрывает сокет и проваливает этап, если срабатывает до того, как произойдёт чистое закрытие само по себе.- Учитывает
abort()— закрывает базовое соединение и отклоняет этап.
Полную аннотированную версию см. в examples/websocket-stage.ts.