Раздел 25 · Effect-TS

Координация: Deferred, Queue, PubSub, Semaphore, Latch

senior~130 мин

открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти

Координация: Deferred, Queue, PubSub, Semaphore, Latch

Сцена · очередь у кассы

Магазин на проходной, один кассир, один покупатель. Если кассир медленнее, чем приходят покупатели, очередь растёт. Что делать с лишними? Развернуть назад (bounded), пускать всех (unbounded, очередь до потолка), выгонять самых старых (sliding) или самых новых (dropping). Это не философия, а дизайн-решения, которые ты принимаешь, когда заводишь канал между producer-ом и consumer-ом.

Backpressure: producer 12/s, consumer 5/s, hwm = 8
producer
flowing
buffer
0/8
consumer
0
0 produced0 consumed0 paused0 tps

Подними producer выше consumer -- буфер заполняется. Когда длина достигает hwm, producer переходит в paused: это и есть backpressure. Опустится буфер ниже половины hwm -- producer возвращается в flowing. Это упрощённая модель: в Node hwm, это порог для возврата false у write(), в Web Streams похожая логика реализована через возврат промиса из pull.

Виджет показывает один из четырёх режимов: bounded c автопаузой, когда буфер достиг hwm (high water mark, верхняя граница буфера). Подними producer над consumer, и увидишь, как буфер заполняется, producer уходит в paused, поток сжимается до скорости consumer-а. Это и есть backpressure: обратная связь идёт от потребителя к производителю через размер буфера. Без неё единственный сценарий, который ты получишь, это OOM на проде через две недели.

Карта урока · что заберёшь домой

Шесть координационных примитивов в Effect:

  • Deferred<A, E>, однократный канал. Producer один раз положит значение, любое число consumer-ов прочитают. Аналог Promise, только под Effect-овским рантаймом и с возможностью прервать ждущих.
  • Queue<A, E>, многоразовый буфер с одной из четырёх стратегий. Любое количество producer-ов и consumer-ов, FIFO-семантика, выбор политики на переполнение.
  • Queue<A, E> с завершением. Второй параметр это канал ошибки, как у Effect: producer может явно закрыть поток (end) или провалить его типизированной ошибкой (fail), и consumer об этом узнает.
  • PubSub<A>, fan-out. Один publish, N независимых subscriber-ов, каждый получает все опубликованные сообщения с момента подписки (плюс короткий хвост через replay).
  • Semaphore, лимит ёмкости. Под капотом счётчик permit-ов, на практике mutex (1 permit), rate-limit (N permit-ов), per-resource ограничение. Базовый сценарий уже видел в 06 · Файберы и concurrency, здесь добавим per-domain паттерн.
  • Latch, ворота. Все ждущие файберы стоят перед latch-ом, на latch.open проходят одновременно. Аналог CountDownLatch из Java, но с open/close и без счётчика.

К концу урока в Pulse появится внутренний bus: UrlQueue, MonitorEvents, DomainLimiter, bootstrapLatch. Pulse станет похож на маленький распределённый стенд внутри одного процесса, и через эту координацию подведём к следующему уроку про транзакции.

Раздел 1 · Deferred, одноразовый канал

Сцена · все ждут одного и того же

В Pulse есть loadConfig: читает YAML с диска, парсит, валидирует через Schema. Десять воркеров стартуют одновременно, каждому нужен конфиг. Что нельзя: десять одновременных loadConfig, десять чтений диска, десять одинаковых ошибок при невалидном файле. Что нужно: один реальный loadConfig, остальные ждут результат, получают его все сразу.

В Promise-мире это пишут через “lazy singleton”:

let configPromise: Promise<PulseConfig> | null = null;
const getConfig = () => {
  if (configPromise === null) configPromise = loadConfig();
  return configPromise;
};

Работает, но есть две беды. Промежуточное состояние “идёт загрузка” не отличить от “загружено”, и ловить ошибку приходится у каждого, кто вызвал, отдельно. И в Promise нельзя отменить ждущих: даже если решил, что конфиг не нужен, все десять await getConfig() будут висеть до конца.

Effect для этого даёт Deferred: одноразовый канал, который один раз получает значение и потом раздаёт его всем подписчикам.

Deferred.make, await, succeed

import { Deferred, Effect } from 'effect';

const program = Effect.gen(function* () {
  const deferred = yield* Deferred.make<number, never>();
  // deferred: Deferred<number, never>, пока без значения

  const reader = Effect.gen(function* () {
    yield* Effect.log('reader: waiting');
    const value = yield* Deferred.await(deferred);
    yield* Effect.log(`reader: got ${value}`);
  });

  const writer = Effect.gen(function* () {
    yield* Effect.sleep('1 second');
    yield* Deferred.succeed(deferred, 42);
    yield* Effect.log('writer: published');
  });

  yield* Effect.all([reader, writer], { concurrency: 'unbounded' });
});

Effect.runPromise(program);
// reader: waiting
// (1 секунда)
// writer: published
// reader: got 42

Deferred.make<A, E>() создаёт пустой канал. Deferred.await(d) приостанавливает текущий файбер, пока канал не получит значение. Deferred.succeed(d, value) кладёт значение, один раз: повторные succeed или fail возвращают false и ничего не делают.

Если ждут несколько файберов, после succeed все просыпаются и получают одно и то же значение. Это и есть свойство, которое отличает Deferred от mutex-а или ad-hoc-Promise: один источник истины раздаётся веером, с типизированной ошибкой в канале E, и любого ждущего можно прервать без оставленного мусора.

Deferred.fail и Deferred.complete

Deferred<A, E> ровно так же умеет провал. Заверши через Deferred.fail(d, error), и все ждущие проснутся с ошибкой E:

import { Data, Deferred, Effect } from 'effect';

class ConfigError extends Data.TaggedError('ConfigError')<{ reason: string }> {}

const program = Effect.gen(function* () {
  const deferred = yield* Deferred.make<PulseConfig, ConfigError>();

  yield* Effect.forkChild(
    Effect.gen(function* () {
      yield* Effect.sleep('500 millis');
      yield* Deferred.fail(deferred, new ConfigError({ reason: 'parse error' }));
    }),
  );

  const result = yield* Deferred.await(deferred);
  // если writer вызвал fail, await поднимет ConfigError в канал E
  return result;
});
// program: Effect<PulseConfig, ConfigError, never>

Тип Deferred<PulseConfig, ConfigError> точно говорит: успех вернёт PulseConfig, неудача ConfigError. У ждущих файберов канал E расширяется ровно на ConfigError, без типа unknown или Error без формы.

Часто хочется не “положить значение”, а “положить результат эффекта”. Для этого есть Deferred.complete(d, effect):

const writer = Deferred.complete(deferred, loadConfig);
// loadConfig: Effect<PulseConfig, ConfigError>
// если упал, deferred получит fail; если успех, succeed

