Файберы и concurrency: fork, join, interrupt, Effect.all, Semaphore, race
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
Файберы и concurrency: fork, join, interrupt, Effect.all, Semaphore, race
Сцена · кухня и пятеро поваров
Шеф (родитель), повара (дочерние файберы). Шеф позвал “стоп смены”, все повара одновременно положили ножи и пошли мыть руки. Никакого “забытого повара”, который доделывает суп уже после смены, потому что он “обещал”. В JavaScript обычный Promise ровно про этого повара: запустил fetch и забыл, отменить нельзя, проигравший в Promise.race тоже доваривает. Effect устроен по-другому: дочерний файбер подчиняется родителю, и когда родитель закончился (или его прервали), дочерние тоже прерываются. Это называется structured concurrency, и она у тебя уже есть, ничего специально включать не надо.
Карта урока
Собираем параллельный пробинг в Pulse: SIGTERM выключает дерево файберов чисто, без забытых fetch-ей и без бесконтрольной нагрузки на сеть. Десять разделов:
- Раздел 1, файбер как runtime-инстанция,
Effect.logи id файбера. - Раздел 2,
Effect.forkChild, дочерний файбер и его lifetime. - Раздел 3, кооперативный шедулинг и suspension points.
- Раздел 4,
Fiber.joinиFiber.await, разница в распространении ошибок. - Раздел 5, прерывание файберов:
Fiber.interrupt,Effect.onInterrupt,Effect.uninterruptible,Effect.interrupt,Effect.runCallback. - Раздел 6,
Effect.forkDetach, escape hatch и почему его почти не пишут. - Раздел 7,
Effect.allи bounded concurrency, своя версия черезSemaphore. - Раздел 8,
Semaphoreкак примитив: mutex, rate-limit. - Раздел 9,
Effect.race, кто быстрее, проигравший прерывается. - Раздел 10, Pulse, параллельный пробинг и graceful shutdown.
Раздел 1 · Файбер, не поток
Что такое файбер
Каждый раз, когда ты запускаешь эффект (Effect.runPromise, Effect.runPromiseExit, runtime.runFork), Effect создаёт файбер, чтобы исполнить его. Файбер, это runtime-сущность, которая владеет одной конкретной программой. В документации и в API он называется Fiber, по-русски дальше везде файбер. Это не поток ОС: V8 однопоточный, никаких контекст-свитчей ядра здесь нет. Не Worker: так как это один и тот же V8-процесс, использует ту же память, что и вся программа. Это легковесный пользовательский поток поверх event loop, ближайший родственник, goroutine в Go или virtual thread в Java.
Размер у файбера маленький, миллион файберов в одном процессе норма. Fork делается дёшево, без системного вызова, просто аллокация в куче.
Effect.log не равно console.log
Внутри файбера ты не пишешь console.log напрямую. Берёшь Effect.log:
import { Effect } from 'effect';
const hello = Effect.gen(function* () {
yield* Effect.log('main fiber');
return 'hello';
});
Effect.runPromise(hello);
// timestamp=2026-05-11T12:00:00.000Z level=INFO fiber=#0 message="main fiber"
Заметь fiber=#0. Это id файбера, который сейчас исполняется. У главного запуска id равен #0. Каждый новый fork получает следующий номер. Это первый и самый дешёвый способ отлаживать concurrency: смотришь на префикс лога и видишь, кто что сделал.
Поверх Effect.log ты получаешь:
- минимальный уровень (через окружение, в проде только
INFO+, локальноDEBUG+); - разные форматтеры (текстовый, structured-JSON);
- свой logger (например,
pinoили Datadog), подключаемый через layer.
Дешёвый отладочный приём: натыкай Effect.log в подозрительные места, смотри на id файбера и порядок строк.
yield* блокирует, как await
Когда внутри Effect.gen(function* () {...}) ты пишешь yield* effect, генератор останавливается до завершения этого эффекта. Аналог await promise:
import { Effect } from 'effect';
const program = Effect.gen(function* () {
yield* Effect.log('beforeSleep');
yield* Effect.sleep('2 seconds');
yield* Effect.log('afterSleep');
return 'hello';
});
Effect.runPromise(program);
// fiber=#0 message="beforeSleep"
// (двухсекундная пауза)
// fiber=#0 message="afterSleep"
Один файбер, шаги идут по порядку. Это тот же ментальный образ, что у async function. Чтобы получить параллельность, нужно явно отделить дочерний файбер, и тут вступает Effect.forkChild.
Что взять с собой
- Файбер, это runtime-инстанция эффекта. Не thread ОС, не Worker, легковесный.
Effect.logпишет в логгер с id файбера, минимальным уровнем и форматтером, не путай сconsole.log.yield* effectблокирует текущий файбер до завершения эффекта.
Раздел 2 · Fork и lifetime файбера
Effect.forkChild, дочерний файбер
В обычном TypeScript ты можешь “забыть” promise, не дожидаясь его:
fetch('/api/track'); // запустили и пошли дальше
await otherWork();
В Effect для этого есть Effect.forkChild. Он берёт эффект, запускает его в новом файбере и сразу возвращает RuntimeFiber:
import { Effect } from 'effect';
Effect.gen(function* () {
yield* Effect.sleep('2 seconds').pipe(
Effect.tap(() => Effect.log('afterSleep')),
Effect.forkChild, // создаёт дочерний fiber
);
yield* Effect.log('parent finishing');
return 'hello';
}).pipe(Effect.runPromise);
Запусти этот код и удивись: ты увидишь только parent finishing. Никакого afterSleep. Где он?
Жизнь дочернего привязана к родителю
Effect.forkChild создаёт файбер, чей lifetime привязан к родительскому файберу. Родитель закончился, и в этот же момент рантайм посылает дочернему Fiber.interrupt. В нашем случае родитель проходит мимо Effect.forkChild за миллисекунду, выводит “parent finishing”, возвращает 'hello', и заканчивается. Дочерний файбер, который только что начал ждать Effect.sleep('2 seconds'), прерывается.
Лечится держанием родителя в живых:
Effect.gen(function* () {
yield* Effect.sleep('2 seconds').pipe(
Effect.tap(() => Effect.log('afterSleep')), // fiber=#1
Effect.forkChild,
);
yield* Effect.log('parent finishing'); // fiber=#0
yield* Effect.sleep('3 seconds'); // родитель переживает ребёнка
return 'hello';
}).pipe(Effect.runPromise);
// fiber=#0 message="parent finishing"
// (через 2 секунды)
// fiber=#1 message="afterSleep"
// (ещё секунда)
// (программа завершается)
Теперь видно оба лога, и видно разные id файберов: #0 родитель, #1 ребёнок. В этот раз структурная отмена сыграла за нас: родитель прожил дольше ребёнка, и тот успел отработать.
Это поведение тот самый дом, в котором живут все Effect-овские параллельные операции. Effect.all, Effect.race, любые комбинаторы concurrency, опираются на ту же гарантию: дочерние не переживут родителя.
Что взять с собой
Effect.forkChildзапускает эффект в дочернем файбере и возвращаетRuntimeFiber<A, E>.- Lifetime дочернего файбера вложен в lifetime родителя. Закрылся родитель, дочерний прерывается.
- Хочешь, чтобы дочерний дожил, держи родителя в живых (sleep, join, race, scope).
Раздел 3 · Кооперативный шедулинг и suspension points
Лог, которого не будет
Тонкая ловушка. Угадай, что выведет этот код:
import { Effect } from 'effect';
Effect.gen(function* () {
yield* Effect.log('beforeSleep').pipe(Effect.forkChild);
yield* Effect.log('parent finishing');
return 'hello';
}).pipe(Effect.runPromise);
Только parent finishing. Дочерний файбер на Effect.log('beforeSleep') так и не запустился, хотя мы его форкнули. И это не баг.
Кооперативная vs вытесняющая многозадачность
Effect использует кооперативный шедулинг. Это значит: текущий файбер держит управление, пока сам не отдаст его. Контроль переходит на другой файбер только в одной из ситуаций:
- Текущий файбер завершился (вернул значение или упал).
- Текущий файбер попал в suspension point (см. ниже).
Effect.forkChild не запускает дочерний файбер немедленно, он его планирует. А когда планировщик доберётся до него, зависит от того, отдаёт ли текущий файбер контроль.
В нашем коде после fork идёт ровно один синхронный Effect.log и return. Текущий файбер успевает закончиться раньше, чем дочерний получит свой квант времени. Дочерний файбер отменяется (см. раздел 2), и кажется, что его не было.
Кажется, что вытесняющая многозадачность проще (“просто запусти, когда сможешь”), но она платит race-conditions: в любой момент тебя могут прервать, любая невалидная промежуточная структура данных видна другим. Effect выбирает кооператив сознательно.
Аналогия · Windows 3.1 vs Windows NT
Windows 3.1 был кооперативным: программа сама вызывала Yield() между задачами. Пока никто не зависал в цикле, система ощущалась плавной. Стоило одному приложению зациклиться, и всё колом, потому что снаружи прервать его было нечем. Windows NT и все современные ОС используют вытесняющую многозадачность: ядро по таймеру и системным прерываниям выкидывает любой поток с CPU в произвольной точке.
Effect-файберы устроены ближе к Windows 3.1, не к NT, и это не лень разработчиков, это ограничение JavaScript. V8 однопоточный, нет сигналов, нет VM-hook-ов на чужой синхронный код, снаружи прервать while (true) нечем. Любая JS-библиотека concurrency (Effect, RxJS, async/await) кооперативна по построению. Effect просто покрывает свой DSL частой сеткой suspension points: пока ты остаёшься внутри Effect-инструкций, ощущение почти как от вытесняющей многозадачности, рантайм успевает дать ход другим файберам, обработать прерывание, проверить таймауты. Вышел в голый sync-цикл без yield, и ты заблокировал шедулер, как Windows 3.1 на зависшем приложении.
К этой границе вернёмся в разделе 5, когда будем разбирать прерывание: файбер, ушедший в синхронный цикл без yield-ов, прервать не получится, рантайм узнает о запросе только когда цикл сам отдаст управление. Если в твоей задаче реально нужна вытесняющая отмена CPU-bound работы, ответ один: вынести её в Worker, который снаружи можно terminate(). Это уже отдельный V8-инстанс, и его рантайм действительно умеет грохнуть.
Плохой sync-цикл и его лечения
Конкретика на код. Вот цикл, который рантайм увидеть не сможет:
const bad = Effect.sync(() => {
for (let i = 0; i < 1e9; i++) {
// тяжёлая работа, ни одного yield
}
});
Пока этот Effect.sync не вернётся, шедулер мёртв. Fiber.interrupt, прилетевший снаружи, дойдёт до bad только когда цикл досчитает до конца. Все остальные файберы стоят в очереди.
Лечения два, обычно используют что-то одно.
1. Разбить работу на шаги через Effect.iterate. Между шагами рантайм видит границу, проверяет interrupt-флаг, переключается на других:
const good = Effect.iterate(0, {
while: (i) => i < 1e9,
body: (i) => Effect.sync(() => i + 1),
});
2. Оставить обычный цикл, но периодически отдавать управление руками через Effect.yieldNow. Полезно, когда переписать на iterate дорого (например, цикл уже завязан на внешний массив):
const explicit = Effect.gen(function* () {
for (let i = 0; i < 1e9; i++) {
if (i % 10_000 === 0) yield* Effect.yieldNow();
doWork(i);
}
});
Шаг “раз в 10 000 итераций” подбирается под характер работы: чаще yield-ить = больше отзывчивости и накладных, реже = меньше накладных, но дольше реакция на прерывание. Точно мерять, не угадывать.
Suspension points
Чтобы вернуть управление рантайму, нужен suspension point. Это любой эффект, который “приостанавливает” исполнение, дав шедулеру шанс. Suspension создают:
| Что | Когда вставлять |
|---|---|
Effect.yieldNow | явный yield, “сейчас в самый раз” |
Effect.sleep('2 seconds') | задержка |
Effect.delay(eff, '500 millis') | задержка эффекта |
Effect.promise(() => ...), Effect.tryPromise | граница в Promise-мир |
Effect.async((cb) => ...) | граница в callback-API |
Fiber.join(fiber), Fiber.await(fiber) | ожидание другого файбера |
| любые комбинаторы поверх перечисленного | косвенно |
Синхронный Effect.log, Effect.succeed, Effect.sync, Effect.flatMap поверх sync-эффектов не дают suspension. Это как await Promise.resolve(value), где у Promise-а уже есть результат, microtask-обхода нет.
Effect.yieldNow, ручной yield
Если форкнул синхронный эффект и хочешь, чтобы он успел запуститься до завершения родителя, поставь после fork явный Effect.yieldNow:
import { Effect } from 'effect';
Effect.gen(function* () {
yield* Effect.log('beforeSleep').pipe(Effect.forkChild);
yield* Effect.yieldNow(); // отдали управление
yield* Effect.log('parent finishing');
}).pipe(Effect.runPromise);
// fiber=#1 message="beforeSleep"
// fiber=#0 message="parent finishing"
Но в реальном коде ты чаще видишь suspension через sleep, promise, ожидание других файберов, и явный yieldNow нужен редко. Он полезен, когда форкнул чисто синхронную работу и хочешь дать ей квант времени до того, как пойти дальше.
Что взять с собой
- Effect шедулит файберы кооперативно: текущий держит контроль до завершения или suspension.
Effect.forkChildне запускает эффект немедленно, он его планирует.- Suspension создают
sleep,delay,promise,tryPromise,async, ожидание других файберов и любые их обёртки. - Если форкнул синхронный эффект, поставь
Effect.yieldNow, чтобы дать ему ход.
Раздел 4 · Fiber.join и Fiber.await
Подождать дочерний файбер
Effect.forkChild отдаёт RuntimeFiber<A, E>. Это как Promise<A> без await: задача запущена, но ты её не ждёшь. Чтобы дождаться результата, бери одно из двух: Fiber.join или Fiber.await. Разница в том, что они делают с ошибкой.
Fiber.join, ошибка пробрасывается
import { Effect, Fiber } from 'effect';
const slow = Effect.gen(function* () {
yield* Effect.sleep('2 seconds');
return 50;
});
const program = Effect.gen(function* () {
const fiber = yield* Effect.forkChild(slow);
const result = yield* Fiber.join(fiber);
// результат типизирован как number
yield* Effect.log(`slow: ${result}`);
});
Fiber.join(fiber) приостанавливает текущий файбер до завершения forked-файбера, потом отдаёт его результат. Если forked-файбер упадёт, ошибка попадёт в канал E родителя, как будто вызвал Effect.fail напрямую.
import { Data, Effect, Fiber } from 'effect';
class CustomError extends Data.TaggedError('CustomError') {}
const slow = Effect.gen(function* () {
yield* Effect.sleep('2 seconds');
return yield* new CustomError();
});
const program = Effect.gen(function* () {
const fiber = yield* Effect.forkChild(slow);
const result = yield* Fiber.join(fiber);
return result;
});
// program: Effect<number, CustomError, never>
CustomError пробрасывается. Канал ошибок родителя расширился до CustomError, как если бы slow был просто yield* slow. Хочешь, чтобы падение ребёнка валило родителя, бери Fiber.join.
Fiber.await, получи Exit без распространения
Иногда ребёнок может упасть, и ты хочешь это обработать, а не пробросить:
import { Effect, Fiber } from 'effect';
const program = Effect.gen(function* () {
const fiber = yield* Effect.forkChild(slow);
const exit = yield* Fiber.await(fiber);
// exit: Exit<number, CustomError>
if (exit._tag === 'Success') {
yield* Effect.log(`got: ${exit.value}`);
} else {
yield* Effect.log('failed, but no propagation');
}
});
// program: Effect<void, never, never>
// канал ошибок никак не расширился
Fiber.await тоже приостанавливает родителя, но возвращает Exit<A, E>. Через Exit.match (см. 05 · Resources) ты ветвишь логику: записать в журнал, отправить уведомление, ретраить. Канал ошибок родителя не расширяется, потому что forked-ошибка не пробрасывается автоматически.
Когда какой брать
| Сценарий | Бери |
|---|---|
| Ошибка ребёнка должна валить родителя | Fiber.join |
| Просто получить результат, ошибки наверх | Fiber.join |
Хочу Exit с разветвлением на Success/Failure/Interrupt | Fiber.await |
| Падение ребёнка не должно автоматически валить родителя | Fiber.await |
Оба ставят suspension point, поэтому через них работает прерывание (см. раздел 5). Разница строго в распространении ошибки.
Что взять с собой
Fiber.join(fiber)ждёт результат, пробрасывает ошибку.Fiber.await(fiber)ждёт результат, отдаётExit<A, E>, ошибку не пробрасывает.- Оба эффекта приостанавливают вызывающий файбер, через них же его можно прервать.
Раздел 5 · Прерывание файберов
Раздел длинный, потому что у прерывания четыре API под одной идеей. Сначала прервём файбер снаружи (Fiber.interrupt), потом изнутри (Effect.interrupt), потом из императивной точки (Effect.runCallback). Параллельно разберём защиту критической секции (Effect.uninterruptible) и хук на cleanup (Effect.onInterrupt). Все они работают на одном механизме: прерывание срабатывает на следующем suspension point.
Fiber.interrupt, у тебя есть ссылка на файбер
Если у тебя в руках RuntimeFiber, его можно прервать:
import { Effect, Fiber } from 'effect';
const slow = Effect.gen(function* () {
yield* Effect.sleep('2 seconds');
return 50;
});
const program = Effect.gen(function* () {
yield* Effect.log('before fork');
const fiber = yield* Effect.forkChild(slow);
yield* Effect.yieldNow();
const exit = yield* Fiber.interrupt(fiber);
// exit: Exit<number, never>
yield* Effect.log('after fork');
});
Что это даёт:
- если файбер уже завершился,
Fiber.interruptтут же возвращает егоExit; - если ещё работает, посылает прерывание и ждёт, пока файбер аккуратно завершится;
- всегда возвращает
Exit, не значение.
Effect.yieldNow нужен потому, что Fiber.interrupt срабатывает на следующем suspension point дочернего файбера. Без явного yield родительский файбер не отдаст управление, и прерывание не дойдёт до дочернего файбера.
Cleanup через Effect.onInterrupt
Effect.onInterrupt (из 05 · Resources) висит на дочернем эффекте и срабатывает, когда тот прерывают:
const slow = Effect.gen(function* () {
yield* Effect.sleep('2 seconds');
return 50;
}).pipe(Effect.onInterrupt(() => Effect.log('interrupted slow')));
При прерывании файбера interrupted slow появится в журнале до того, как файбер окончательно умрёт. Это рабочий механизм cleanup-а, на нём же построено Effect.acquireRelease и закрытие scope.
Прерывание это suspension-зависимая операция
Без suspension point прерывание не дойдёт. Если файбер крутит синхронный цикл без yield-ов, его невозможно прервать чисто:
const tight = Effect.sync(() => {
for (let i = 0; i < 1e9; i++) {
// ничего не yield-им
}
});
Чтобы прерывание сработало, нужно либо разбить цикл на куски и yield-ить периодически, либо использовать встроенные комбинаторы (Effect.repeat, Effect.forEach, любые async-эффекты). На практике tight лучше переписать через Stream с шагами.
Глубже ·
Causeпосле прерывания.ExitпослеFiber.interruptнесётCause.Sequentialс вложенными interrupt-ами: внешний слой это сам файбер, внутренний это, например,Effect.sleep, на котором он сидел. На практике хватаетExit.match(success/failure) иCause.isInterrupt(cause)(это было прерывание?). Подробнее проCauseв 03 · Errors.
Effect.uninterruptible, защита критической секции
Иногда есть участок, который нельзя прерывать посередине: запись в файл, транзакция, освобождение ресурса. Заворачивай его в Effect.uninterruptible:
const slow = Effect.gen(function* () {
yield* Effect.sleep('2 seconds');
return 50;
}).pipe(Effect.uninterruptible);
Даже если Fiber.interrupt прилетит на этот файбер, он сначала закончит slow, и только потом отработает прерывание (если оно ещё актуально). Внутри Effect почти все finalizer-ы автоматически uninterruptible: иначе они могли бы прерваться на середине cleanup-а, что хуже отсутствия cleanup-а.
Не злоупотребляй: чем шире uninterruptible-блок, тем дольше реагирует SIGTERM. Размечай ровно “критическая запись”, не “всё подряд”.
Effect.interrupt, прервать самого себя
А что, если ссылки на файбер нет, но какая-то логика говорит “пора заканчивать”? Бери Effect.interrupt:
import { Effect } from 'effect';
const repeat = Effect.log('hello').pipe(
Effect.delay('300 millis'),
Effect.repeat({ times: 20 }),
);
const program = Effect.gen(function* () {
yield* Effect.forkChild(repeat);
yield* Effect.sleep('2 seconds');
return yield* Effect.interrupt;
}).pipe(Effect.onInterrupt(() => Effect.log('interrupted')));
Effect.runPromise(program);
// hello, hello, ... (примерно 6 раз за 2 секунды)
// interrupted
Effect.interrupt прерывает текущий файбер. Через structured concurrency прерывание автоматически дойдёт и до его дочерних файберов (repeat тоже остановится).
Важная деталь: пиши return yield* Effect.interrupt, не просто yield* Effect.interrupt. Effect.interrupt это эффект, который “падает” как Effect.fail, и компилятор отметит, что после него код недостижим. Без return редактор может думать, что выполнение продолжится.
Effect.never плюс Effect.runCallback, отмена из императивного мира
В React, в обработчике события, в любой императивной точке у тебя нет файбера, чтобы его прервать. Effect даёт Effect.runCallback, который запускает программу и возвращает функцию-cancel:
import { Effect } from 'effect';
const repeat = Effect.log('hello').pipe(
Effect.delay('300 millis'),
Effect.repeat({ times: 20 }),
);
const program = Effect.gen(function* () {
yield* Effect.forkChild(repeat);
return yield* Effect.never; // блокирует fiber вечно
}).pipe(Effect.onInterrupt(() => Effect.log('interrupted')));
const cancel = Effect.runCallback(program);
// где-то позже
setTimeout(() => cancel(), 2000);
Effect.never это эффект, который никогда не возвращает (как while (true) {}, только без цикла). Программа крутится, пока её не прервут. Effect.runCallback отдаёт функцию cancel, которая шлёт Fiber.interrupt главному файберу. Идеально под React:
useEffect(() => {
const cancel = Effect.runCallback(myProgram);
return cancel; // на unmount компонента вернёт cleanup
}, []);
При unmount или смене зависимостей React дёрнет cleanup, и файбер аккуратно прервётся вместе со всеми дочерними. Никаких “забытых” подписок и таймеров.
Каскадное прерывание дерева
Перед таблицей “что когда брать”, полезно увидеть структурную отмену в действии на одном маленьком примере. Сцена банальная: HTTP-обработчик породил три параллельные подзадачи, клиент отвалился, нам надо аккуратно всё отменить.
import { Effect, Fiber } from 'effect';
const subtask = (name: string, ms: number) =>
Effect.gen(function* () {
yield* Effect.addFinalizer(() => Effect.log(` [${name}] cleanup`));
yield* Effect.log(` [${name}] started`);
yield* Effect.sleep(`${ms} millis`);
yield* Effect.log(` [${name}] finished`);
}).pipe(Effect.scoped);
const handleRequest = Effect.gen(function* () {
yield* Effect.addFinalizer(() => Effect.log('[request] cleanup'));
yield* Effect.log('[request] started');
yield* Effect.all(
[subtask('db-query', 3000), subtask('cache', 2000), subtask('audit-log', 1000)],
{ concurrency: 'unbounded' },
);
}).pipe(Effect.scoped);
const program = Effect.gen(function* () {
const fiber = yield* Effect.forkChild(handleRequest);
yield* Effect.sleep('500 millis');
yield* Effect.log('>>> client disconnected, killing request <<<');
yield* Fiber.interrupt(fiber);
});
Effect.runPromise(program);
// [request] started
// [db-query] started
// [cache] started
// [audit-log] started
// >>> client disconnected, killing request <<<
// [db-query] cleanup
// [cache] cleanup
// [audit-log] cleanup
// [request] cleanup
Что произошло шаг за шагом:
- Прибили root,
Fiber.interruptпоставил флаг прерывания наhandleRequest. handleRequestсидел вEffect.all, на следующем suspension point flag дошёл и до него, и до всех трёх форкнутых подзадач.- На каждой подзадаче отработал её finalizer (
[name] cleanup). - Когда внутренние scope-ы закрылись, отработал finalizer самого
handleRequest.
Порядок чистки строго LIFO, как using в C#, bracket в Haskell или RAII в C++. Это и есть structured concurrency на одном экране: тебе не пришлось ни таскать AbortController через сигнатуры, ни вручную перебирать дочерние файберы, ни вычищать таймеры. Pulse в разделе 10 это та же конструкция, только finalizer-ы закрывают реальные AbortController, а в роли root-а выступает обработчик SIGTERM.
Когда что использовать
| Нужно | Бери |
|---|---|
есть RuntimeFiber, прервать его | Fiber.interrupt(fiber) |
| прервать “себя и детей” из логики внутри | return yield* Effect.interrupt |
| прервать программу из imperative-кода (React, обработчик события) | Effect.runCallback(program) плюс возвращённый cancel |
| запустить и не давать завершиться, пока не прервут | Effect.never |
| защитить критическую секцию | Effect.uninterruptible(eff) |
| зацепить cleanup на прерывание | Effect.onInterrupt(() => eff) |
Что взять с собой
Fiber.interrupt(fiber)ждёт окончания, возвращаетExit.- Прерывание срабатывает на suspension points, без них файбер не прервёшь.
Effect.uninterruptibleзащищает критические секции.return yield* Effect.interruptпрерывает текущий файбер и дочерних.Effect.runCallback(program)отдаётcancel-функцию для imperative-кода (React useEffect cleanup).
Раздел 6 · Effect.forkDetach, escape hatch
Когда Effect.forkChild мало
Effect.forkChild привязывает дочерний файбер к родительскому файберу. Закрылся родитель, ребёнок прерван. Это правильное умолчание, но иногда мешает: например, ты хочешь heartbeat, который переживёт текущий HTTP-обработчик и продолжит работать.
В 05 · Resources мы для этого использовали Effect.forkScoped: файбер привязан к scope, а не к родительскому файберу. Scope живёт сколько надо (например, всё время жизни ManagedRuntime), файбер переживёт обработчик, но всё равно подчиняется shutdown-у.
Вторая опция, Effect.forkDetach, вешает файбер на корневой файбер всей программы:
import { Effect } from 'effect';
const repeat = Effect.log('hello').pipe(
Effect.delay('300 millis'),
Effect.repeat({ times: 20 }),
);
Effect.gen(function* () {
yield* Effect.forkDetach(repeat);
yield* Effect.sleep('2 seconds');
return yield* Effect.interrupt;
}).pipe(Effect.onInterrupt(() => Effect.log('interrupted')));
// hello, hello, ... продолжает идти даже после interrupted
// (пока не закончит свои 20 повторений или не убьют процесс)
Daemon-файбер не прерывается, когда падает родитель. Он продолжает крутиться, пока не закончит сам или пока не помрёт весь процесс.
Почему это почти всегда ловушка
- отцепленный файбер не подчиняется shutdown-у твоей подсистемы;
- ресурсы, которые он держит, освобождаются непредсказуемо;
- теряется structured concurrency, ради которой ты и взял Effect;
- отлаживается мерзко: “почему процесс не падает на SIGTERM” обычно про забытый
forkDetach.
В реальном коде Effect.forkDetach нужен редко: когда ты пишешь свой ManagedRuntime поверх raw-runtime, или когда тебе действительно нужен глобальный фоновый процесс на всю программу (например, единый flusher метрик). И даже тогда чаще правильнее завести Layer с ресурсным конструктором и forkScoped поверх root-scope, чем отцеплять файбер.
Запомни таблицу:
| Что | Кому подчиняется |
|---|---|
Effect.forkChild | родительскому файберу |
Effect.forkScoped | scope (см. 05 · Resources) |
Effect.forkIn(eff, scope) | явному scope |
Effect.forkDetach | корневому файберу программы |
Умолчание для одноразовых параллельных задач, Effect.forkChild. Умолчание для фоновых процессов сервиса, Effect.forkScoped. forkDetach это escape hatch, отмеченный в голове как “если я его написал, я знаю зачем”.
Что взять с собой
Effect.forkDetachотвязывает файбер от родителя и привязывает к корневому файберу программы.- Это escape hatch, его легко написать, но он ломает structured concurrency.
- Если тебе нужен фоновый процесс, который должен пережить родителя, бери
Effect.forkScopedилиEffect.forkIn(eff, scope)с явным scope.
Раздел 7 · Effect.all и bounded concurrency
Сцена · десять задач, два слота
В Pulse у тебя массив URL, и пробивать их по одному это 10×500ms задержки. Хочется параллельно. Но не безбрежно: если URL сто, ты не хочешь стрелять в сеть сотней fetch-ей одновременно. Нужен лимит.
Effect даёт это в Effect.all с опцией concurrency. Дальше мы посмотрим на три режима, потом соберём свой all через Semaphore, чтобы понять механику, и потом сравним с тем, как Effect.all работает на практике (он чуть умнее).
Тестовая задача
import { Effect, Array as Arr } from 'effect';
const makeTask = (index: number) =>
Effect.gen(function* () {
yield* Effect.sleep('200 millis');
yield* Effect.log(`task ${index} finished`);
});
Задача спит 200мс и пишет в лог. Теперь будем гонять её через Effect.all.
Шаг 1 · Sequential, по умолчанию
const program = Effect.gen(function* () {
const tasks = Arr.makeBy(10, (i) => makeTask(i));
yield* Effect.all(tasks);
});
Effect.runPromise(program);
// task 0 finished
// (200ms)
// task 1 finished
// (200ms)
// ...
// общее время: ~2 секунды
Effect.all(tasks) без опций исполняет задачи последовательно, как обычный for-цикл. Полезно, когда между задачами есть зависимость, и бесполезно для пробинга независимых URL.
Шаг 2 · Unbounded, “запусти всё разом”
yield* Effect.all(tasks, { concurrency: 'unbounded' });
// все 10 задач параллельно
// общее время: ~200мс
Все 10 задач запускаются одновременно. Время равно времени самой медленной. Хорошо для маленьких массивов, плохо для сотен задач (сетевые соединения, file descriptors, rate-limit-ы сторонних API).
Шаг 3 · Bounded, явный лимит
yield* Effect.all(tasks, { concurrency: 2 });
// одновременно работают не больше 2 задач
// task 0, task 1 (параллельно)
// task 2, task 3 (параллельно)
// ...
// общее время: ~1 секунда
concurrency: 2, лимит на одновременно работающие задачи. Каждая закончившаяся освобождает слот для следующей.
Третьего значения у опции нет: только число и 'unbounded'. Если хочется унаследовать лимит от родителя, просто не передавай concurrency вовсе, наследование теперь и есть поведение по умолчанию.
Своя версия all через Semaphore
Чтобы понять, что внутри bounded-концерренси, соберём свой all. Без options-аргумента, с фиксированным лимитом, для тренировки.
import { Effect, Fiber, Semaphore } from 'effect';
const myAll = <A, E, R>(
effects: ReadonlyArray<Effect.Effect<A, E, R>>,
options: { readonly concurrency: 'unbounded' | number },
): Effect.Effect<Array<A>, E, R> =>
Effect.gen(function* () {
const sem =
options.concurrency === 'unbounded'
? null
: yield* Semaphore.make(options.concurrency);
const fibers: Array<Fiber.RuntimeFiber<A, E>> = [];
for (const effect of effects) {
const wrapped = sem !== null ? sem.withPermits(1)(effect) : effect;
const fiber = yield* Effect.forkChild(wrapped);
fibers.push(fiber);
}
const results: Array<A> = [];
for (const fiber of fibers) {
results.push(yield* Fiber.join(fiber));
}
return results;
});
Что мы сделали:
- Создали semaphore с N permit-ами (N = concurrency).
- На каждом эффекте обернули его в
sem.withPermits(1). Это значит: эффект перед стартом возьмёт 1 permit, на завершении вернёт. - Форкнули все эффекты сразу.
- Ждём результат каждого через
Fiber.join.
Когда permit-ы закончатся, следующие файберы будут ждать на withPermits, пока кто-то не вернёт свой permit. Это и есть bounded concurrency.
В чём настоящий Effect.all хитрее
Наша версия работает, но менее эффективна. Покажу разницу на сценарии, где одна из задач падает.
import { Effect, Array as Arr } from 'effect';
const makeTask = (index: number) =>
Effect.gen(function* () {
yield* Effect.sleep('200 millis');
if (index === 5) {
return yield* Effect.fail('five died');
}
yield* Effect.log(`task ${index} finished`);
}).pipe(Effect.onInterrupt(() => Effect.sync(() => console.log(`task ${index} interrupted`))));
Гоняем через нашу myAll(tasks, { concurrency: 5 }):
task 0 finished
task 1 finished
task 2 finished
task 3 finished
task 4 finished
task 6 interrupted
task 7 interrupted
task 8 interrupted
task 9 interrupted
Что произошло: мы форкнули все 10 файберов сразу, но 5 из них тут же ушли спать на withPermits без permit-а. Когда задача 5 упала, надо отменить остальные. Те, что ещё ждут permit, тоже запущены, и им приходит interrupt. Мы получаем 4 лишних “interrupted”.
Гоняем через Effect.all(tasks, { concurrency: 5 }):
task 0 finished
И всё. Вторые пять даже не запускались. Effect.all считает, сколько файберов реально нужно (не больше concurrency), кладёт остальные в очередь, и запускает их по мере освобождения слотов. Когда задача 5 упала, очередь просто выкидывается, никаких лишних interrupt-ов.
Преимущества родного Effect.all:
- меньше файбер-объектов;
- меньше мусора в Cause при падении (не tree из interrupt-ов всех ждущих);
- меньше нагрузки на runtime-шедулер.
Сейчас тебе достаточно знать, что наш myAll это хорошая ментальная модель (“каждой задаче надо permit”), но на проде ты всегда пишешь Effect.all(tasks, { concurrency: N }).
Effect.forEach, тот же интерфейс
Когда задачи генерируются по входу, удобнее Effect.forEach:
// ProbeError собирается в [03 · Tagged-ошибки](/l/25-effect/03-errors)
// это объединение классов NetworkError, TimeoutError, HttpStatusError из probe.
type ProbeError = NetworkError | TimeoutError | HttpStatusError;
const probeAll = (urls: ReadonlyArray<string>) =>
Effect.forEach(urls, (url) => probe(url), { concurrency: 5 });
// Effect<Array<ProbeResult>, ProbeError, ProbeRequirements>
Семантика та же, что у Effect.all, но без промежуточного urls.map(probe). Лишний раз не аллоцируешь массив эффектов.
Что взять с собой
Effect.all(tasks)без опций исполняет последовательно.Effect.all(tasks, { concurrency: 'unbounded' })запускает всё параллельно.Effect.all(tasks, { concurrency: N })ограничивает одновременно работающие задачи.- Под капотом bounded-concurrency похож на semaphore, но Effect.all держит очередь незапущенных задач, не fork-ает их про запас.
- Для коллекций удобнее
Effect.forEach(input, f, { concurrency }), чтобы не писатьinput.map(f).
Раздел 8 · Semaphore как примитив
Самостоятельное применение
Semaphore это не только хелпер для Effect.all. Это первоклассный примитив, который ты применяешь, когда нужно ограничить доступ к ресурсу или сериализовать критическую секцию.
import { Semaphore } from 'effect';
const sem = yield* Semaphore.make(N);
const protected = sem.withPermits(k)(effect);
Semaphore живёт в собственном модуле, import { Semaphore } from 'effect', и тип экземпляра зовётся Semaphore.Semaphore. Semaphore.make(N) создаёт semaphore с N permit-ами. sem.withPermits(k)(effect) оборачивает эффект так, чтобы он на старте брал k permit-ов, на завершении возвращал. Если permit-ов не хватает, эффект приостанавливается до освобождения (это suspension point, прерывание сработает корректно).
Применение 1 · Mutex для refresh-токена
В клиенте к API одна функция: получить актуальный токен. Если токен протух, его надо обновить. Беда: пять параллельных запросов одновременно увидели “протух” и одновременно запустили refresh. На сервере пять refresh-ей, четыре лишних, одна race condition.
Лечится semaphore с одним permit-ом, это и есть mutex:
import { Effect, Semaphore } from 'effect';
const makeTokenService = Effect.gen(function* () {
const refreshSem = yield* Semaphore.make(1);
const refresh = (current: string) =>
refreshSem.withPermits(1)(
Effect.gen(function* () {
yield* Effect.log('refreshing token');
const next = yield* callRefreshAPI(current);
return next;
}),
);
return { refresh };
});
Из пяти параллельных вызовов refresh permit-у достанется одному, остальные подождут. Когда первый закончит, второй увидит обновлённый токен и (если его кеширует, например, через Ref) сразу заберёт его без второго refresh. Пятеро гарантированно, один реальный сетевой вызов.
Применение 2 · Rate-limit на API-клиенте
API не любит больше 5 одновременных подключений. Отлично, заворачиваем все вызовы в semaphore с 5 permit-ами:
import { Effect, Semaphore } from 'effect';
const makeAPIClient = Effect.gen(function* () {
const apiSem = yield* Semaphore.make(5);
const get = (url: string) =>
apiSem.withPermits(1)(Effect.tryPromise(() => fetch(url)));
return { get };
});
Теперь сколько бы параллельных мест в коде ни вызывали client.get(url), шестой и далее подождут, пока один из первых пяти не завершится. Совершенно неважно, кто и где вызывает: semaphore сидит в сервисе, обтекает все точки.
Применение 3 · Защита локальной структуры
Если у тебя есть Ref или мутируемая структура в памяти, и ты делаешь над ней многошаговую операцию (прочитал, посчитал, записал обратно), используй semaphore с одним permit-ом, чтобы между read и write никто не вклинился. Альтернатива, тащить Ref-операции через Ref.update (атомарную), но не любая операция выражается одним update.
Несколько permit-ов сразу
withPermits(k) берёт k permit-ов сразу. Это полезно, когда задача “тяжёлее” других. Например, на API-клиенте лимит 5, а одна операция “большая выгрузка” сама по себе тратит слоты на трёх соединениях:
const heavyExport = apiSem.withPermits(3)(
Effect.gen(function* () {
// три запроса серией, как одна "тяжёлая" задача
}),
);
heavyExport возьмёт сразу 3 permit-а, оставив только 2 на остальной мир. Когда закончится, отдаст все три обратно. Так semaphore работает не только как “не больше N задач”, а как “общая ёмкость X, у каждой задачи свой вес”.
Что взять с собой
Semaphore.make(N)создаёт счётный semaphore с N permit-ами.sem.withPermits(k)(eff)оборачивает эффект так, чтобы он держалkpermit-ов на время исполнения.- Mutex, это semaphore с 1 permit-ом.
- Rate-limit на сервисе, это semaphore с N permit-ами вокруг каждого вызова.
withPermits(k)для k > 1 это вес задачи в общей ёмкости.
Раздел 9 · Effect.race, кто быстрее
Сцена · primary и secondary
У тебя есть основной endpoint и резервный. Хочешь: пробую оба, отдаю результат первого, кто ответил. Второй прерываю, нечего ему висеть.
Как это пишется на голом Promise
Сначала наивный вариант, тот, который чаще всего и пишут:
const slowFetch = (id: string, ms: number, signal: AbortSignal) =>
new Promise<string>((resolve) => {
console.log(`[${id}] started`);
setTimeout(() => {
console.log(`[${id}] finished`); // выполнится ВСЕГДА
resolve(`result-${id}`);
}, ms);
});
const controller = new AbortController();
const winner = await Promise.race([
slowFetch('A', 100, controller.signal),
slowFetch('B', 5000, controller.signal),
]);
controller.abort();
console.log('winner:', winner);
signal приняли в сигнатуру и забыли. Победитель отдал результат, проигравший продолжает крутиться, через 5 секунд напечатает finished. Утечка на ровном месте, потому что AbortController это не магия, это канал. Сам по себе он не отменяет ничего, его сигнал кто-то должен прочитать и среагировать.
Правильный вариант, где функция действительно реагирует на abort, выглядит так:
const slowFetch = (id: string, ms: number, signal: AbortSignal) =>
new Promise<string>((resolve, reject) => {
if (signal.aborted) {
return reject(new DOMException('Aborted', 'AbortError'));
}
console.log(`[${id}] started`);
const timerId = setTimeout(() => {
signal.removeEventListener('abort', onAbort);
console.log(`[${id}] finished`);
resolve(`result-${id}`);
}, ms);
const onAbort = () => {
clearTimeout(timerId);
console.log(`[${id}] cleanup`);
reject(new DOMException('Aborted', 'AbortError'));
};
signal.addEventListener('abort', onAbort, { once: true });
});
const controller = new AbortController();
try {
const winner = await Promise.race([
slowFetch('A', 100, controller.signal),
slowFetch('B', 5000, controller.signal),
]);
controller.abort();
console.log('winner:', winner);
} catch (e) {
if ((e as Error).name !== 'AbortError') throw e;
}
Чтобы оно работало, пришлось:
- Проверить
signal.abortedна входе (вдруг abort прилетел до старта). - Сохранить
timerId, чтобы было что чистить. - Подписаться на
abort, в обработчике погасить таймер и завершить промис ошибкойAbortError. - Удалить listener на success-пути, иначе утечка памяти, signal живёт дольше задачи.
- Помнить, что
Promise.raceне отменяет проигравших сам, дёрнутьcontroller.abort()руками после успеха. - Отличить
AbortErrorот настоящих ошибок и проглотить его вcatch.
Шесть шагов ручной обвязки на одну функцию с одним ресурсом (таймером). Дальше масштабируется так: на каждом слое абстракции тащим signal через сигнатуру; в каждой функции, которая что-то открывает (fetch, transaction, file handle), повторяем обвязку с её собственным API отмены; если cancellation-scope-ов несколько (timeout запроса плюс abort от пользователя), берём AbortSignal.any([...]) и переплетаем их руками.
То же самое через Effect.race
import { Effect } from 'effect';
const slowTask = (id: string, ms: number) =>
Effect.gen(function* () {
yield* Effect.log(`[${id}] started`);
yield* Effect.addFinalizer(() => Effect.log(`[${id}] cleanup`));
yield* Effect.sleep(`${ms} millis`);
yield* Effect.log(`[${id}] finished`);
return `result-${id}`;
}).pipe(Effect.scoped);
const program = Effect.race(slowTask('A', 100), slowTask('B', 5000));
Effect.runPromise(program).then(console.log);
// [A] started
// [B] started
// [A] finished
// [A] cleanup // scope победителя закрылся нормально
// [B] cleanup // scope проигравшего закрылся по interrupt
// result-A
Cleanup отрабатывает симметрично: и у победителя через нормальное закрытие scope, и у проигравшего через interrupt. addFinalizer это одна строка вместо listener-а плюс addEventListener плюс removeEventListener плюс проверки signal.aborted плюс ветки в catch. Дерево файберов и распространение interrupt-а сделал рантайм; типы у Effect.race не расширились новыми ошибками, потому что AbortError тут как сущности нет.
Запомни ментальный образ: AbortController относится к файберу примерно как goto к структурному программированию. Технически на goto можно написать всё, но любая чуть более сложная программа быстро превращается в обвязку.
На практике
Чаще всего ты пишешь не slowTask, а реальный запрос. Контракт ровно тот же:
import { Effect } from 'effect';
const primary = Effect.tryPromise(() => fetch('https://primary.api/health'));
const secondary = Effect.tryPromise(() => fetch('https://secondary.api/health'));
const fastest = Effect.race(primary, secondary);
// Effect<Response, UnknownException, never>
Когда один из двух файберов завершается успехом, второй автоматически получает Fiber.interrupt. Через Effect.onInterrupt или Effect.addFinalizer ты ему ещё и cleanup повесишь (например, controller.abort() на реальном fetch-е через Effect.acquireRelease, как в Pulse в разделе 10).
Effect.raceAll для массива
Если соревнующихся больше двух, бери Effect.raceAll:
const result = Effect.raceAll([primary, secondary, tertiary]);
Семантика та же: первый успех выигрывает, остальные прерываются.
Особенности и грабли
- Если первый завершившийся упал, race ещё подождёт остальных. Race останавливается на первом успехе, не на первой реакции.
- Хочешь “первый ответ, неважно успех или ошибка”, смотри
Effect.raceFirst(есть в API). - Race с
Effect.sleepэто идиома для timeout:Effect.race(work, Effect.sleep('5 seconds').pipe(Effect.flatMap(() => Effect.fail(new Timeout())))). Чище, правда, написатьEffect.timeout(work, '5 seconds'), он ровно об этом. - Если оба упали,
Effect.raceотдастCause, в массивеreasonsкоторого лежат обе ошибки. Плоская структура из 03 · Tagged-ошибки удобна ровно тут: параллельные провалы просто складываются рядом.
Что взять с собой
Effect.race(a, b)запускает два файбера, отдаёт результат первого успеха, второй прерывает.Effect.raceAll(effects)тот же шаблон для массива.- Race останавливается на первом успехе, упавших ждёт.
- Через structured concurrency проигравший автоматически получает interrupt, нет “забытых fetch-ей”.
Раздел 10 · Pulse · вклад этого урока
Параллельный пробинг массива URL
В Pulse у тебя в конфиге список URL. До этого урока цикл пробинга шёл последовательно. Теперь:
// pulse-<nick>/src/program.ts
import { Effect } from 'effect';
import { HttpService } from './services/http.ts';
import { Storage } from './services/storage.ts';
const probeOne = (url: string) =>
Effect.gen(function* () {
const http = yield* HttpService;
const storage = yield* Storage;
const controller = yield* Effect.acquireRelease(
Effect.sync(() => new AbortController()),
(c) => Effect.sync(() => c.abort()),
);
const response = yield* http.get(url, controller.signal);
yield* storage.append({ _tag: 'ProbeSuccess', url, status: response.status });
}).pipe(
Effect.scoped,
Effect.timeout('5 seconds'),
Effect.catch((err) =>
Effect.gen(function* () {
const storage = yield* Storage;
yield* storage.append({ _tag: 'ProbeFailure', url, reason: String(err) });
}),
),
Effect.onInterrupt(() => Effect.log(`probe ${url} interrupted`)),
);
export const probeAll = (urls: ReadonlyArray<string>) =>
Effect.forEach(urls, probeOne, { concurrency: 5 });
Ключевое:
Effect.forEachсconcurrency: 5, одновременно не больше пяти живых fetch-ей;Effect.scopedоборачивает один пробинг, локальные ресурсы (AbortController) живут только на время одного URL;Effect.timeout('5 seconds')страхует от висящих ответов: по истечении срока эффект падает встроенной ошибкойTimeoutError, и её подхватывает нашEffect.catchниже. Если нужен не провал, аOption, естьEffect.timeoutOption, а если своя доменная ошибка,Effect.timeoutOrElse;Effect.onInterruptпишет в журнал, что пробинг прервали (например, на SIGTERM);- падение одного пробинга не валит остальные: мы перехватили его через
Effect.catchи записали какProbeFailure.
Если хочешь “первый отвечающий из набора зеркал” вместо “пробей все”, меняешь forEach на Effect.raceAll:
const probeAny = (urls: ReadonlyArray<string>) =>
Effect.raceAll(urls.map(probeOne));
Один URL ответил, остальные прерваны, нет лишних запросов.
SIGTERM прерывает всё дерево
Из 05 · Resources у тебя в cli/watch.ts уже есть обработчик SIGTERM, который зовёт runtime.dispose(). Что происходит на нём с параллельным пробингом:
- SIGTERM прилетает в процесс.
- Обработчик зовёт
Fiber.interrupt(fiber)на главном файбере. - Через structured concurrency прерывание идёт по дереву:
Effect.forEach=> каждый активныйprobeOne=> его внутреннийhttp.get. - На каждом активном пробинге срабатывает
Effect.onInterrupt, пишет “probe X interrupted”. - Finalizer-ы локальных scope-ов закрывают
AbortController(на нём сидит реальныйfetch, прерываниеcontroller.abort()отменяет HTTP-запрос). runtime.dispose()закрывает scopeMainLive, finalizer-ыStorageдописывают буфер и закрывают handle.process.exit(0).
Без structured concurrency пришлось бы вручную носить AbortController через все промежуточные функции, помнить, что после ошибки в одном пробинге надо отменить остальные, проверять, что таймауты не срабатывают после shutdown-а. С Effect это гарантия по типу: если ты завернул работу в Effect.scoped и форкаешь через Effect.forkChild/forEach, тебе достаточно один раз нажать “стоп смены”.
Тест на отмену
В 13 · Testing и code style разберём TestClock. Здесь короткая идея проверки: запускаешь probeAll с моком HttpService, который возвращает Effect.sleep('10 seconds'), через TestClock.adjust('1 second') доводишь время до момента, прерываешь главный файбер, и ассертом проверяешь, что:
- у каждого
probeOneв журнале есть строка “interrupted”; - ни у одного нет успешного
ProbeSuccessилиProbeFailure; Storage.appendне вызывался для финального события.
Это и есть тест “никаких висящих fetch-ей после shutdown-а”. В обычном Promise-мире такой тест писать почти невозможно, потому что отмена в Promise это договор, а не контракт. В Effect это часть рантайма.
Что взять с собой
- Pulse пробивает URL параллельно через
Effect.forEach(urls, probeOne, { concurrency: N }). - Один пробинг это
Effect.scopedповерхacquireReleaseдляAbortController. - Падение одного пробинга не валит остальные, перехватываешь через
Effect.catchи пишешь в журнал. - SIGTERM через
runtime.dispose()прерывает всё дерево файберов, никаких висящих сетевых запросов.
Финал · чек-лист
ДЗ
Дальше
Следующий урок · 07 · Координация: Deferred, Queue, PubSub, Latch. Когда у тебя есть параллельные файберы, им надо обмениваться сигналами и значениями. Deferred, это однократный канал, Queue буферизованный, PubSub веерный, Latch ворота для старта-стопа. Дальше 08 · STM показывает, как делать составную атомарную операцию над несколькими Ref/Queue без mutex-ов, и 09 · Stream собирает всё это в типизированную ленту с backpressure.
Контекст по event loop удобно перечитать в 05-async/01 · Granularity. Резюме по Effect-овским каналам ошибок и по Cause в 03 · Errors. Resources, Scope и forkScoped в 05 · Resources.