Runtime и Schedule: retry, repeat, timeout, race с fallback
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
Runtime и Schedule: retry, repeat, timeout, race с fallback
Сцена · расписание дворника
У дворника есть план: подметать каждые два часа, после дождя добавить ещё один обход, в снегопад переходить на лопату. План это не сам дворник: он не подметает, он отвечает на один вопрос: “сейчас работать, и через сколько следующий заход”. Дворник, это рантайм. Он берёт план, делает шаг, ждёт, делает следующий, и так пока кто-то не скажет “смена закончилась”.
В Effect эту пару зовут Schedule и ManagedRuntime. Schedule описывает план, рантайм его исполняет. У Schedule богатый набор комбинаторов: экспонента, фибоначчи, fixed, spaced, min, max, jittered. У рантайма главное умение, переиспользоваться: один граф зависимостей, одна точка инициализации, несколько точек входа.
В этом уроке мы собираем всю обвязку Pulse, без которой он остался бы линейным скриптом: расписание для каждого монитора, ретраи с экспонентой и потолком, таймаут на запрос, race с резервным URL, общий рантайм для CLI и для HTTP-сервера. К концу твой Pulse умеет жить долго, переживать кратковременные ошибки и корректно сворачиваться на SIGTERM.
Карта урока · что заберёшь домой
Семь разделов.
- Раздел 1, рантайм подробно: что внутри, какие три метода запуска, и зачем делить один рантайм между точками входа.
- Раздел 2,
Schedule<Out, In, Error, Env>как машина решений. Базовые расписания:recurs,spaced,fixed,windowed,exponential,fibonacci,forever. - Раздел 3, композиция:
min,max,addDelay,jittered,upTo,modifyDelay. Каноничный рецепт “экспонента с потолком и шумом” собирается из трёх вызовов. - Раздел 4,
Effect.retryпротивEffect.repeat. Почему оба часто нужны вместе. Выборочный retry по тегам через опциюwhile. - Раздел 5,
Effect.timeout,Effect.race,Effect.raceAll. Fallback на резервный URL, fallback на cached value, hedged-запрос с задержкой. - Раздел 6, календарные и сглаживающие расписания:
Schedule.cron, плюсthrottleиdebounceна стримах. - Раздел 7, Pulse: per-target Schedule из конфига, политика retry на NetworkError, таймаут 5 секунд, race с резервным URL, repeat по интервалу. Один
Runtimeдля CLI и для HTTP-сервера.
К концу урока ты собираешь Schedule из трёх кусков и читаешь его как одно расписание; отличаешь retry и repeat, и используешь их в одной программе; применяешь Effect.timeout без Promise.race(setTimeout); разделяешь один ManagedRuntime между двумя точками входа без дублирования зависимостей.
Раздел 1 · Рантайм подробно
Что лежит внутри рантайма
ManagedRuntime<R, ER> это не одна вещь, а контейнер:
type ManagedRuntime<R, ER> = {
context: () => Promise<Context.Context<R>>; // граф зависимостей (Layer построил)
memoMap: Layer.MemoMap; // мемоизация слоёв
scope: Scope.Closeable; // общий скоуп всего графа
// плюс runPromise / runFork / runSync / dispose
};
contextэто то, что положилLayer. Когда твой код пишетyield* HttpService, рантайм достаёт реализацию отсюда.memoMapследит за тем, чтобы один и тот же Layer, встретившийся в графе дважды, построился один раз (см. 04 · Services и Layer).scopeэто тот самый скоуп, который закроется наdispose()и погасит все finalizer-ы.
Отдельного типа Runtime<R>, который можно было собрать руками из контекста и флагов, в Effect 4 нет. Рантайм теперь ровно один вид, и строится он из Layer-а: ManagedRuntime.make(MainLive). Модуль Runtime в библиотеке остался, но занимается он другим: Runtime.makeRunMain и логика корректного завершения процесса.
Три точки входа
Запустить эффект можно тремя способами:
import { Effect } from 'effect';
// 1. Promise-based, асинхронное ожидание результата
await Effect.runPromise(Effect.succeed(42));
// 2. Fiber-based, не ждём результата, получаем хэндл fiber-а
const fiber = Effect.runFork(Effect.sleep('5 seconds'));
// можно потом yield* Fiber.interrupt(fiber)
// 3. Sync-only, бросает, если внутри есть async-граница
Effect.runSync(Effect.sync(() => 42));
Все три требуют, чтобы R был never: у голого запуска сервисов взять неоткуда. Если в R сидит что-то ещё, TypeScript остановит:
declare const probe: Effect.Effect<Response, NetworkError, HttpService>;
Effect.runFork(probe);
// ^^^^^^^ ошибка: HttpService не покрыт
Запуск с контекстом
Дальше два пути. Первый, если контекст у тебя уже есть готовым значением, это *With-варианты:
import { Context, Effect } from 'effect';
const services = Context.make(HttpService, httpServiceImpl);
Effect.runForkWith(services)(probe); // ok, R покрыт
Второй, и в реальном коде основной, это ManagedRuntime.make(MainLive). Он собирает контекст из Layer-а сам, с честной мемоизацией и с dispose(), и даёт те же три метода прямо на себе: runtime.runPromise, runtime.runFork, runtime.runSync.
Почему один рантайм на программу
Сцена из урока 04 повторяется: каждый Effect.runPromise(eff.pipe(Effect.provide(MainLive))) строит граф заново. Под нагрузкой это новый HttpService, новый Storage, новый file handle. CLI-команда переживает (она и так одна), HTTP-сервер падает.
Pulse в финальном виде имеет одну точку сборки и две точки входа:
┌──────────────────────────────────────────┐
│ MainLive Layer │
│ HttpService · Storage · Logger · Clock │
│ MonitorEvents · SlaState (TxRef) │
│ Probe · Schedule · BatchedDnsResolver │
└────────┬─────────────────────────┬───────┘
│ │
┌──────────┴──────────┐ ┌─────────┴──────────┐
│ pulse watch │ │ embedded HTTP │
│ (Terminal) │ │ (node:http) │
└─────────────────────┘ └────────────────────┘
▲
│
один ManagedRuntime,
один MemoMap,
один Storage handle
Внутри обеих команд работают одни и те же Probe, Schedule, MonitorEvents. Различается то, как они отдают результат наружу: одна рисует ANSI-таблицу в терминал, вторая отдаёт JSON по HTTP. С точки зрения Effect это одна программа с двумя выходами. Второй выход мы сейчас соберём руками на node:http, чтобы не тащить вперёд HTTP-модули Effect. В уроке 12 этот же выход превратится в полноценный pulse serve на NodeHttpServer, рантайм при этом не поменяется.
Шаг 1 · runtime.ts, единственный экземпляр
// pulse-<nick>/src/runtime.ts
import { ManagedRuntime } from 'effect';
import { MainLive } from './main.ts';
export const runtime = ManagedRuntime.make(MainLive);
ManagedRuntime.make(layer) лениво строит граф на первом runPromise. Если хочешь принудительно прогреть его при запуске процесса, вызови:
await runtime.runPromise(Effect.void);
Это полезно для серверов, где первый HTTP-запрос не должен платить за подключение к базе, открытие файла, прогрев DNS-кеша.
Шаг 2 · две точки входа
// pulse-<nick>/src/cli/watch.ts
import { runtime } from '../runtime.ts';
import { watchProgram } from './watch-program.ts';
await runtime.runPromise(watchProgram);
// pulse-<nick>/src/main.ts
import http from 'node:http';
import { runtime } from './runtime.ts';
import { statusProgram } from './http/status.ts'; // обычный Effect, отдаёт JSON-снимок
// embedded HTTP-вход: каждый запрос это просто ещё одна программа на том же runtime
const server = http.createServer((_req, res) => {
runtime.runPromise(statusProgram).then((snapshot) => {
res.setHeader('content-type', 'application/json');
res.end(JSON.stringify(snapshot));
});
});
server.listen(8080);
Оба пути используют тот же объект runtime. Граф MainLive собирается ровно один раз, Storage открывает свой events.jsonl ровно один раз, MonitorEvents живёт ровно в одном экземпляре. Чем больше команд и поверхностей появляется в проекте, тем меньше у тебя случайных дубликатов. Внутри обработчика каждый запрос остаётся Effect-программой на общем рантайме, никакого Effect.provide(MainLive) на запрос. Тут это голый node:http, без HTTP-модулей Effect. В уроке 12 заменим ручной createServer на NodeHttpServer, а runtime оставим как есть.
Шаг 3 · dispose() на SIGTERM
В 05 · Resources и Scope и 06 · Файберы и concurrency мы говорили про корректное завершение. Кратко: у ManagedRuntime есть метод dispose(), который закрывает Scope, держащий весь Layer. Все finalizer-ы (file handle, fork-нутые фоновые файберы) отрабатывают, и программа выходит чисто.
import process from 'node:process';
import { runtime } from './runtime.ts';
const shutdown = async () => {
await runtime.dispose();
process.exit(0);
};
process.on('SIGINT', shutdown);
process.on('SIGTERM', shutdown);
В CLI на effect/unstable/cli (урок 12) эта обвязка пишется иначе, через NodeRuntime.runMain, который сам слушает сигналы и зовёт dispose. До тех пор, пока ты пишешь маленькие точки входа руками, помни про эти три строчки.
Что взять с собой
- Рантайм это
(Context, MemoMap, Scope)плюс методы запуска. Собирается черезManagedRuntime.make(layer); отдельного собираемого руками типаRuntime<R>в v4 нет. - Три способа запуска:
runPromise(Promise),runFork(Файбер, без ожидания),runSync(бросает на async-границе). - Без своего рантайма
Effect.runPromise(eff.pipe(Effect.provide(MainLive)))каждый раз пересобирает граф. СManagedRuntimeграф один на программу. - Один
runtimeобслуживает все точки входа (CLI и HTTP-сервер). Это и есть “одна программа с двумя поверхностями”. runtime.dispose()на остановке обязательно, иначе scoped-ресурсы не закроются.
Раздел 2 · Schedule как машина решений
Идея словами
Schedule<Out, In, Error, Env> это машина, которая на каждом тике говорит одно из двух: “продолжаем, подожди столько-то и повтори” или “хватит”. На входе она видит результат предыдущего шага (In), на выходе отдаёт что-то полезное в лог или в каскад (Out), на каждом шаге внутри посчитала задержку.
Чтобы машина начала тикать, её надо поднять в Effect.retry, Effect.repeat, Effect.schedule, Stream.schedule или похожий комбинатор. Сам по себе Schedule это рецепт, не процесс.
Сигнатура:
type Schedule<Out, In = unknown, Error = never, Env = never> = ...
Out, что машина возвращает на каждом тике. Часто это число итераций или сама задержка.In, что машина видит на входе. УEffect.retryэтоE(тип ошибки), уEffect.repeatэтоA(тип успеха). ЧерезInпишутся “ретраить только сетевые ошибки”, “repeat-ить, пока результат меньше 10”.Error, ошибка самого расписания. Появляется, если ты повесил на него эффектныйtapилиmodifyDelay, который умеет падать.Env, зависимости расписания. Почти всегдаnever, но если вScheduleвлез какой-нибудь сервис, он попадёт сюда.
Четвёртый параметр это то, чем сигнатура отличается от привычной по старым материалам: раньше их было три, без канала ошибки.
Базовые расписания
Все они импортируются из effect/Schedule:
import { Schedule } from 'effect';
// 1. recurs(N), ровно N повторов и стоп
const fiveTimes = Schedule.recurs(5);
// Out = number (0, 1, 2, 3, 4), In = unknown, R = never
// 2. forever, бесконечно
const endless = Schedule.forever;
// 3. spaced(d), пауза d между завершением предыдущего и стартом следующего
const everyMinute = Schedule.spaced('1 minute');
// 4. fixed(d), пауза d между стартами (то есть, если работа заняла 30s, ждём ещё 30)
const everyMinuteFromStart = Schedule.fixed('1 minute');
// 5. windowed(d), тики выравниваются по границам окон длины d
const alignedToMinute = Schedule.windowed('1 minute');
// 6. exponential(start, factor), задержки 100, 200, 400, 800...
const exp = Schedule.exponential('100 millis', 2.0);
// 7. fibonacci(start), задержки start, start, 2*start, 3*start, 5*start...
const fib = Schedule.fibonacci('100 millis');
Из этих семи штук собирается почти всё, что встречается в продакшене. Линейного расписания (start, 2*start, 3*start) отдельного кирпича нет: если нужно именно оно, собирается через Schedule.modifyDelay поверх recurs.
Шаг 1 · spaced против fixed
Эта пара путается чаще всего. Картина различия:
spaced('30 seconds'), пауза между завершениями
┌─────┐ ┌────┐ ┌──────┐
│ 5s │═════════════════════│ 4s │══════════════════│ 10s │═════
└─────┘ 30 seconds └────┘ 30 seconds └──────┘ 30s
старт-1 завершилось, старт-2 завершилось, старт-3
ждём 30s ждём 30s
fixed('30 seconds'), пауза между стартами
┌─────┐ ┌────┐ ┌─────
│ 5s │══════════════════════════│ 4s │══════════════════│ 10s
└─────┘ ждём до 30s от старта-1 └────┘ до 30s от старта-2
старт-1 30s старт-2 30s старт-3
spaced это “после того, как закончил, посиди 30s и снова возьмись”. fixed это “запускай в начале каждой 30-секундной клетки”. Если работа занимает 35 секунд, fixed('30 seconds') стартует следующую сразу же, как закончилась предыдущая (он не успевает в свою клетку, и подтягивается к следующей). spaced всегда даёт ровно 30 секунд тишины между запусками.
Для пробинга URL правильнее fixed: тебе нужна частота “раз в минуту”, независимо от того, сколько отвечал сервер. Для рассылки писем “не присылать чаще, чем раз в день” удобнее spaced.
Шаг 2 · exponential и jittered
exponential('100 millis', 2.0) отдаёт 100, 200, 400, 800, 1600, 3200, … миллисекунд. Без потолка это быстро улетит в минуты, поэтому обычно накрывают его вторым расписанием через Schedule.min (см. раздел 3).
import { Schedule } from 'effect';
const backoff = Schedule.exponential('100 millis', 2.0);
// задержки: 100ms, 200ms, 400ms, 800ms, 1.6s, 3.2s, ...
Голый exponential плох тем, что десять файберов после одновременного провала будут ретраиться одновременно. Их новый “кулак” запросов прилетит ровно через 100ms, потом 200, и каждая волна синхронна. Это thundering herd, и сервер ложится повторно.
Schedule.jittered ломает синхронность: каждая задержка умножается на случайный коэффициент, по умолчанию в диапазоне от 0.8 до 1.2, то есть от 80 до 120 процентов исходного интервала:
const backoffJittered = Schedule.exponential('100 millis', 2.0).pipe(Schedule.jittered);
// 100ms превращается в [80ms..120ms] (по умолчанию), 200 в [160..240] и так далее
Нужен другой разброс, бери Schedule.modifyDelay и считай задержку сам. По умолчанию jittered использует Random из стандартных сервисов, в тестах его подменяют своим слоем (04 · Default services).
Шаг 3 · Schedule.tap, лог каждого тика
Когда отлаживаешь свою политику, удобно видеть, что и когда машина решила:
import { Effect, Schedule } from 'effect';
const noisy = Schedule.exponential('100 millis').pipe(
Schedule.tap((meta) => Effect.log(`schedule tick, output=${String(meta.output)}`)),
);
tap это маленький “вынюхиватель”: он ничего не меняет, просто на каждом тике дёргает твой эффект. Внутрь приходит не голый выход, а объект-метаданные шага: там и output, и input, и номер повтора. В проде обычно убирают, в разработке оставляют, пока расписание не уляжется в голове.
Что взять с собой
Schedule<Out, In, Error, Env>это машина решений: “продолжаем, и через сколько”, или “хватит”.- Сам по себе ничего не запускает, нужен
Effect.retry/repeat/scheduleилиStream.schedule. - Базовых элементов мало:
recurs,forever,spaced,fixed,windowed,exponential,fibonacci. Из них собирается всё. spaced(d)это пауза между завершениями,fixed(d)это пауза между стартами. Для регулярного опроса бериfixed.- Сырая
exponentialбезjitteredэто рецепт thundering herd.
Раздел 3 · Композиция Schedule
Зачем композировать
Расписание из одного куска редко бывает достаточным. Реалистичное правило звучит так: “ретраить с экспонентой, но не больше пяти раз, с потолком на 30 секунд, и с jitter”. Это четыре свойства. По одному комбинатору на каждое.
В Effect композиторы Schedule спроектированы как алгебра: пять комбинаторов, всё остальное это их сочетания. Освоить эту пятёрку, и расписание становится таким же лёгким, как массив через map/filter/reduce.
Главные комбинаторы
import { Schedule } from 'effect';
// 1. min, взять меньшую задержку; живёт, пока живо хотя бы одно плечо
Schedule.min([Schedule.exponential('100 millis'), Schedule.spaced('30 seconds')]);
// экспонента, но не медленнее 30 секунд между попытками
// 2. max, взять большую задержку; живёт, пока живы оба плеча
Schedule.max([Schedule.exponential('100 millis'), Schedule.spaced('1 second')]);
// экспонента, но не быстрее секунды между попытками
// 3. upTo, ограничить расписание по времени или по числу повторов
Schedule.exponential('100 millis').pipe(Schedule.upTo({ duration: '30 seconds' }));
Schedule.exponential('100 millis').pipe(Schedule.upTo({ times: 5 }));
// 4. addDelay, прибавить задержку к расчётной
Schedule.exponential('100 millis').pipe(Schedule.addDelay(() => Effect.succeed('50 millis')));
// 5. jittered, рандомизировать задержку
Schedule.exponential('100 millis').pipe(Schedule.jittered);
Запомнить вот что: min и max берут список расписаний и складывают их решения. min это “или, и побыстрее”, max это “и, и помедленнее”. upTo/addDelay/jittered это унарные модификаторы.
Имена тут говорящие, и это заметное улучшение против того, что было раньше. Прошлые версии Effect предлагали intersect и union (плюс синоним either), и по имени было невозможно вспомнить, какой из них берёт минимум задержки, а какой максимум. min и max называют ровно то, что делают.
Шаг 1 · “пять попыток с экспонентой”
Каноничный сетевой retry это сборка из двух кусков:
import { Schedule } from 'effect';
const networkRetry = Schedule.exponential('200 millis', 2.0).pipe(
Schedule.jittered, // ломаем синхронность
Schedule.upTo({ times: 5 }), // максимум 5 ретраев
);
// итог: задержки [~100..200, ~200..400, ~400..800, ~800..1600, ~1600..3200] ms,
// после пятого провала останавливается
Читай как фразу: “экспонента с jitter, не больше пяти повторов”. Раньше ограничение по счёту приходилось писать как пересечение с recurs(5), и такое пересечение заодно меняло тип выхода расписания на пару. upTo({ times }) ограничивает только счётчик, не трогая ни задержку, ни Out.
Шаг 2 · потолок на задержку
exponential уходит в космос за несколько шагов. Часто хочется сказать “не больше 30 секунд между попытками”:
const cappedExp = Schedule.min([
Schedule.exponential('200 millis'),
Schedule.spaced('30 seconds'),
]);
min на каждом тике берёт минимальную задержку из плеч. Пока экспонента ниже 30 секунд, побеждает она. Как только переползла, начинает побеждать spaced('30 seconds'). На графике это типичная кривая “сначала растёт, потом упирается в полку”.
Альтернативный способ записать потолок:
const cappedExp = Schedule.exponential('200 millis').pipe(
Schedule.modifyDelay(({ output }) =>
Effect.succeed(Duration.toMillis(output) > 30_000 ? '30 seconds' : output),
),
);
Длиннее, но если потолок не круглый или зависит от обстоятельств, иногда удобнее. Обрати внимание, что modifyDelay теперь принимает метаданные шага и возвращает Effect: задержку можно считать асинхронно и с доступом к сервисам.
Шаг 3 · “пять попыток, экспонента, jitter, потолок 30 секунд”
Собираем всё разом:
import { Duration, Schedule } from 'effect';
const sturdyRetry = Schedule.min([
Schedule.exponential('200 millis', 2.0),
Schedule.spaced(Duration.seconds(30)), // потолок
]).pipe(
Schedule.jittered, // шум
Schedule.upTo({ times: 5 }), // максимум 5 ретраев
);
Это и есть тот самый “production retry policy”, который ты будешь копировать из проекта в проект. У него четыре свойства, и каждое читается по одному комбинатору.
Шаг 4 · условный retry через опцию while
Часть ошибок ретраить не нужно. HTTP 404 это не временный сбой, повторять смысла нет. Сетевая ошибка или таймаут это типичный кандидат на ретрай. Условие живёт не на расписании, а на самом Effect.retry:
import { Effect } from 'effect';
import type { NetworkError, TimeoutError, HttpStatusError } from './errors.ts';
type ProbeError = NetworkError | TimeoutError | HttpStatusError;
const retryTransient = <A, R>(effect: Effect.Effect<A, ProbeError, R>) =>
effect.pipe(
Effect.retry({
schedule: sturdyRetry,
while: (error) =>
error._tag === 'NetworkError' ||
error._tag === 'TimeoutError' ||
(error._tag === 'HttpStatusError' && error.status >= 500),
}),
);
while сужает политику: продолжаем, пока ошибка попадает под условие. На первом 4xx или нетипичном теге машина останавливается, и Effect.retry отдаёт ошибку наружу. В том же объекте живут until (зеркальное условие) и times (ограничение по счёту без отдельного расписания), так что для простых случаев schedule можно вообще не писать:
effect.pipe(Effect.retry({ times: 3, while: isTransient }));
Раньше это условие вешалось на само расписание через Schedule.whileInput. Теперь его там нет, и это к лучшему: политика “что ретраить” и политика “как долго ждать” перестали смешиваться в одном значении.
Шаг 5 · Schedule.map и Schedule.passthrough
Иногда хочется превратить Out (счётчик повторов) в более полезную штуку. Например, в саму задержку или в строку для лога:
import { Schedule } from 'effect';
const explained = Schedule.exponential('100 millis').pipe(
Schedule.map((delay) => `retry in ${String(delay)}`),
);
// Out теперь string
Schedule.passthrough пускает In наружу как Out. Полезно, когда хочется в каскаде tap увидеть саму ошибку:
const debugRetry = Schedule.recurs(3).pipe(
Schedule.passthrough,
Schedule.tap((meta) => Effect.logWarning(`retrying after: ${String(meta.output)}`)),
);
Schedule.passthrough это маленький трюк, не первая пятёрка. Но когда впервые понадобится посмотреть, на какой именно ошибке сработал retry, ты будешь рад его знать.
Что взять с собой
- Пять главных комбинаторов:
min,max,upTo,addDelay,jittered. Плюс модификаторыmap,passthrough,modifyDelay,tap. - “Production retry” это три строки:
min([экспонента, потолок]),jittered,upTo({ times }). - Выборочный retry по типу ошибки живёт в опциях
Effect.retry, а не на расписании. Не ретраишь 4xx, ретраишь NetworkError и 5xx. - Сам
Scheduleничего не запускает. Это только описание, как вLayer.
Раздел 4 · Effect.retry против Effect.repeat
Идея словами
Effect.retry(eff, schedule), политика на падение. Покаeffотдаёт ошибку, рантайм ждёт по расписанию и повторяет. На первом успехе возвращаетA. Если расписание истощилось, отдаёт последнюю ошибку наружу.Effect.repeat(eff, schedule), политика на успех. ПокаeffотдаётA, рантайм ждёт и повторяет. На первой ошибке отдаёт её наружу. Расписание определяет, когда остановить успешное повторение (например,Schedule.recurs(10), “повтори 10 раз и хватит”).
В реальном коде эти два часто комбинируются: каждый отдельный запрос имеет retry на NetworkError, а вся серия запросов repeat-ится по расписанию “раз в минуту”.
const probeRound = probe(url).pipe(
Effect.retry(networkRetry), // выживание одного раунда
);
const probeForever = probeRound.pipe(
Effect.repeat(Schedule.fixed('1 minute')), // регулярность
);
Прочти про себя: “пробу делаем с ретраями, всю пробу повторяем каждую минуту”. Это типичная пара “защита и регулярность”.
Шаг 1 · Effect.retry детально
import { Effect, Schedule } from 'effect';
import { probe } from './probe.ts';
const safeProbe = probe('https://github.com').pipe(
Effect.retry(
Schedule.exponential('200 millis').pipe(
Schedule.jittered,
Schedule.upTo({ times: 5 }),
),
),
);
// safeProbe: Effect<Response, ProbeError, never>
// тип ошибки тот же: после 5 ретраев последняя ошибка прилетает наружу
Что важно. retry не меняет тип ошибки, только меняет вероятность увидеть её. Если первый запрос дал NetworkError, дальше идёт пауза, второй запрос, третий. Когда расписание кончилось, последняя NetworkError пробрасывается.
Тот же retry с выборочным условием:
const safeProbeTagged = probe('https://github.com').pipe(
Effect.retry({
schedule: Schedule.exponential('200 millis').pipe(Schedule.upTo({ times: 5 })),
while: (error) =>
error._tag === 'NetworkError' ||
error._tag === 'TimeoutError' ||
(error._tag === 'HttpStatusError' && error.status >= 500),
}),
);
В объектной форме доступны четыре поля: schedule, while, until, times. Любое можно опустить, так что “три раза без пауз” это просто Effect.retry({ times: 3 }).
Шаг 2 · Effect.retryOrElse, fallback после исчерпания
Если все ретраи провалились, иногда хочется не упасть с ошибкой, а вернуть значение по умолчанию или сделать что-то ещё. Effect.retryOrElse для этого:
import { Effect, Schedule } from 'effect';
const probeWithFallback = probe('https://github.com').pipe(
Effect.retryOrElse(
Schedule.recurs(3),
(error, attempt) =>
Effect.gen(function* () {
yield* Effect.logError(`probe failed after ${attempt} attempts: ${String(error)}`);
return new Response(null, { status: 503 });
}),
),
);
// probeWithFallback: Effect<Response, never, never>
// тип ошибки превратился в never, потому что fallback всегда возвращает Response
retryOrElse(schedule, onFail), где onFail получает последнюю ошибку и номер попытки. Возвращаемое значение onFail-а заменяет результат всей операции.
Шаг 3 · Effect.repeat детально
Тот же самый трюк, но на успехе:
import { Effect, Schedule } from 'effect';
const tenTimes = Effect.gen(function* () {
yield* Effect.log('tick');
return Date.now();
}).pipe(Effect.repeat(Schedule.recurs(10)));
// tenTimes: Effect<number, never, never>
// каждые 10 итераций возвращает последнее значение из эффекта
Schedule.recurs(10) это “повторить десять раз”, значит всего эффект выполнится 11 раз (один первый плюс 10 повторов). Это та же самая семантика, что в 09 · Stream: recurs(N) это “ещё N тиков после первого”.
Регулярные пробы:
const probeForever = probe('https://github.com').pipe(
Effect.repeat(Schedule.fixed('1 minute')),
);
// probeForever: Effect<never, ProbeError, never>
// тип A это never, потому что цикл не завершается сам по себе (fixed работает forever).
// прервать его можно через Fiber.interrupt снаружи.
fixed не имеет встроенного “стоп”, цикл бесконечен. Завершение приходит снаружи через Fiber.interrupt или runtime.dispose(). На уровне типов это видно по A = never: программа ничего не вернёт штатно.
Шаг 4 · repeat плюс retry
Каноничный паттерн “регулярно стучи в URL, при ошибке ретраить”:
import { Effect, Schedule } from 'effect';
const networkRetry = Schedule.exponential('200 millis').pipe(
Schedule.jittered,
Schedule.upTo({ times: 5 }),
);
const probeOnce = probe('https://github.com').pipe(Effect.retry(networkRetry));
const probeForever = probeOnce.pipe(Effect.repeat(Schedule.fixed('1 minute')));
Порядок важен. retry внутри, repeat снаружи. Тогда каждый отдельный заход переживает кратковременные сбои, а внешний цикл задаёт регулярность.
Если перевернёшь, Effect.repeat(probe).pipe(Effect.retry(...)), семантика становится другой и редко полезной: вся бесконечная серия ретраится как один эффект.
Шаг 5 · короткая форма без Schedule
Иногда хочется маленький частный случай, без сборки полной Schedule. У Effect.repeat ровно такая же объектная форма, как у retry:
import { Effect } from 'effect';
// просто N раз
Effect.repeat(action, { times: 4 }); // выполнит action 5 раз (один плюс 4 повтора)
// пока не выполнится условие на результате
Effect.repeat(askForToken, { until: (token) => token.length > 0 });
Отдельных функций repeatN и repeatUntil в v4 нет: они схлопнулись в поля times и until этого объекта. Удобно для маленьких циклов опроса, где не нужны экспоненты и потолки.
Что взять с собой
Effect.retryэто политика на падение. Тип ошибки не меняется, последний провал пробрасывается.Effect.repeatэто политика на успех. На первой ошибке цикл прерывается.- В паре они работают вместе:
probe.pipe(Effect.retry(rNet), Effect.repeat(rTime)). Внутри защита, снаружи регулярность. retryOrElseспасает от пробрасывания ошибки: после исчерпания расписания подставляет дефолт.- У обоих есть объектная форма с полями
schedule,while,until,times, и для простых циклов её хватает безSchedule.
Раздел 5 · timeout, race, fallback
Effect.timeout
Сетевой запрос без таймаута это бомба замедленного действия. Соединение зависло, файбер висит, ресурс не освобождается. Effect.timeout решает это типизированно:
import { Effect } from 'effect';
const probe5s = probe('https://github.com').pipe(Effect.timeout('5 seconds'));
// probe5s: Effect<Response, ProbeError | TimeoutError, never>
В E появляется TimeoutError (стандартный класс из Effect, живёт в модуле Cause). На границе ты ловишь его через catchTag('TimeoutError', ...) или Effect.catch. Если хочешь свою ошибку или вообще другое поведение вместо падения, есть Effect.timeoutOrElse:
import { Effect } from 'effect';
import { TimeoutError } from './errors.ts';
const probe5sTyped = probe('https://github.com').pipe(
Effect.timeoutOrElse({
duration: '5 seconds',
orElse: () => Effect.fail(new TimeoutError({ url: 'https://github.com', deadline: 5000 })),
}),
);
// E теперь ProbeError | TimeoutError, со своим тегом
Обрати внимание на форму orElse: это не “какую ошибку подставить”, а “какой эффект выполнить вместо просроченного”. Поэтому подставить туда можно что угодно, не только падение: значение из кеша, заглушку, поход к другому источнику. Именно поэтому старая пара timeoutFail плюс timeoutFailCause схлопнулась в один комбинатор.
В Pulse у нас уже есть свой TimeoutError (см. 03 · Errors), и ловится он по тегу, наравне с остальными ошибками домена. Это аккуратнее для нашей системы ошибок.
Альтернативный сценарий, “если не успели, верни Option<None>”:
import { Effect } from 'effect';
const probeOption = probe('https://github.com').pipe(Effect.timeoutOption('5 seconds'));
// probeOption: Effect<Option<Response>, ProbeError, never>
// None означает "не успели за 5 секунд", без ошибки
timeoutOption уносит таймаут из канала ошибок в Option. Подходит, когда таймаут это нормальный исход, а не сбой.
Итого три варианта на выбор, и все три названы по своему поведению: timeout падает встроенной ошибкой, timeoutOption отдаёт Option, timeoutOrElse выполняет твой запасной эффект.
Как работает таймаут внутри
В отличие от Promise.race(action, sleep(d)) (где action продолжит фоновую работу после “проигрыша”), Effect.timeout прерывает внутренний файбер. На прерывании срабатывает finalizer scoped-ресурса (например, AbortController.abort() для fetch), и реальный сетевой запрос отменяется. Это та самая разница, которую мы обсуждали в 06 · Файберы и concurrency.
В Pulse это значит: если запрос к github.com висит на 30 секунд, через 5 секунд Effect.timeout прерывает файбер, fetch получает abort, сокет закрывается, и Pulse не копит висящие соединения.
Effect.race, первый успешный
Если у твоего URL есть резервный адрес (CDN-зеркало, второй регион), естественная политика, “запрашивай оба, бери первый успешный ответ”:
import { Effect } from 'effect';
const fastest = Effect.race(
probe('https://primary.example.com'),
probe('https://fallback.example.com'),
);
// fastest: Effect<Response, ProbeError, never>
// первый успешный побеждает, второй прерывается
Семантика: оба эффекта стартуют одновременно. Тот, кто первым отдал успех, побеждает, проигравший прерывается. Если оба упали, итог это последняя ошибка (детали см. в Effect.raceAll).
Маленькая тонкость, “первый успешный” против “первый ответивший”. Effect.race это “first success wins”. Если победитель провалился, рантайм ждёт результата от проигравшего. Это удобно: ты не получаешь ошибку, пока есть хоть один шанс на успех.
Effect.raceAll, N мирроров
То же самое, но для произвольного числа эффектов:
import { Effect } from 'effect';
const fastestOfMany = Effect.raceAll([
probe('https://us-east.api.example.com'),
probe('https://us-west.api.example.com'),
probe('https://eu-central.api.example.com'),
]);
Полезно для CDN-проб, multiroute-запросов, “проверить, что хоть один регион жив”.
Hedged-запрос, race с задержкой
Чистый race запускает второй запрос немедленно. Это в два раза больше трафика на каждый probe. Часто хочется: “стартуй основной, если за 200ms не успел, запусти второй параллельно”:
import { Effect } from 'effect';
const hedged = Effect.race(
probe('https://primary.example.com'),
probe('https://fallback.example.com').pipe(Effect.delay('200 millis')),
);
Это hedged-запрос: трафик удваивается только тогда, когда основной запрос реально медленный. На p50 это бесплатно, на p99 спасает.
В Pulse это укладывается без швов: первичный и резервный URL уже лежат в конфиге монитора, добавить Effect.delay('200 millis') к резервному стоит одной строки.
Effect.catch, fallback на cached value
Иногда правильная стратегия после провала, не race и не retry, а вернуть последний известный ответ из кеша. Отдельного orElse в v4 нет, его роль играет Effect.catch, “если упал, сделай вот это”:
import { Cache, Effect, Option } from 'effect';
import { ProbeCache } from './cache.ts';
const probeWithCache = (url: string) =>
probe(url).pipe(
Effect.catch(() =>
Effect.gen(function* () {
const cache = yield* ProbeCache;
const last = yield* Cache.getOption(cache, url);
if (Option.isNone(last)) return yield* Effect.fail(new NoCacheError({ url }));
return last.value;
}),
),
);
catch ловит любую ожидаемую ошибку основного эффекта и подставляет результат фолбэка. Тип ошибки изменится: вместо ProbeError будет то, чем может упасть сам фолбэк.
Имя catch тут не случайно короткое: это основной ловец ошибок в v4, и рядом с ним по одному принципу названы catchCause, catchDefect, catchTag, catchFilter. Если в старом коде видишь Effect.catchAll, это ровно он (03 · Errors).
Effect.orElseSucceed(value) это короткий вариант для константного фолбэка:
const probeOrEmpty = probe('https://x').pipe(
Effect.orElseSucceed(() => new Response(null, { status: 503 })),
);
// тип E теперь never
Чек “когда timeout, когда race, когда orElse”
| Сценарий | Инструмент |
|---|---|
| Запрос может зависнуть, нужен жёсткий потолок по времени | Effect.timeout или Effect.timeoutOrElse |
| Есть несколько эквивалентных адресов, нужен любой | Effect.race или Effect.raceAll |
| Есть основной плюс резервный, не хочется удваивать трафик всегда | Effect.race плюс Effect.delay на резервном |
| После провала достаточно вернуть последнее известное значение | Effect.catch или Effect.orElseSucceed |
| Транзиентные ошибки, по которым стоит подождать и попробовать снова | Effect.retry со Schedule |
Эти пять рычагов закрывают почти все стратегии устойчивости. В Pulse мы используем четыре из пяти одновременно, разнесённые по слоям.
Что взять с собой
Effect.timeout(d)прерывает внутренний файбер, добавляет встроенныйTimeoutErrorвE.timeoutOrElseвыполняет твой запасной эффект.timeoutOptionубирает таймаут изEи кладёт вOption.- В отличие от
Promise.race(setTimeout), таймаут реально отменяет работу, finalizer-ы срабатывают. Effect.race(a, b)это first success wins: победитель отдаётся, проигравший прерывается. Если победитель упал, ждём проигравшего.- Hedged request это
race(a, b.pipe(delay)): трафик удваивается только на медленном хвосте. Effect.catch/orElseSucceedэто fallback на провал без race и без retry. Удобно для кеша.
Раздел 6 · Календарь, throttle и debounce
Разделы 2 и 3 закрывают расписания вида “повторяй через столько-то”. Осталось два класса задач, которые через spaced и exponential не выражаются.
Шаг 1 · Календарные расписания
“Каждый день в три часа ночи” это не интервал: между запусками может быть 23, 24 или 25 часов, если по дороге переводили время. Для таких случаев есть Schedule.cron:
import { Effect, Schedule } from 'effect';
const nightly = Schedule.cron('0 3 * * *', 'UTC'); // каждый день в 03:00 UTC
const job = cleanupOldEvents.pipe(Effect.repeat(nightly));
Schedule.cron принимает выражение прямо строкой, парсить заранее не нужно. Если предпочитаешь распарсить один раз и переиспользовать, есть Cron.parseUnsafe('0 3 * * *', 'UTC'), и получившийся Cron тот же Schedule.cron тоже примет.
Обрати внимание на второй аргумент: часовой пояс это часть расписания, а не деталь окружения. “Каждый день в три ночи по Москве” и “в три ночи по UTC” это разные расписания, и на границе перехода времени они разъезжаются.
Ещё деталь для типов: у cron-расписания третий параметр не never, а CronParseError. Выражение разбирается лениво, и кривая строка это ошибка расписания, а не исключение при сборке. Именно ради таких случаев в Schedule и появился канал ошибки.
Отдельных кирпичей вроде “только по вторникам” или “только в 15 часов” в Effect 4 нет: весь календарь выражается через Schedule.cron. Если нужно совместить календарь с обычным интервалом, комбинируй их Schedule.max из раздела 3: например, “каждый час, но не чаще, чем разрешает cron”.
Практическая деталь: календарное расписание в одном процессе это не распределённый планировщик. Если сервис работает в трёх копиях, задача выполнится трижды. Нужна одна, значит нужен внешний замок (advisory lock в Postgres, запись в таблице с уникальным ключом на дату).
Шаг 2 · throttle и debounce
Обе штуки про входящий темп, а не про повтор, поэтому живут на Stream, а не на Schedule.
Stream.throttle пропускает не больше N единиц за интервал:
const limited = source.pipe(
Stream.throttle({
cost: (batch) => batch.length,
units: 2,
duration: '100 millis',
strategy: 'shape',
}),
);
cost объясняет, сколько единиц стоит порция. Порция это обычный массив, поэтому чаще всего это просто batch.length, но может быть и вес в байтах. strategy: 'shape' придерживает лишнее, растягивая поток во времени, а 'enforce' выбрасывает превышение. Первое для своих очередей, второе для защиты от злоупотребления.
Stream.debounce отдаёт элемент, только если после него была пауза:
const settled = Stream.make(1, 2, 3).pipe(Stream.debounce('20 millis'));
// оставил: [3]
Три элемента подряд без пауз схлопнулись в один, последний. Ровно то, что нужно для поля поиска, для реакции на всплеск событий, для пересчёта после серии изменений.
Как выбирать, чтобы не путать: throttle ограничивает частоту (пропускаю не чаще, чем), debounce ждёт тишины (реагирую, когда перестали дёргать), а RateLimiter из 19 · HttpClient и обвязка API ограничивает свои исходящие запросы к чужому API. Три разные задачи, три разных инструмента.
Раздел 7 · Pulse · полная политика probe-а
Сцена · что мы собираем
В Pulse у каждого монитора есть:
id, идентификатор;primaryUrlиfallbackUrl, основной и резервный адреса;interval, как часто стучать (из конфига урока 02);timeoutMs, на сколько ждать ответ;- состояние
SlaStateнаTxRef(урок 08), которое решает, какой URL сейчас активен.
К концу этого раздела мы собираем probeWithPolicy, который:
- Берёт snapshot
SlaState, понимает, на какой URL идти; - Запускает HTTP-запрос с таймаутом 5 секунд (через
Effect.timeoutOrElse); - Если основной URL не успел за 200ms, параллельно стартует резервный (hedged);
- На NetworkError/TimeoutError/5xx делает retry с экспонентой, jitter, потолком 30 секунд, максимум 5 попыток;
- Записывает результат через
Sla.recordSuccessилиSla.recordFailure; - Возвращает результат;
- Внешний цикл повторяет всё это через
Schedule.fixed(interval); - Прерывание (
Fiber.interrupt/SIGTERM) корректно сворачивает все вложенные файберы.
Шаг 1 · политика retry в отдельном файле
// pulse-<nick>/src/runtime/schedule.ts
import { Duration, Schedule } from 'effect';
import type { ProbeError } from '../errors.ts';
/** Само расписание: экспонента с потолком, шумом и лимитом попыток. */
export const retryWithBackoff = Schedule.min([
Schedule.exponential('200 millis', 2),
Schedule.spaced(Duration.seconds(30)), // потолок 30s
]).pipe(
Schedule.jittered, // шум
Schedule.upTo({ times: 5 }), // максимум 5 ретраев
);
// retryWithBackoff: Schedule<Duration, unknown, never, never>
/** Что именно ретраить. Живёт отдельно от того, как долго ждать. */
export const isTransient = (error: ProbeError): boolean =>
error._tag === 'NetworkError' ||
error._tag === 'TimeoutError' ||
(error._tag === 'HttpStatusError' && error.status >= 500);
Расписание вынесено в отдельный файл, чтобы его можно было использовать в нескольких местах (probe, DNS-резолвер, whois) и тестировать отдельно.
Заметь, что политика распалась на две независимые части: расписание не знает про типы ошибок Pulse, а предикат не знает про задержки. Соединяются они в точке применения, в опциях Effect.retry. Раньше предикат приходилось вшивать внутрь расписания, и такое расписание было уже не переиспользовать: оно намертво привязывалось к одному типу ошибки.
Шаг 2 · hedged probe с таймаутом
// pulse-<nick>/src/probe/hedged.ts
import { Effect } from 'effect';
import { HttpClient } from '../services/http-client.ts';
import { TimeoutError } from '../errors.ts';
const probeOne = (url: string) =>
Effect.gen(function* () {
const http = yield* HttpClient;
return yield* http.get(url);
});
export const hedgedProbe = (primaryUrl: string, fallbackUrl: string) =>
Effect.race(
probeOne(primaryUrl),
probeOne(fallbackUrl).pipe(Effect.delay('200 millis')),
).pipe(
Effect.timeoutOrElse({
duration: '5 seconds',
orElse: () => Effect.fail(new TimeoutError({ url: primaryUrl, deadline: 5000 })),
}),
);
// hedgedProbe: (primary, fallback) => Effect<Response, ProbeError, HttpClient>
Что здесь происходит. Основной URL стартует немедленно. Резервный, с задержкой 200ms. Если основной успел за 200ms, резервный отменяется не родившись. Если медлит, резервный успевает запуститься параллельно, и Effect.race берёт первый успешный ответ. Сверху общий таймаут 5 секунд, после него оба прерываются и наружу летит TimeoutError.
Шаг 3 · вписываем SlaState и журнал
// pulse-<nick>/src/probe/with-policy.ts
import { Effect, PubSub, Result } from 'effect';
import { MonitorEventsPubSub } from '../concurrency/coordination.ts';
import { Sla } from '../concurrency/sla-state.ts';
import { hedgedProbe } from './hedged.ts';
import { isTransient, retryWithBackoff } from '../runtime/schedule.ts';
import type { MonitorTarget } from '../config.ts';
export const probeWithPolicy = (target: MonitorTarget) =>
Effect.gen(function* () {
const sla = yield* Sla;
const bus = yield* MonitorEventsPubSub;
const state = yield* sla.snapshot;
const [activeUrl, backupUrl] =
state.active === 'primary'
? [target.primaryUrl, target.fallbackUrl]
: [target.fallbackUrl, target.primaryUrl];
const attempt = hedgedProbe(activeUrl, backupUrl).pipe(
Effect.retry({ schedule: retryWithBackoff, while: isTransient }),
);
const outcome = yield* attempt.pipe(Effect.result);
if (Result.isSuccess(outcome) && outcome.success.status < 500) {
yield* sla.recordSuccess;
yield* PubSub.publish(bus, {
_tag: 'ProbeSuccess',
targetId: target.id,
url: activeUrl,
status: outcome.success.status,
});
return outcome.success;
}
const next = yield* sla.recordFailure;
const error = Result.isFailure(outcome)
? outcome.failure
: new Error(`status ${outcome.success.status}`);
yield* PubSub.publish(bus, {
_tag: 'ProbeFailure',
targetId: target.id,
url: activeUrl,
sla: next,
});
return yield* Effect.fail(error);
});
Снаружи probeWithPolicy это одно действие, которое всегда либо записывает успех в журнал и возвращает Response, либо записывает провал и пробрасывает ошибку. Внутри он включает в себя SlaState-снапшот, hedged-запрос, retry-политику, маршрутизацию active/backup, и публикацию событий. Каждый кусок мы уже собирали в отдельных уроках; здесь они работают вместе.
Две детали формы, которые повторяются по всему Pulse. Effect.result заворачивает исход в Result (тот же тип, что раньше звался Either), и читается он через Result.isSuccess плюс поля .success и .failure. Публикация в шину это PubSub.publish(bus, event), потому что PubSub в Effect 4 это данные, а операции живут в модуле (07 · Координация).
Шаг 4 · регулярный цикл через Effect.repeat
// pulse-<nick>/src/probe/loop.ts
import { Duration, Effect, Schedule } from 'effect';
import { probeWithPolicy } from './with-policy.ts';
import type { MonitorTarget } from '../config.ts';
export const probeLoop = (target: MonitorTarget) =>
probeWithPolicy(target).pipe(
Effect.catch(() => Effect.void), // не валим цикл на провале
Effect.repeat(Schedule.fixed(Duration.millis(target.intervalMs))),
);
// probeLoop: (target) => Effect<never, never, ...services>
Что важно. Effect.catch(() => Effect.void) это обязательный шаг: без него любая ошибка из probeWithPolicy уронит цикл, и repeat прервётся. Мы уже зафиксировали ошибку в журнале, опубликовав событие в шину, поэтому здесь нам важна только сторона “цикл живёт”. Если опустить catch, при первом же 5xx repeat остановится, и мониторинг перестанет работать.
И ещё раз про границу: Effect.catch ловит ожидаемые ошибки, те, что стоят в типе E. Дефекты (Effect.die, непойманное исключение из чужой библиотеки) сюда не попадут и убьют файбер монитора. Это правильное поведение: ошибка в коде должна быть видна, а не проглочена вечным циклом.
Schedule.fixed(Duration.millis(target.intervalMs)) это “запускай probe в начале каждой клетки intervalMs”. Если сама проба заняла 4 секунды при интервале 60, следующая стартует через 56. Если заняла 70 секунд, следующая стартует сразу же.
Шаг 5 · форкаем по одному файберу на монитор
// pulse-<nick>/src/probe/manager.ts
import { Effect } from 'effect';
import { probeLoop } from './loop.ts';
import type { MonitorTarget } from '../config.ts';
export const startMonitors = (targets: ReadonlyArray<MonitorTarget>) =>
Effect.forEach(targets, (target) => Effect.forkScoped(probeLoop(target)));
// startMonitors: (...) => Effect<Fiber<never, never>[], never, Scope + services>
Каждый монитор живёт в собственном файбере. Effect.forkScoped (см. 04 · Services и Layer) привязывает файбер к Scope-у, который держит MainLive. На runtime.dispose() все эти файберы получают interrupt, finalizer-ы отрабатывают, fetch-и отменяются, и Pulse выходит без висящих сокетов.
Шаг 6 · общий рантайм и две точки входа
// pulse-<nick>/src/runtime.ts
import { ManagedRuntime } from 'effect';
import { MainLive } from './main.ts';
export const runtime = ManagedRuntime.make(MainLive);
// pulse-<nick>/src/cli/watch.ts
import { Effect } from 'effect';
import process from 'node:process';
import { Config } from '../config.ts';
import { startMonitors } from '../probe/manager.ts';
import { runtime } from '../runtime.ts';
const program = Effect.gen(function* () {
const config = yield* Config;
yield* startMonitors(config.targets);
yield* Effect.never; // главный fiber висит, дети работают
});
process.on('SIGINT', async () => {
await runtime.dispose();
process.exit(0);
});
await runtime.runPromise(Effect.scoped(program));
// pulse-<nick>/src/main.ts
import http from 'node:http';
import { runtime } from './runtime.ts';
import { statusProgram } from './http/status.ts'; // та же сборка MainLive
// embedded HTTP-вход на общем runtime, без HTTP-модулей Effect
const server = http.createServer((_req, res) => {
runtime.runPromise(statusProgram).then((snapshot) => {
res.setHeader('content-type', 'application/json');
res.end(JSON.stringify(snapshot));
});
});
server.listen(8080);
Один runtime. Внутри одна сборка MainLive, один Storage, один SlaState, один пул мониторов (стартует через pulse watch). HTTP-вход не запускает свои мониторы заново, он читает из общего MonitorEvents. Это и есть та “одна программа, две поверхности”, о которой шла речь в уроке 01. В уроке 12 этот ручной node:http станет pulse serve на NodeHttpServer, добавятся /api/monitors и SSE /events, но точка сборки и runtime остаются ровно эти.
Шаг 7 · тест на политику через TestClock
Тест проверяет: за 10 минут при интервале 60 секунд должно случиться ровно 10 проб. С реальными sleep это 10 минут ожидания, с TestClock (13 · Testing) это 0 миллисекунд:
import { describe, it } from '@effect/vitest';
import { Duration, Effect, Ref, Schedule } from 'effect';
import { TestClock } from 'effect/testing';
import { expect } from 'vitest';
describe('probeLoop', () => {
it.effect('runs probe ten times over ten minutes at 60s interval', () =>
Effect.gen(function* () {
const count = yield* Ref.make(0);
const probe = Ref.update(count, (n) => n + 1);
const loop = probe.pipe(Effect.repeat(Schedule.fixed(Duration.seconds(60))));
yield* Effect.forkChild(loop);
yield* TestClock.adjust('10 minutes');
yield* Effect.yieldNow;
const result = yield* Ref.get(count);
// первое выполнение плюс 10 повторов через каждые 60 секунд за 10 минут
// = 11 запусков (или 10, в зависимости от границы окна)
expect(result).toBeGreaterThanOrEqual(10);
expect(result).toBeLessThanOrEqual(11);
}).pipe(Effect.provide(TestClock.layer())),
);
});
Три вещи, которые в этом тесте выглядят иначе, чем в старых примерах. TestClock приезжает обычным слоем из подмодуля effect/testing, отдельного TestContext нет. Effect.fork называется Effect.forkChild: в имени теперь видно, что файбер привязан к родителю. И Effect.yieldNow это значение, а не функция, вызывать его скобками не надо.
Тонкость: первое выполнение происходит немедленно, Schedule.fixed ставит следующее через 60 секунд от старта. За 10 минут проходит граница [старт + 10 * 60 секунд]. Сколько именно проб попадёт внутрь окна, зависит от того, считаешь ли ты сам старт. На практике пиши тест с диапазоном, не с равенством на единицу.
Что взять с собой
- Политика probe-а это композиция:
hedgedProbe(race с задержкой) внутриEffect.timeoutOrElse(потолок) внутриEffect.retry({ schedule, while })снаружи. - Регулярность приходит снаружи через
Effect.repeat(Schedule.fixed(interval))плюсEffect.catchдля защиты цикла. - Каждый монитор это свой файбер, форкается через
Effect.forkScoped, прерывается наruntime.dispose(). - Один
ManagedRuntimeобслуживает иpulse watch, и embedded HTTP-вход наnode:http. Один Layer, один MemoMap, один Storage. В уроке 12 на том же рантайме подниметсяpulse serve. TestClock.adjust(d)делает тест на политику детерминированным и мгновенным.
Финал · чек-лист
ДЗ
Дальше
Следующий урок · 12 · Production-обвязка: CLI, Terminal, HTTP-сервер. После расписания и таймаутов у Pulse уже есть полноценное ядро. Остаётся надеть на него две поверхности: CLI на effect/unstable/cli (живой ANSI-дашборд по pulse watch) и HTTP-сервер на effect/unstable/http плюс @effect/platform-node (/status, /api/monitors, SSE /events). Обе поверхности используют один и тот же runtime из этого урока.
Контекст по ManagedRuntime и Layer-у живёт в 04 · Services и Layer. Schedule в стримах подробнее в 09 · Stream, TestClock для тестов на расписание в 13 · Testing. Если эта политика probe-а напоминает тебе про circuit breaker, это потому что она и есть его естественное расширение: SlaState (урок 08) держит счётчик подряд идущих провалов, retryWithBackoff (этот урок) выживает на одном раунде, hedged-запрос экономит p99.