complete запустит эффект и поместит его Exit в deferred. Это нужная форма для bootstrap: ты не знаешь заранее, успех или провал, и не хочешь дублировать try-логику.

Bootstrap через Deferred

Возвращаемся к десяти воркерам и одному loadConfig. Делаем сервис:

import { Context, Deferred, Effect, Layer } from 'effect';

class Config extends Context.Service<Config>()('Pulse/Config', {
  make: Effect.gen(function* () {
    const deferred = yield* Deferred.make<PulseConfig, ConfigError>();

    // один раз запустим загрузку в фоне
    yield* Effect.forkChild(Deferred.complete(deferred, loadConfig));

    return {
      getConfig: Deferred.await(deferred),
    };
  }),
}) {
  static readonly layer = Layer.effect(this, this.make);
}

Десять воркеров вызывают yield* config.getConfig, и все получают один и тот же результат от единственной загрузки. Если упало, все десять синхронно увидят ConfigError, можно ретраить на уровне сервиса.

Если завернуть loadConfig через Effect.cached (см. 11 · Runtime и Schedule), получится тот же эффект с TTL: значение обновляется по таймеру, читатели не блокируются. Deferred это атомарный кирпичик, на котором cached и подобные кэши построены.

Прерывание ждущего

Deferred.await это suspension point (см. 06 · Файберы и concurrency). Если файбер, который сидит в await, получает Fiber.interrupt, он аккуратно завершается, не оставляя за собой состояния. Сам Deferred это просто значение в памяти, никаких finalizer-ов он не требует.

Что взять с собой

  • Deferred<A, E>, одноразовый канал. Один раз succeed/fail/complete, любое число await-ов.
  • Все ждущие просыпаются на одном значении, последующие await получают его без ожидания.
  • Deferred.complete(d, effect) берёт эффект и кладёт его Exit в канал.
  • Используй для bootstrap, кешированных вычислений, сигналов готовности.
  • await это suspension point, прерывается чисто.

Раздел 2 · Queue и четыре стратегии backpressure

Что такое Effect Queue

Queue<A> это многоразовый FIFO-буфер. У него два главных метода: Queue.offer(q, value) положить, Queue.take(q) взять. Любое количество файбер-producer-ов и файбер-consumer-ов работают с одной очередью одновременно, безопасность по типу гарантирована рантаймом.

Создаётся через одну из четырёх фабрик:

ФабрикаПоведение на полной очередиЧто использовать
Queue.bounded(N)offer приостанавливается, ждёт местаproducer быстрый, нельзя терять
Queue.unbounded()никогда не полна, растёт без пределамалый фиксированный объём, доверие к источнику
Queue.sliding(N)выкидывает самый старый элемент, offer не блокируетсяхочется свежие данные, старые не нужны (метрики, последние события UI)
Queue.dropping(N)выкидывает самый новый (offer возвращает false), не блокируетсяважно сохранить уже принятое, перегруз отбросить

Оба метода отдают Effect-ы:

const offer = Queue.offer(q, value); // Effect<boolean, never, never>
const take = Queue.take(q);          // Effect<A, never, never>

take приостанавливает файбер, если очередь пуста, и просыпается, когда offer положит элемент. Queue.shutdown(q) закрывает очередь: все ждущие take-и просыпаются с прерыванием, последующие offer/take тоже завершаются прерыванием. Это атомарный finalizer, через Effect.acquireRelease его удобно повесить на scope сервиса.

Стратегия 1 · bounded, базовый случай

import { Effect, Queue } from 'effect';

const program = Effect.gen(function* () {
  const queue = yield* Queue.bounded<number>(2);

  const producer = Effect.gen(function* () {
    for (let index = 1; index <= 5; index += 1) {
      yield* Effect.log(`offer ${index}`);
      yield* Queue.offer(queue, index);
    }
    yield* Effect.log('producer done');
  });

  const consumer = Effect.gen(function* () {
    for (let index = 0; index < 5; index += 1) {
      yield* Effect.sleep('500 millis');
      const value = yield* Queue.take(queue);
      yield* Effect.log(`take ${value}`);
    }
  });

  yield* Effect.all([producer, consumer], { concurrency: 'unbounded' });
});

Effect.runPromise(program);

Что увидишь в журнале:

offer 1
offer 2
offer 3                 <- producer ушёл в suspended, очередь заполнена
(500 ms, consumer проснулся)
take 1
offer 3                 <- producer проснулся, положил
offer 4                 <- снова suspended
(500 ms)
take 2
offer 4
...

offer на полной очереди не возвращает false, не теряет элемент, не падает. Он приостанавливает файбер-producer до освобождения слота. Это ключевая разница с буферами в чистом JavaScript (Array.push всегда сразу): backpressure встроен в семантику метода, ничего вручную не разруливать.

Стратегия 2 · unbounded, никогда не блокирует, иногда падает по памяти

const queue = yield* Queue.unbounded<number>();
yield* Queue.offer(queue, 1); // никогда не блокирует

Полезен в одном случае: ты уверен, что общее число элементов ограничено внешним фактором (например, в пайплайн поступают только сообщения с фиксированного входа, чьё число не превышает N). Если эта уверенность ошибочна, рано или поздно процесс падает в OOM.

В реальном Pulse unbounded почти не встречается: всегда есть hwm, всегда есть смысл в bounded.

Стратегия 3 · sliding, выкидываем старое

Сценарий: dashboard показывает последние 100 событий. Если их прилетает 1000, нужны последние 100, а не первые. Это sliding:

const queue = yield* Queue.sliding<Event>(100);

yield* Queue.offer(queue, event); // если очередь полна, выкидывает самый старый

offer всегда возвращает true, никогда не блокирует. Старые элементы тихо теряются, рантайм их не сохраняет. Подходит для метрик, телеметрии, последних статусов: тебя интересует “сейчас”, не “история”.

Стратегия 4 · dropping, выкидываем новое

const queue = yield* Queue.dropping<Job>(50);

const accepted = yield* Queue.offer(queue, job); // false, если очередь полна

Здесь, наоборот, на переполнении новый элемент не принимается, offer возвращает false. Это полезно, когда уже принятые задания обязательно должны быть обработаны, а лишние можно просто отбросить (с логом или 429 Too Many Requests клиенту).

Сценарий “rate-limit на входной API”: если очередь job-ов забита, новые запросы получают 429, мы не ставим их в очередь “на потом”, чтобы не накапливать долг.

Сравнение в одной таблице

producer:  □□□□□□□□□□□□□□□□□  (быстро)
consumer:  □ □ □ □ □ □ □ □ □  (медленно)

bounded(2):    producer ждёт места, consumer всегда видит непустую
                ничего не теряем, поток замедляется до consumer

unbounded:     producer не ждёт, очередь растёт
                ничего не теряем, рискуем по памяти

sliding(2):    producer не ждёт, старые выкидываются
                видим "хвост" событий, начало пропало

dropping(2):   producer не ждёт, новые отбрасываются
                видим "голову" событий, конец пропал

Выбор стратегии это дизайн, не оптимизация. Сначала ответь “что для меня важнее, не терять или не блокировать”, потом подставь фабрику.

Паттерн · producer и consumer как файберы

Полный шаблон:

import { Effect, Queue } from 'effect';

const makeProcessor = Effect.gen(function* () {
  const queue = yield* Queue.bounded<Event>(64);

  const consumer = Effect.gen(function* () {
    while (true) {
      const event = yield* Queue.take(queue);
      yield* processEvent(event); // обработка
    }
  }).pipe(Effect.forever);

  // запускаем consumer на scope сервиса
  yield* Effect.forkScoped(consumer);

  return {
    submit: (event: Event) => Queue.offer(queue, event),
  };
});

Сервис отдаёт наружу один метод submit. Внутри файбер-consumer крутит take в цикле, на scope сервиса, чтобы при runtime.dispose() он остановился (см. 05 · Resources).

Если consumer-ов нужно несколько (например, четыре worker-а на bounded(64)):

yield* Effect.forEach(
  Array.from({ length: 4 }, (_, i) => i),
  () => Effect.forkScoped(consumer),
  { concurrency: 'unbounded' },
);

Четыре файбера конкурируют на одном take, каждое сообщение получит ровно один из них. Это и есть простой work-stealing pool на четыре файбера.

Грабли · take без shutdown

Бесконечный цикл Queue.take сам никогда не завершится, ему нужен внешний сигнал. Без Effect.forkScoped (или другого scope-bound fork) этот цикл переживёт прерывание родителя только если ты случайно завернул его в forkDetach. Через scope ты получаешь автоматическое: на закрытии scope take прервётся, цикл завершится.

Queue.shutdown(q) тоже разбудит ждущих, и take отдаст Cause.interrupt. Какой механизм брать, решает уровень: для подсистемного завершения хватает scope, для именно очереди (например, “хватит принимать новые события, доработаем имеющиеся”) берёшь Queue.shutdown.

Пачкой, сразу несколько · takeAll, takeUpTo, takeBetween

Кроме take (одно сообщение, блокируется до прихода) у очереди есть три батч-варианта:

const all = yield* Queue.takeAll(queue);
// NonEmptyArray<A>: всё, что лежит. Ждёт, пока не появится хотя бы один элемент.

const someOf = yield* Queue.takeN(queue, 10);
// Array<A>, ровно до 10 элементов.

const between = yield* Queue.takeBetween(queue, 5, 20);
// Array<A>, ждёт, пока не накопится хотя бы 5, потом отдаст всё доступное вплоть до 20.

const rest = yield* Queue.collect(queue);
// Array<A>: всё до самого завершения очереди.

Все они отдают обычные массивы, а не Chunk: в v4 Chunk почти целиком ушёл из пользовательского API. Stream.fromQueue под капотом крутит именно батч-чтение: минимум один элемент ждёт, дальше забирает столько, сколько есть. Из этого получается естественный backpressure в стриме: следующий “вытащить” происходит только после того, как текущий чанк прошёл через mapEffect.

unsafeOffer, синхронный путь

Если ты в callback-коде или просто не хочешь блокировать файбер на полной bounded-очереди, есть queue.unsafeOffer(value). Это не эффект, а обычная функция:

const accepted = Queue.offerUnsafe(queue, 42); // boolean
// true: положили
// false: bounded переполнено или dropping отбросил

Обрати внимание на форму имени. В v4 приставку unsafe перенесли в суффикс: offerUnsafe, endUnsafe, makeUnsafe, sizeUnsafe. Правило одно на всю библиотеку, поэтому небезопасный вариант всегда видно в конце имени, рядом друг с другом в автодополнении.

Никакой backpressure, никаких ожиданий. Подходит, когда нужно прокинуть Effect-овский Queue в синхронный API (например, в обработчик event-emitter-а из не-Effect-кода) и потеря на переполнении это допустимое поведение. Для bounded-очереди с гарантией доставки используй обычный Queue.offer под Effect-рантаймом.

Effect.forever против while(true)

В примере выше consumer ходил через while (true). Эквивалент через комбинатор:

const consumer = Queue.take(queue).pipe(
  Effect.tap(processEvent),
  Effect.forever,
);

Разница тонкая, но важная: Effect.forever между итерациями вставляет Effect.yieldNow, явно отдавая управление планировщику Effect. В while (true) ты отвечаешь за yield-точки сам. Внутри Queue.take уже сидит suspension, так что обычный consumer-цикл работает корректно. Но если кто-нибудь потом добавит в тот же цикл чисто синхронную обработку (например, Console.log без Effect.log), while (true) может забить файбер в горячем цикле, а Effect.forever страхует.

Для job-processor-ов и долгоживущих consumer-ов предпочитай Effect.forever: меньше вероятность отстрелить себе ногу позже.

Что взять с собой · Queue

  • Queue.bounded(N), backpressure, ничего не теряем, producer ждёт.
  • Queue.unbounded(), никогда не блокирует, может разорвать память.
  • Queue.sliding(N), держит последние N, старые выкидываются.
  • Queue.dropping(N), держит первые принятые N, новые отбрасываются.
  • Стратегия это дизайн: что важнее, не терять или не блокировать.
  • takeAll/takeN/takeBetween/collect забирают пачкой, не по одному, и отдают обычные массивы; Stream.fromQueue так и работает под капотом.
  • Queue.offerUnsafe это синхронный путь без backpressure, для мостов в не-Effect-API. Суффикс Unsafe в v4 единый на всю библиотеку.
  • Effect.forever гарантирует yield-точку между итерациями, while (true) доверяет твоему циклу.
  • Consumer-ы запускают через Effect.forkScoped, чтобы scope управлял их жизнью.
  • Queue.shutdown разбудит ждущих, take отдаст interrupt.

Раздел 3 · Завершение очереди: end, fail, Cause.Done

Сцена · producer однажды заканчивает

Queue отлично подходит для бесконечных потоков: event bus, цикл пробинга в Pulse, входящий API-трафик. Никто не “заканчивает”. Но бывают другие сценарии. Producer дочитал файл и сообщает “всё, больше не будет”. Producer упал на parse-error и хочет передать причину consumer-у через тот же канал, а не через побочный Deferred<Cause> или ad-hoc-Ref<Option<Error>>.

Раньше это решалось костылями: отдельный Ref<boolean> “done”, sentinel-значение в очереди, Queue.shutdown плюс ловля interrupt. Получалась рассыпуха, и легко было забыть один из путей. В v4 у очереди появился второй параметр типа, и вся эта история закрывается ей самой.

Queue<A, E>, второй параметр

import { Cause, Effect, Queue } from 'effect';

const program = Effect.gen(function* () {
  // без завершения: E = never, обычная бесконечная очередь
  const plain = yield* Queue.unbounded<number>();

  // с нормальным завершением: в канал ошибки кладём Cause.Done
  const finite = yield* Queue.bounded<number, Cause.Done>(10);

  // с собственной ошибкой
  const withErrors = yield* Queue.unbounded<number, ParseError | Cause.Done>();
});

Второй параметр это канал ошибки, ровно как E у Effect. Внутри него живёт и твоя доменная ошибка, и специальный маркер Cause.Done, которым рантайм обозначает нормальное завершение. Читается это буквально: “очередь чисел, которая может закончиться” и “очередь чисел, которая может закончиться или сломаться разбором”.

Все четыре стратегии на месте, плюс общая фабрика с опциями:

const suspending = yield* Queue.bounded<number, Cause.Done>(10);
const dropping = yield* Queue.dropping<number, Cause.Done>(10);
const sliding = yield* Queue.sliding<number, Cause.Done>(10);
const custom = yield* Queue.make<number, Cause.Done>({ capacity: 10, strategy: 'sliding' });

offer, offerAll, take ведут себя как раньше:

yield* Queue.offer(queue, 1);
yield* Queue.offerAll(queue, [2, 3, 4]);
const value = yield* Queue.take(queue);

end, нормальное завершение

const program = Effect.gen(function* () {
  const queue = yield* Queue.bounded<number, Cause.Done>(10);

  yield* Queue.offer(queue, 1);
  yield* Queue.offer(queue, 2);
  yield* Queue.end(queue); // сигнал "больше не будет"

  const accepted = yield* Queue.offer(queue, 3);
  // accepted: false, после end новые offer-ы отбрасываются

  const first = yield* Queue.take(queue); // 1
  const second = yield* Queue.take(queue); // 2
  // третий take упадёт с Cause.Done: буфер разобран, очередь закрыта
});

Сообщения, положенные до end, не теряются. Consumer сначала получит всё, что было до сигнала, и только потом наткнётся на завершение. Само завершение приходит по каналу ошибки, поэтому его нельзя случайно пропустить: компилятор помнит про Cause.Done в типе.

collect, забрать всё до конца

Раз завершение живёт в типе, у очереди есть операция “читай, пока не закончится”:

const all = yield* Queue.collect(queue);
// Array<number>, всё до завершения. Cause.Done уже вычтен из канала ошибки.

Это и есть consumer-цикл в одну строку, без внешних флагов и без ручной ловли done. Если очередь завершилась не через end, а через fail, collect поднимет твою ошибку.

fail, типизированная ошибка через канал

import { Cause, Data, Effect, Queue } from 'effect';

class ProcessingError extends Data.TaggedError('ProcessingError')<{
  readonly reason: string;
}> {}

const program = Effect.gen(function* () {
  const queue = yield* Queue.unbounded<number, ProcessingError | Cause.Done>();

  yield* Queue.offer(queue, 1);
  yield* Queue.offer(queue, 2);
  yield* Queue.fail(queue, new ProcessingError({ reason: 'parser blew up' }));

  // буфер забирается раньше ошибки
  const first = yield* Queue.take(queue); // 1
  const second = yield* Queue.take(queue); // 2

  // следующий take упадёт с типизированной ошибкой
  const error = yield* Queue.take(queue).pipe(Effect.flip);
  // error: ProcessingError { reason: 'parser blew up' }
});

Семантика “что лежит, то отдадим” сохраняется. Producer положил два сообщения, потом упал, consumer обработает оба, и только следующий take поднимет ошибку в канал E. Это правильное поведение для пайплайнов: то, что уже принято, не теряется, но провал виден явно и с понятным типом.

Для полного Cause (defect, несколько причин сразу, прерывание) есть failCause:

yield* Queue.failCause(queue, Cause.die('fatal'));

Мост в Stream

Stream.fromQueue(queue) понимает завершение сам:

import { Cause, Effect, Queue, Stream } from 'effect';

const stream = Stream.fromQueue(queue);
// Stream<A, E>, завершится на Queue.end, упадёт на Queue.fail

Заметь сигнатуру: из канала ошибки потока Cause.Done вычтен. Для очереди это ошибка, для потока это просто конец, и типы это отражают. Обратный мост, Stream.toQueue(stream), отдаёт scoped-очередь, которую наполняет поток.

Job-processor на этом пишется так:

import { Cause, Console, Context, Data, Effect, Layer, Queue, Stream } from 'effect';

class JobProcessorError extends Data.TaggedError('JobProcessorError')<{
  readonly reason: string;
}> {}

class JobProcessor extends Context.Service<JobProcessor>()('Pulse/JobProcessor', {
  make: Effect.gen(function* () {
    const queue = yield* Queue.unbounded<number, JobProcessorError | Cause.Done>();

    yield* Effect.forkScoped(
      Stream.fromQueue(queue).pipe(
        Stream.mapEffect((value) => Console.log('processing', value), {
          concurrency: 2,
        }),
        Stream.runDrain,
      ),
    );

    return {
      submit: (value: number) => Queue.offer(queue, value),
      end: Queue.end(queue),
      fail: (reason: string) => Queue.fail(queue, new JobProcessorError({ reason })),
    };
  }),
}) {
  static readonly layer = Layer.effect(this, this.make);
}

Три ручки наружу: submit, end, fail. Типизированное завершение, никакого “shutdown это interrupt”.

Unsafe-варианты для callback-кода

Для мостов в не-Effect-API есть синхронные варианты, все с суффиксом Unsafe:

Queue.offerUnsafe(queue, value);        // boolean
Queue.offerAllUnsafe(queue, [1, 2, 3]); // оставшиеся элементы
Queue.endUnsafe(queue);                 // синхронный end
Queue.failCauseUnsafe(queue, cause);    // синхронный fail

end против shutdown

Две операции закрывают очередь, но по-разному, и путать их не стоит:

СвойствоQueue.end(q)Queue.shutdown(q)
Что видит consumerсначала весь буфер, потом Cause.Doneнемедленный interrupt
Уже принятые сообщениядочитываютсятеряются
Видно ли в типеда, Cause.Done в Eнет
Когда братьproducer закончил работу штатноподсистему гасят, дочитывать нечего

Когда брать что:

  • Поток без естественного конца (event bus, бесконечный probe-цикл), обычная Queue<A> без второго параметра, закрытие через scope.
  • Поток с явным завершением (читаем файл, обрабатываем batch, агент закрывает сессию), Queue<A, Cause.Done> плюс end.
  • Producer может упасть с типизированной ошибкой, и эту ошибку должен увидеть consumer, Queue<A, MyError | Cause.Done> плюс fail.

Что взять с собой · завершение очереди

  • У Queue<A, E> в v4 есть канал ошибки, и завершение живёт прямо в типе.
  • Queue.end(q) закрывает очередь штатно, consumer дочитывает буфер и получает Cause.Done.
  • Queue.fail(q, error) и Queue.failCause(q, cause) закрывают её с причиной, тоже после разбора буфера.
  • Queue.collect(q) забирает всё до завершения одной операцией, без внешних флагов.
  • Stream.fromQueue вычитает Cause.Done из канала ошибки: для потока это просто конец.
  • end дочитывает буфер, shutdown прерывает ждущих. Разные инструменты, разные ситуации.

Раздел 4 · PubSub, fan-out

Чем PubSub отличается от Queue

Queue.take это work-stealing: каждое сообщение получает ровно один consumer. Это правильно для пула обработчиков, но плохо, если хочешь, чтобы каждый подписчик увидел каждое сообщение. Сценарий: один источник MonitorEvent, три независимых потребителя, журналирующий, отправляющий в Slack, индексирующий в SQLite. Все три должны увидеть все события.

Аналогия из радиоэфира: один диктор говорит в микрофон, тысячи приёмников слышат то же самое. У Queue сообщение получает один курьер из пула, у PubSub сигнал слышит каждый, кто настроен на волну.

В Queue это решается тремя очередями и duplicate-логикой в публикаторе. Эту работу за тебя делает PubSub:

import { Effect, PubSub } from 'effect';

const program = Effect.gen(function* () {
  const pubsub = yield* PubSub.bounded<string>(8);

  const subscriber = (name: string) =>
    Effect.gen(function* () {
      const sub = yield* PubSub.subscribe(pubsub);
      // sub: PubSub.Subscription<string>, личный канал этого подписчика
      while (true) {
        const message = yield* PubSub.take(sub);
        yield* Effect.log(`${name}: ${message}`);
      }
    }).pipe(Effect.scoped, Effect.forever);

  // три параллельных подписчика
  yield* Effect.all([
    Effect.forkScoped(subscriber('logger')),
    Effect.forkScoped(subscriber('slack')),
    Effect.forkScoped(subscriber('indexer')),
  ]);

  yield* Effect.sleep('100 millis'); // дать подписчикам подключиться
  yield* PubSub.publish(pubsub, 'hello');
  yield* PubSub.publish(pubsub, 'world');
});

Что увидишь:

logger: hello
slack: hello
indexer: hello
logger: world
slack: world
indexer: world

Каждое из двух сообщений приходит трём подписчикам. Внутри PubSub это реализовано через отдельный буфер на каждого подписчика: при subscribe заводится новая подписка, на publish сообщение копируется в каждую активную.

Обрати внимание на тип того, что вернул subscribe. Это не Queue, а собственный тип PubSub.Subscription<A>, и читают из него функциями того же модуля: PubSub.take(sub), PubSub.takeAll(sub), PubSub.takeUpTo(sub, n), PubSub.takeBetween(sub, min, max). Разделение честное: очередь и подписка ведут себя по-разному (у подписки, например, есть окно реплея), и раньше общий тип это скрывал.

publishAll, batch-публикация

Если у тебя пачка сообщений, не зови publish в цикле. Есть publishAll, который кладёт сразу всё:

yield* PubSub.publishAll(pubsub, [1, 2, 3, 4, 5]);
// Effect<boolean>, true если все сообщения приняты всеми подписчиками

Поведение по стратегиям такое же, как у одиночного publish: на bounded блокируется на самом медленном подписчике, на sliding теснит старое в очередях у медленных, на dropping лишнее теряется. Преимущество перед циклом одно: меньше файбер-точек переключения, меньше работы scheduler-у.

Scope подписки

subscribe это scoped-эффект (см. 05 · Resources), его сигнатура Effect<PubSub.Subscription<A>, never, Scope>. Когда scope подписки закрывается, рантайм отписывает подписчика и освобождает его буфер. В нашем коде Effect.scoped оборачивает один subscriber-цикл, так что при interrupt-е файбера scope закроется, очередь освободится, последующие publish уже не будут пытаться туда писать.

Это критично для долгоживущих систем: подписчики приходят и уходят, без scope-а ты получишь утечку буферов в PubSub.

Что значит “bounded” у PubSub

У PubSub.bounded(N) есть тонкая семантика: publish приостанавливается, пока самый медленный подписчик не освободит место в своей очереди. Это классическая ловушка fan-out систем: один тормозящий подписчик блокирует всех. На практике для журналов и метрик обычно берут PubSub.sliding(N) (медленный подписчик пропускает старые сообщения, остальные не страдают) или PubSub.unbounded() с понятным риском по памяти.

const pubsub = yield* PubSub.sliding<MonitorEvent>(256);
// publish никогда не блокирует, медленный subscriber теряет старые

Решение зависит от семантики данных. Для финансовых событий sliding неприемлем, нужен bounded плюс мониторинг отстающих. Для UI-обновлений sliding это ровно то, что надо.

Изоляция падений

Если один подписчик упал (его файбер завершился с ошибкой), это не влияет на остальных. PubSub держит только их активные очереди, ошибка на стороне одного потребителя локальна:

const flaky = Effect.gen(function* () {
  const sub = yield* PubSub.subscribe(pubsub);
  while (true) {
    const message = yield* PubSub.take(sub);
    if (message === 'crash') return yield* Effect.fail('boom');
    yield* Effect.log(`flaky: ${message}`);
  }
}).pipe(Effect.scoped);

При publish('crash') flaky упадёт и его файбер завершится. Через scope подписка снимется, её буфер освободится. Остальные подписчики продолжают получать сообщения как ни в чём не бывало.

Это другая важная разница с Queue.take пулом: там ошибка одного файбера оставляет несбалансированный пул и некому обрабатывать сообщения, которые висели в полёте. В PubSub каждый подписчик независим.

Replay · последние N для опоздавших

По умолчанию PubSub не помнит истории: подписчик видит только сообщения, опубликованные после его subscribe. Подписался поздно, всё, что было до этого, проехало мимо.

Иногда это неудобно. Открыл UI-клиент, хочется сразу увидеть последние 10 статусов, а не пустой экран до следующего тика. Подключился к чату, хочется прочитать последние реплики. Для этого PubSub принимает опцию replay:

const pubsub = yield* PubSub.unbounded<number>({ replay: 3 });

yield* PubSub.publishAll(pubsub, [1, 2, 3, 4, 5]);

// поздний подписчик
const subscriber = yield* PubSub.subscribe(pubsub);
const recent = yield* PubSub.takeAll(subscriber);
// recent: [3, 4, 5], последние 3 из реплея

takeAll отдаёт обычный массив, а не Chunk: в v4 из пользовательского API Chunk почти везде ушёл в пользу нативных массивов.

Внутри PubSub держит sliding-буфер из последних N сообщений. На каждый новый subscribe рантайм сначала наливает в новую подписку эти N (или меньше, если опубликовано меньше N), и только потом начинает доставлять “живые” сообщения.

replay сочетается со всеми стратегиями:

PubSub.bounded<Event>({ capacity: 256, replay: 10 });
PubSub.sliding<Event>({ capacity: 256, replay: 10 });

Тонкость: replay: N это минимальная гарантия, а не максимальный лимит. Если у тебя unbounded и в системе осталось больше N сообщений, поздний подписчик получит их все. Чтобы строго ограничить хвост, бери bounded или sliding-стратегию, тогда старое физически вытесняется.

В Pulse это пригодится для UI-клиента, который подключается через WebSocket: первое, что он видит, это хвост последних MonitorEvent-ов, чтобы интерфейс не загружался пустым.

Когда Queue, когда PubSub

СценарийЧто брать
Несколько обработчиков, каждое сообщение должен взять ровно одинQueue плюс пул consumer-ов
Одно сообщение должен увидеть каждый подписчикPubSub
Один producer, один consumer, backpressureQueue.bounded
Один producer, неизвестное число подписчиков, появляются/уходят на летуPubSub.sliding или bounded плюс мониторинг
Лента событий, новые подписчики должны видеть последние N сообщенийPubSub.sliding(N) или PubSub.bounded({ capacity, replay })

По умолчанию PubSub не хранит истории: подписчик видит только сообщения, опубликованные после его subscribe. Если нужен короткий хвост, передавай replay: N (см. выше). Для сложных стратегий (фильтрация, замешивание реплея с другим источником, длинная история) бери Stream, см. 09 · Stream.

Мост в Stream

PubSub это примитив, на котором Stream строит свой fan-out:

import { Effect, PubSub, Stream } from 'effect';

// PubSub -> Stream, любой подписчик
const stream = Stream.fromPubSub(pubsub);

// Stream -> PubSub, источник в bus
const pubsub = yield* Stream.toPubSub(myStream, 10);

Есть и третий мост, Stream.fromSubscription(sub): он делает поток из уже открытой подписки, когда subscribe ты сделал сам и хочешь дочитать её как Stream. В уроке 09 · Stream этим воспользуются Stream.broadcast (раздать поток нескольким подписчикам) и Stream.share (ленивый старт плюс reference counting). Оба под капотом крутят PubSub ровно по тому, что ты уже видел в этом разделе.

Что взять с собой

  • PubSub<A>, fan-out: каждый подписчик получает каждое сообщение.
  • PubSub.publishAll для пачки сообщений, не зови publish в цикле.
  • subscribe это scoped-эффект, отписка происходит на закрытии scope.
  • Стратегия bounded блокирует publish на самом медленном подписчике, осторожно.
  • sliding, для метрик и UI; bounded, когда терять нельзя.
  • Падение одного подписчика не валит остальных, изоляция по подписке.
  • По умолчанию истории нет, но replay: N отдаст опоздавшим короткий хвост.
  • Stream.fromPubSub, Stream.fromSubscription, Stream.toPubSub это мосты в Stream, на которых построены broadcast и share.

Раздел 5 · Semaphore: per-domain лимит

Короткое напоминание

В 06 · Файберы и concurrency Semaphore.make(N) уже разобрали как mutex (1 permit) и rate-limit на сервисе (N permit-ов). Здесь возьмём один новый паттерн, который не помещается в “один semaphore на всё”.

Сцена · разные домены, разные лимиты

В Pulse у тебя 100 URL, и они принадлежат разным доменам. Один общий semaphore с 5 permit-ами даст тебе 5 одновременных запросов на всё. Но к одному example.com хочется не больше 4 параллельных, а к cdn.fast.io хоть сотню. Один глобальный semaphore не различает домены, и медленный example.com блокирует быстрые.

Решение: per-domain semaphore. Заводишь HashMap<string, Semaphore>, на каждый новый домен лениво создаёшь свой semaphore с нужным лимитом, и каждый запрос обёртывает себя в semaphore своего домена.

Шаг 1 · сервис

import { Context, Effect, HashMap, Layer, Ref, Semaphore } from 'effect';

class DomainLimiter extends Context.Service<DomainLimiter>()('Pulse/DomainLimiter', {
  make: Effect.gen(function* () {
    const semaphores = yield* Ref.make(HashMap.empty<string, Semaphore.Semaphore>());
    const PER_DOMAIN_LIMIT = 4;

    const semaphoreFor = (domain: string) =>
      Effect.gen(function* () {
        const map = yield* Ref.get(semaphores);
        const existing = HashMap.get(map, domain);
        if (existing._tag === 'Some') return existing.value;

        const next = yield* Semaphore.make(PER_DOMAIN_LIMIT);
        yield* Ref.update(semaphores, (m) => HashMap.set(m, domain, next));
        return next;
      });

    const withDomainSlot = <A, E, R>(domain: string, eff: Effect.Effect<A, E, R>) =>
      Effect.gen(function* () {
        const sem = yield* semaphoreFor(domain);
        return yield* sem.withPermits(1)(eff);
      });

    return { withDomainSlot };
  }),
}) {
  static readonly layer = Layer.effect(this, this.make);
}

Что внутри:

  • Ref<HashMap<string, Semaphore>>, мутируемая мапа domain => semaphore.
  • semaphoreFor(domain), читает мапу, если есть, возвращает; если нет, создаёт новый и записывает.
  • withDomainSlot(domain, effect), оборачивает эффект в withPermits(1) соответствующего semaphore.

Шаг 2 · Использование в probe

import { Effect } from 'effect';

import { DomainLimiter } from './services/domain-limiter.ts';
import { HttpService } from './services/http.ts';

const probe = (url: string) =>
  Effect.gen(function* () {
    const http = yield* HttpService;
    const limiter = yield* DomainLimiter;
    const domain = new URL(url).hostname;

    return yield* limiter.withDomainSlot(domain, http.get(url));
  });

probe('https://example.com/a') и probe('https://example.com/b') соревнуются за permit-ы одного example.com-semaphore. probe('https://cdn.fast.io/x') идёт через свой, независимый. На 100 URL одного домена одновременно ровно 4 fetch-а, остальные стоят в очереди на permit.

Шаг 3 · Race condition в инициализации

Внимательный читатель увидит проблему: между Ref.get(semaphores) и Ref.update(semaphores, ...) другой файбер может создать второй semaphore для того же домена. Итог: два разных semaphore-а на один домен, лимит сломан.

Проще всего починить через Ref.modify, атомарную read-then-write операцию:

const semaphoreFor = (domain: string) =>
  Effect.gen(function* () {
    // быстрый путь: уже есть
    const map = yield* Ref.get(semaphores);
    const existing = HashMap.get(map, domain);
    if (existing._tag === 'Some') return existing.value;

    // медленный путь: создаём и атомарно вписываем
    const fresh = yield* Semaphore.make(PER_DOMAIN_LIMIT);
    return yield* Ref.modify(semaphores, (current) => {
      const winner = HashMap.get(current, domain);
      if (winner._tag === 'Some') return [winner.value, current]; // кто-то опередил
      return [fresh, HashMap.set(current, domain, fresh)];
    });
  });

Ref.modify(ref, fn) атомарно: читает значение, вызывает fn, результатом [returnValue, nextState] записывает новое состояние и возвращает returnValue. Если другой файбер успел записать раньше, мы это увидим внутри modify и просто отдадим его semaphore вместо своего. “Лишний” семафор останется без подписчиков и будет собран GC.

В 08 · Транзакции этот же паттерн пишется через TxRef внутри Effect.tx, и читается куда чище. Идея одна: read и write это одна транзакция, не две операции.

Что взять с собой

  • HashMap<string, Semaphore>, per-domain лимит на параллелизм.
  • Semaphore.make(N) создаёт semaphore лениво, по первому запросу домена.
  • Ref.modify атомарно читает и записывает, защищает от двух одновременных создателей.
  • В 08 · Транзакции тот же паттерн станет одной транзакцией без явного “modify”.

Раздел 6 · Latch, ворота для одновременного старта

Сцена · “не стартуем, пока не загрузили”

В Pulse worker-ы пробинга стартуют как файберы, но первое, что им нужно, это загруженный конфиг. Сейчас (см. раздел 1) эту задачу решает Deferred<PulseConfig, ConfigError>: воркер await-ит, получает значение, идёт работать.

Иногда тебе значение не нужно, тебе нужен сигнал. Конфиг уже лежит в Ref-е, его обновляет отдельный файбер, а worker-ы должны просто знать “стартовать пора”. Заводить Deferred<void, never> неудобно: значение всё равно есть, тип говорит “я что-то отдаю”. Чище, когда инструмент сам по себе про “ворота открыты или закрыты”.

Effect для этого даёт Latch: ворота, которые можно закрывать и открывать сколько угодно раз, и await на закрытом latch-е приостанавливает файбер до open.

Latch.make, await, open, close

У Latch, как и у Semaphore, свой модуль в ядре: import { Latch } from 'effect', тип экземпляра Latch.Latch. Раньше и то и другое создавалось методами самого Effect, теперь у каждого примитива координации своя дверь.

import { Effect, Latch } from 'effect';

const program = Effect.gen(function* () {
  const latch = yield* Latch.make(false);
  // false: закрыт; true: открыт сразу

  const worker = (id: number) =>
    Effect.gen(function* () {
      yield* Effect.log(`worker ${id}: waiting`);
      yield* latch.await;
      yield* Effect.log(`worker ${id}: go!`);
    });

  // три воркера ждут на latch
  yield* Effect.forkChild(worker(1));
  yield* Effect.forkChild(worker(2));
  yield* Effect.forkChild(worker(3));

  yield* Effect.sleep('1 second');
  yield* Effect.log('main: opening latch');
  yield* latch.open;

  yield* Effect.sleep('100 millis');
});

Effect.runPromise(program);
// worker 1: waiting
// worker 2: waiting
// worker 3: waiting
// main: opening latch
// worker 1: go!
// worker 2: go!
// worker 3: go!

Все три worker-а просыпаются одновременно, на одном latch.open. Это и есть свойство, которого нет у Deferred: latch можно закрыть обратно через latch.close, и следующие await снова приостановятся.

Сравнение Latch против Deferred

СвойствоDeferred<A, E>Latch
Несёт значениеда, Aнет
Несёт ошибкуда, Eнет
Сколько раз срабатываетодинсколько угодно (open/close/open)
Можно закрыть обратнонетда, latch.close
Сценарийразовый bootstrap, кеш одного значенияпереключаемые “ворота”: maintenance, pause, восстановление

Если нужно “отдать значение и забыть”, это Deferred. Если нужны ворота, которые открываются и закрываются по жизненному циклу подсистемы, это Latch.

latch.whenOpen, “сделай это, когда откроется”

Часто хочется не просто ждать, а выполнить эффект ровно в момент открытия. Для этого есть latch.whenOpen(eff):

const probe = (url: string) =>
  Effect.gen(function* () {
    return yield* latch.whenOpen(realProbe(url));
  });

whenOpen(eff) это сахар над latch.await *> eff. Если latch уже открыт, eff выполняется сразу. Если закрыт, файбер приостанавливается до open, и затем запускает eff. Удобно, когда у тебя десяток мест с одинаковым “ждём готовности”.

latch.release, разовое прохождение

latch.release это “впусти текущих ждущих, оставь ворота закрытыми”. Один раз пропускает накопленных, после следующий await снова блокируется. Полезно для “тиков” в тестах: вместо пуска времени через TestClock (см. 13 · Testing) ты вручную дёргаешь latch и ровно один шаг проходит.

Bootstrap-сценарий в Pulse

Применяем к pulse-<nick>:

import { Context, Effect, Latch, Layer } from 'effect';

class Bootstrap extends Context.Service<Bootstrap>()('Pulse/Bootstrap', {
  make: Effect.gen(function* () {
    const ready = yield* Latch.make(false);

    // эффект, который инициализирует всё нужное и открывает latch
    const init = Effect.gen(function* () {
      yield* Effect.log('bootstrap: loading config');
      yield* loadConfig;
      yield* Effect.log('bootstrap: warming caches');
      yield* warmCaches;
      yield* ready.open;
      yield* Effect.log('bootstrap: ready');
    });

    yield* Effect.forkScoped(init);

    return { ready };
  }),
}) {
  static readonly layer = Layer.effect(this, this.make);
}

Любой код, которому нужна готовность системы, пишет:

const probeAll = Effect.gen(function* () {
  const bootstrap = yield* Bootstrap;
  yield* bootstrap.ready.await;
  // дальше реальный пробинг, конфиг уже загружен, кеши прогреты
  yield* doProbing;
});

Если стартуют 50 worker-ов одновременно, все 50 встанут на ready.await, и все 50 одновременно стартуют после ready.open. Без polling-а: нет лишних wake-up по таймеру, нет race condition между загрузкой и стартом, рантайм не тратит CPU на проверку флага.

И заметь: ready.close (если когда-нибудь захочешь приостановить пробинг для maintenance) аккуратно “запрёт” все worker-ы на следующем ready.await, без необходимости вручную трекать “включён/выключен” флаг и сигнализировать его всем активным файберам.

Что взять с собой

  • Latch, перезаряжаемые ворота. Открыл, закрыл, открыл, сколько угодно.
  • latch.await приостанавливает до open, whenOpen(eff) запускает эффект на открытии.
  • Не несёт значения, в отличие от Deferred. Берёшь, когда нужен сигнал, не данные.
  • В Pulse latch это естественный способ выразить “стартуем после bootstrap” и “ставим на паузу для maintenance”.

Раздел 7 · Pulse · вклад этого урока

Шаг 1 · MonitorEvents через PubSub

Раньше события MonitorEvent писал только Storage (см. 05 · Resources). Теперь у нас будет несколько подписчиков: storage пишет в JSONL, console-printer пишет в stdout, future-Slack-notifier шлёт уведомления. Заводим bus:

// pulse-<nick>/src/services/monitor-events.ts
import { Context, Effect, Layer, PubSub } from 'effect';

import type { MonitorEvent } from '../events.ts';

export class MonitorEvents extends Context.Service<MonitorEvents>()('Pulse/MonitorEvents', {
  make: Effect.gen(function* () {
    const pubsub = yield* Effect.acquireRelease(
      PubSub.bounded<MonitorEvent>(256),
      PubSub.shutdown,
    );

    return {
      publish: (event: MonitorEvent) => PubSub.publish(pubsub, event),
      subscribe: PubSub.subscribe(pubsub),
    };
  }),
}) {
  static readonly layer = Layer.effect(this, this.make);
}

subscribe это Effect<PubSub.Subscription<MonitorEvent>, never, Scope>, scope подписки уйдёт в scope подписчика. publish обычный эффект, его вызывает probe после успеха или ошибки. Сам PubSub мы заворачиваем в Effect.acquireRelease, чтобы на закрытии scope он корректно погасил всех ждущих. Отдельного Layer.scoped для этого не нужно: Layer.effect сам вычитает Scope из требований конструктора.

Шаг 2 · urlsToProbeQueue через Queue

Поток URL для пробинга идёт из конфига (через loadConfig) и из CLI-команды pulse enqueue url. Внутри сервиса один общий buffered queue:

import { Context, Effect, Layer, Queue } from 'effect';

export class UrlQueue extends Context.Service<UrlQueue>()('Pulse/UrlQueue', {
  make: Effect.gen(function* () {
    const queue = yield* Effect.acquireRelease(Queue.bounded<string>(1024), Queue.shutdown);
    return {
      enqueue: (url: string) => Queue.offer(queue, url),
      take: Queue.take(queue),
    };
  }),
}) {
  static readonly layer = Layer.effect(this, this.make);
}

bounded(1024) означает: если кто-то навалит 10 000 URL одной командой, enqueue приостановится после 1024-го, пока worker-ы не разгребут. CLI получит “буфер полон, подождите” в виде ровно правильного типа: enqueue вернёт Effect, который ждёт места.

Worker-ы запускаются в MainLive:

import { Effect, Result } from 'effect';

import { DomainLimiter } from './services/domain-limiter.ts';
import { MonitorEvents } from './services/monitor-events.ts';
import { HttpService } from './services/http.ts';
import { UrlQueue } from './services/url-queue.ts';

const worker = Effect.gen(function* () {
  const queue = yield* UrlQueue;
  const limiter = yield* DomainLimiter;
  const http = yield* HttpService;
  const bus = yield* MonitorEvents;

  while (true) {
    const url = yield* queue.take;
    const domain = new URL(url).hostname;
    const outcome = yield* limiter
      .withDomainSlot(domain, http.get(url))
      .pipe(Effect.result);

    if (Result.isSuccess(outcome)) {
      yield* bus.publish({ _tag: 'ProbeSuccess', url, status: outcome.success.status });
    } else {
      yield* bus.publish({ _tag: 'ProbeFailure', url, reason: String(outcome.failure) });
    }
  }
}).pipe(Effect.forever);

// четыре worker-а на scope MainLive
export const startWorkers = Effect.forEach(
  [1, 2, 3, 4],
  () => Effect.forkScoped(worker),
  { concurrency: 'unbounded' },
);

Что мы получили:

  • провал пробинга не роняет worker: Effect.result превращает его в значение Result, у которого успех лежит в .success, а ошибка в .failure (это тот самый переименованный Either из 03 · Tagged-ошибки);
  • поток URL развязан от пробинга через Queue, producer и consumer работают в своём ритме;
  • четыре worker-а тащат из одной очереди, work-stealing бесплатно;
  • DomainLimiter следит, чтобы один домен не получил больше 4 параллельных fetch-ей;
  • результат уходит в MonitorEvents, сколько угодно подписчиков независимо обработают.

Шаг 3 · bootstrapLatch для старта

До загрузки конфига worker-ы не должны начинать. Latch сидит в bootstrap-сервисе из раздела 6:

const worker = Effect.gen(function* () {
  const bootstrap = yield* Bootstrap;
  yield* bootstrap.ready.await; // <-- ворота

  const queue = yield* UrlQueue;
  // ...как раньше
});

Все четыре worker-а после форка моментально упираются в ready.await, висят без CPU-нагрузки, и проходят одновременно после ready.open в init-файбере.

А ready.close пригодится в будущих сценариях maintenance: останавливаем приём новых URL (UrlQueue.shutdown или ready.close плюс задержка после take), доделываем то, что в полёте, спокойно перезагружаем подсистему.

Шаг 4 · graceful shutdown

runtime.dispose() (см. 05 · Resources) на SIGTERM закроет scope MainLive. Что произойдёт по шагам:

  1. Worker-файберы получают Fiber.interrupt (через scope finalizer от forkScoped).
  2. Queue.take в каждом из них просыпается с interrupt-ом, цикл forever завершается.
  3. Finalizer-ы UrlQueue и MonitorEvents (через PubSub.shutdown под капотом) разбудят возможных оставшихся take-ов с interrupt.
  4. Подписчики MonitorEvents (storage, console-printer) тоже завершают свои циклы через interrupt на PubSub.take.
  5. Finalizer Storage дописывает буфер и закрывает file handle.
  6. process.exit(0).

Никаких “забытых сообщений”, никаких висящих fetch-ей, никаких protruding handle-ов. Координация через примитивы Effect и закрытие через scope, как одно целое.

Что взять с собой

  • MonitorEvents через PubSub, любой подписчик независимо обрабатывает каждое событие.
  • UrlQueue через Queue.bounded, естественный backpressure от пайплайна.
  • DomainLimiter через HashMap<string, Semaphore>, per-domain лимит.
  • bootstrapLatch через Latch, синхронный старт всех worker-ов.
  • Все примитивы регистрируются на scope MainLive, runtime.dispose() закрывает их в правильном порядке.

Финал · чек-лист

Чек-листготово

ДЗ

Дальше

Следующий урок · 08 · Транзакции над общим состоянием. Координация через примитивы решает много, но как только тебе нужна составная атомарная операция над несколькими ячейками (перевести деньги между двумя счетами, переключить SLA-state на основании счётчика провалов), ручные mutex-ы и Ref.modify рассыпаются. Ответ Effect это семейство TxRef, TxQueue, TxHashMap плюс Effect.tx: оптимистичные транзакции, которые прогоняют обычный Effect.gen, сверяют прочитанное при коммите и перезапускают его при конфликте.

Дальше 09 · Stream собирает Queue плюс backpressure плюс forEach в типизированную ленту с операторами (map, filter, groupBy, mapEffect с concurrency). По сути это тот же producer/consumer из этого урока, только декларативно.

Контекст по structured concurrency и forkScoped живёт в 05 · Resources и 06 · Файберы и concurrency. Базовый Ref и почему он не справляется с “read-then-write” операциями подробнее в 08 · Транзакции.