Раздел 25 · Effect-TS

Production-обвязка: CLI, Terminal, HTTP-сервер

middle-senior~140 мин

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

Production-обвязка: CLI, Terminal, HTTP-сервер

Сцена · бортовой компьютер

Самолёт устроен так: одно ядро (двигатели, гидравлика, авионика) и две приборные панели (кабина пилотов и наземный пульт обслуживания). Обе панели смотрят на то же ядро, обе умеют отдавать команды. Их различает не логика, а форма: пилоты дёргают рычаги, наземные инженеры тыкают в терминал на тележке.

К концу урока 11 ядро Pulse уже работает. У нас есть probe, Schedule, MainLive, MonitorEventsPubSub, SlaState на TxRef. Не хватает только поверхностей, через которые с ядром общается внешний мир. В этом уроке наденем на ядро две: CLI для оператора (тот, кто сидит за терминалом) и HTTP-сервер для машин (другие сервисы, дашборды, браузерные клиенты).

Обе поверхности это Effect-программы, оба provide тот же MainLive, оба запускаются через NodeRuntime.runMain. На реальном проекте pulse watch крутится у оператора, pulse serve поднят как systemd-сервис, оба видят одно состояние, потому что делят один и тот же Layer-граф. Это и есть основная мысль урока: точки входа меняются, ядро остаётся.

Где всё это лежит

Один организационный момент, без которого дальше будет путаница в импортах. В Effect 4 отдельных пакетов под CLI и HTTP больше нет. То, что раньше ставилось как @effect/cli и @effect/platform, влилось в ядро effect, в unstable-модули:

что нужнооткуда импортировать
Command, Flag, Argument, Prompteffect/unstable/cli
HttpRouter, HttpServerResponse, HttpMiddlewareeffect/unstable/http
HttpApi, HttpApiGroup, HttpApiEndpoint, HttpApiBuildereffect/unstable/httpapi
Terminal, FileSystem, Path, Stdioпрямо из effect
NodeRuntime, NodeServices, NodeHttpServer@effect/platform-node

Отдельным пакетом остались только те куски, которые физически привязаны к конкретному рантайму: @effect/platform-node (Node), @effect/platform-bun (Bun) и так далее. Всё остальное едет с одной версией всей экосистемы, и рассинхронизации версий между effect и @effect/cli больше не бывает.

Слово unstable в пути означает ровно то, что написано: у этих модулей API может поменяться между минорными версиями. На практике это самые молодые части библиотеки, и в ядре они лежат именно для того, чтобы дозреть.

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

Восемь разделов:

  1. Раздел 1, effect/unstable/cli: Argument, Flag, валидация через Schema.
  2. Раздел 2, Subcommand и сборка иерархии команд без копи-паста.
  3. Раздел 3, Terminal: raw mode, ANSI, корректный выход через finalizer Scope.
  4. Раздел 4, Prompt: интерактивный setup без сторонних библиотек ввода.
  5. Раздел 5, HttpApi и HttpApiBuilder: декларативное описание API со схемами входа и выхода.
  6. Раздел 6, NodeHttpServer и NodeRuntime.runMain: как Layer-сервер ложится на Node-процесс.
  7. Раздел 7, SSE как Stream.fromPubSub плюс HttpServerResponse.stream. Эндпоинт /events без ручной работы с заголовками.
  8. Раздел 8, один MainLive, два entry-point-а. src/cli/commands.ts и src/http/server.ts делят Layer и Runtime; в проде это два процесса, в тестах одна программа.

К концу урока ты описываешь команду через Argument и Flag за пять минут, поднимаешь raw-mode UI с правильным выходом по SIGINT, объясняешь, почему SSE в Effect это не отдельный модуль, а просто Stream плюс HttpServerResponse.stream, и собираешь CLI и HTTP-сервер на одном Layer без дублей.

Раздел 1 · Argument и Flag, декларативное описание команды

Почему process.argv это плохо

Стандартный способ разобрать аргументы в Node это process.argv плюс commander или yargs. Проблем у такого подхода две. Первая: парсинг и валидация разнесены. Сначала ты получаешь строку, потом сам пишешь parseInt, URL.canParse, if (port < 1 || port > 65535). Каждая команда заводит свои проверки, ошибки получаются разные по форме, тесты на парсинг копируются.

Вторая, более тонкая: типы аргументов теряются. После commander ты сидишь с объектом { url: string, interval?: string, expectStatus?: string }, где string-овость это всё, что компилятор знает. Тащить дальше в типобезопасный код приходится с проверками руками.

CLI-модуль Effect решает это иначе. Argument и Flag описывают аргументы декларативно, валидация прицепляется через Schema из урока 02, и в handler команды приходит уже типизированный объект.

Имена тут стоит запомнить сразу, потому что они говорящие и это смена относительно старых материалов: позиционный аргумент это Argument, а --флаг это Flag. Раньше они назывались Args и Options, и по названию Options было совершенно непонятно, что речь про флаги командной строки, а не про объект настроек.

Простейшая команда

// pulse-<nick>/src/cli/commands.ts
import { Effect } from 'effect';
import { Argument, Command } from 'effect/unstable/cli';

const urlArgument = Argument.string('url');

export const probeCommand = Command.make('probe', { url: urlArgument }, ({ url }) =>
  Effect.gen(function* () {
    // url: string, типизированный
    yield* Effect.log(`пинг ${url}`);
  }),
);

Command.make(name, args, handler). Первый параметр это имя, второй это объект с Argument и Flag, третий это функция, которая получает уже разобранные значения. Handler возвращает Effect, и при запуске CLI этот Effect провязывается с тем рантаймом, который ты подсунул.

Заметь форму конструктора: Argument.string('url'), а не объект с полем name. Имя это первый позиционный параметр, и такая же форма у Flag.string('interval').

Уже на этом уровне видна разница с commander: тип url это string, не string | undefined. Если пользователь забудет передать аргумент, CLI сам напечатает help и выйдет с кодом 1, до handler-а дело не дойдёт.

Флаги с дефолтами и алиасами

import { Effect } from 'effect';
import { Argument, Command, Flag } from 'effect/unstable/cli';

const url = Argument.string('url');
const interval = Flag.string('interval').pipe(
  Flag.withAlias('i'),
  Flag.withDefault('30s'),
);
const expectStatus = Flag.integer('expect-status').pipe(Flag.withDefault(200));

export const addCommand = Command.make(
  'add',
  { url, interval, expectStatus },
  ({ url, interval, expectStatus }) =>
    Effect.gen(function* () {
      yield* Effect.log(`добавляю ${url} с интервалом ${interval}, ждём ${expectStatus}`);
    }),
);

Flag.string и Flag.integer уже отличаются типом результата: string отдаёт string, integer отдаёт number и автоматически прерывает запуск с понятной ошибкой, если в --expect-status abc пришло не число. Flag.withDefault(value) снимает с типа Option<A> обёртку: вместо Option<number> в handler приходит number.

Flag.withAlias('i') добавляет короткий ключ. pulse monitor add https://example.com -i 30s и pulse monitor add https://example.com --interval 30s это одно и то же.

Отдельно про файловые флаги, потому что тут поменялась не только приставка. У Flag.file(name, options) есть один булев параметр mustExist, и по умолчанию он выключен:

// файл может ещё не существовать, первый запуск его создаст
const editConfig = Flag.file('config').pipe(Flag.withDefault('./pulse.config.json'));

// файл обязан быть на месте, иначе CLI откажется запускаться
const readConfig = Flag.file('config', { mustExist: true }).pipe(
  Flag.withDefault('./pulse.config.json'),
);

Раньше на этом месте было трёхзначное { exists: 'yes' | 'no' | 'either' } с дефолтом “файл обязан существовать”. Если переносишь старый код, обрати внимание: дефолт поменялся на противоположный, и молча пропущенный mustExist: true превратит “проверяю наличие” в “не проверяю”.

Валидация через Schema

Сюда же подключается Schema из урока 02. Если хочется, чтобы url был не просто строкой, а валидным URL, или чтобы --interval парсился как длительность, прицепляешь Schema:

import { Effect, Schema, SchemaGetter, SchemaIssue } from 'effect';
import { Argument, Flag } from 'effect/unstable/cli';

const Url = Schema.String.check(
  Schema.isPattern(/^https?:\/\//, {
    description: 'URL должен начинаться с http:// или https://',
  }),
);

const Interval = Schema.String.pipe(
  Schema.decodeTo(Schema.Number, {
    decode: SchemaGetter.transformOrFail((input: string) => {
      const match = /^(\d+)(s|m|h)$/.exec(input);
      if (!match) {
        return Effect.fail(new SchemaIssue.InvalidValue({ message: 'ожидается 30s, 5m или 1h' }));
      }
      const [, num, unit] = match;
      const ms = Number(num) * (unit === 's' ? 1000 : unit === 'm' ? 60_000 : 3_600_000);
      return Effect.succeed(ms);
    }),
    encode: SchemaGetter.transform((ms: number) => `${ms / 1000}s`),
  }),
);

const url = Argument.string('url').pipe(Argument.withSchema(Url));
const interval = Flag.string('interval').pipe(
  Flag.withSchema(Interval),
  Flag.withDefault(30_000),
);

Теперь в handler приходит url: string (уже проверенный) и interval: number (миллисекунды). Невалидные значения CLI отлавливает до handler-а, печатает понятную ошибку, выходит с кодом 1. Это та же дисциплина “валидация на входе” из урока 02, только применённая к командной строке.

Форма записи схем тут ровно та, что разбиралась в уроке 02, и по ней стоит сверяться, если переносишь старый код. Проверки перестали быть pipe-комбинаторами и уехали внутрь .check(...), а их имена получили приставку is: Schema.pattern стал Schema.isPattern, Schema.int стал Schema.isInt. Преобразование одного типа в другой теперь пишется как Schema.decodeTo(Куда, { decode, encode }), где обе функции берутся из модуля SchemaGetter. Модуля ParseResult больше нет: неудачу возвращаешь обычным Effect.fail с ошибкой из SchemaIssue.

Variadic и optional

const urls = Argument.string('url').pipe(Argument.variadic); // string[]
const oneToFive = Argument.string('url').pipe(Argument.variadic({ min: 1, max: 5 }));
const name = Flag.string('name').pipe(Flag.optional); // Option<string>

Argument.variadic собирает все позиционные аргументы в массив. pulse remove a.com b.com c.com даст urls: ['a.com', 'b.com', 'c.com']. Границы задаются опциями { min, max }. Flag.optional оставляет Option<A> без withDefault, и в handler приходит Option<string>, который ты разбираешь через Option.match или Option.getOrUndefined.

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

  • Argument.string, Argument.integer это позиционные аргументы; Flag.string, Flag.integer это именованные флаги.
  • Имя передаётся первым параметром: Argument.string('url'), Flag.string('interval').
  • withDefault, withAlias, optional, variadic это модификаторы для типичных кейсов; Flag.file управляет проверкой наличия через булев mustExist (по умолчанию выключен).
  • Валидация прицепляется через Argument.withSchema / Flag.withSchema, ошибки парсинга превращаются в человекочитаемый help.
  • В handler команды приходит уже типизированный объект; никакого string | undefined тащить в логику не надо.

Раздел 2 · Subcommand, иерархия команд

Зачем дерево команд

pulse это не одна команда, а семейство. Оператор пишет:

pulse monitor add https://github.com --interval 30s
pulse monitor list
pulse monitor remove https://github.com
pulse watch
pulse serve --port 8080
pulse init

У monitor есть три подкоманды (add, list, remove). У корневого pulse есть пять (monitor, watch, serve, init, replay). Если описывать это руками через if/else по argv[2], получится та же каша, что и в любом git-обёртке без библиотеки.

Subcommand это механизм собрать такое дерево декларативно.

Сборка дерева

// pulse-<nick>/src/cli/commands.ts
import { Command } from 'effect/unstable/cli';

const monitorGroup = Command.make('monitor').pipe(
  Command.withDescription('Управление списком отслеживаемых URL'),
  Command.withSubcommands([monitorAdd, monitorList, monitorRemove]),
);

export const pulse = Command.make('pulse').pipe(
  Command.withDescription('Сторож URL для самодельщиков'),
  Command.withSubcommands([initCommand, monitorGroup, watch, serveCommand]),
);
// pulse-<nick>/src/cli/index.ts
import { NodeRuntime, NodeServices } from '@effect/platform-node';
import { Effect } from 'effect';
import { Command } from 'effect/unstable/cli';

import { MainLive } from '../layers.ts';
import { pulse } from './commands.ts';

const cli = Command.run({ version: '0.1.0' })(pulse);

cli.pipe(
  Effect.provide(MainLive),
  Effect.provide(NodeServices.layer),
  NodeRuntime.runMain,
);

Что здесь происходит:

  • Command.make('pulse') без handler-а это узел без действия. Без подкоманды CLI напечатает help и выйдет с кодом 1.
  • Command.withSubcommands(children) присоединяет дочерние команды.
  • Command.run({ version })(root) возвращает готовый Effect, а не функцию от argv. Аргументы командной строки он берёт сам, из сервиса Stdio. Отсюда и порядок аргументов: сначала настройки запуска, потом команда.
  • Effect.provide(MainLive) подсовывает все сервисы Pulse (Storage, HttpClient, MonitorEventsPubSub, Sla из урока 08). NodeServices.layer добавляет платформенные сервисы (Terminal, FileSystem, Path, Stdio), без них CLI не соберётся по типам.
  • NodeRuntime.runMain запускает программу с правильной обработкой сигналов: SIGINT и SIGTERM прерывают корневой файбер, тот вызывает все finalizer-ы scope, процесс выходит с правильным кодом.

Две вещи, на которых спотыкаются при переносе старого кода. Раньше запуск писался как Command.run(root, { name, version })(process.argv): команда шла первым аргументом, а argv передавался руками. Теперь process.argv не передаётся вовсе, и это не косметика: чтение аргументов стало обычной зависимостью от сервиса, а значит в тестах его можно подменить, не трогая глобальный process. И NodeContext.layer переименовался в NodeServices.layer.

Бесплатное приложение

Раз уж запуск проходит через Command.run, ты бесплатно получаешь набор глобальных флагов, которые не надо объявлять:

  • --help на любом уровне дерева;
  • --wizard, интерактивный режим, который проводит по всем флагам команды вопрос за вопросом;
  • --completions <shell>, генерация скрипта автодополнения для bash, zsh или fish;
  • --log-level, уровень логирования на один запуск.

--wizard особенно приятен на командах с десятком флагов: вместо чтения справки пользователь просто отвечает на вопросы, а CLI в конце печатает готовую строку запуска.

Общие флаги через родителя

Бывает, что у всех подкоманд есть общий флаг, например --config. Дублировать его во всех Command.make накладно, и для этого есть Command.withSharedFlags:

import { Command, Flag } from 'effect/unstable/cli';

const config = Flag.file('config').pipe(
  Flag.withDefault('./pulse.json'),
  Flag.withDescription('Путь к конфиг-файлу'),
);

export const pulse = Command.make('pulse').pipe(
  Command.withSharedFlags({ config }),
  Command.withSubcommands([monitorGroup, watch]),
);

Флаг объявлен один раз, а разобранное значение приходит в handler каждой дочерней команды. В Pulse мы этим не пользуемся (конфиг у нас объявлен на самих подкомандах, потому что у watch и у add разные требования к наличию файла), но знать про возможность полезно.

Коды возврата через канал ошибок

Главная фишка Effect.provide(MainLive): сервисы внутри handler-ов это типизированный канал. Если handler упал в типизированную ошибку, эта ошибка автоматически попадает в NodeRuntime.runMain, который вернёт ненулевой код. Никакого process.exit(1) в коде писать не надо:

import { Effect } from 'effect';
import { Argument, Command } from 'effect/unstable/cli';

import { Storage, StorageError } from '../services/storage.ts';

export const removeCommand = Command.make(
  'remove',
  { url: Argument.string('url') },
  ({ url }) =>
    Effect.gen(function* () {
      const storage = yield* Storage;
      const removed = yield* storage.remove(url);
      if (!removed) {
        return yield* Effect.fail(new StorageError({ reason: 'not-found', url }));
      }
      yield* Effect.log(`удалил ${url}`);
    }),
);

При Effect.fail(new StorageError(...)) NodeRuntime.runMain напечатает stack trace типизированной ошибки и выйдет с кодом 1. При успехе выйдет с 0. Это совместимо со скриптами и CI: pulse monitor remove unknown.com && echo "удалил" отработает корректно.

Ошибки самого разбора аргументов (не хватает позиционного, неизвестный флаг, невалидная схема) приходят отдельным типом CliError. Раньше он назывался ValidationError; если ловишь его по имени, поправь.

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

  • Command.withSubcommands([...]) собирает дерево, help печатается автоматически.
  • Command.run({ version })(root) возвращает готовый Effect, аргументы он читает из сервиса Stdio, process.argv руками не передаётся.
  • Глобальные --help, --wizard, --completions, --log-level приезжают бесплатно.
  • NodeRuntime.runMain сам ставит обработчики SIGINT/SIGTERM и возвращает правильный exit-код по результату handler-а.
  • Типизированные ошибки из канала автоматически превращаются в exit-код 1, отдельно process.exit не нужен.

Раздел 3 · Terminal, raw mode, и корректный выход через Scope

Что значит “живая таблица”

Команда pulse watch это та самая ANSI-таблица, как у k9s, btop, top: она занимает весь терминал, перерисовывается каждые N миллисекунд, реагирует на клавиши (q для выхода, r для ручного refresh), и при Ctrl+C чисто восстанавливает экран. Без библиотеки типа ink или blessed.

Что для этого нужно от терминала:

  • альтернативный экран (escape-последовательность \x1b[?1049h входит, \x1b[?1049l выходит), чтобы после выхода вернуть пользователю его prompt без мусора.
  • raw mode, чтобы реагировать на каждое нажатие клавиши, а не ждать Enter.
  • курсор и очистку экрана, чтобы перерисовать кадр.
  • Финализатор, который при выходе (Ctrl+C, исключение, нормальное завершение) гарантированно вернёт raw mode в обычный, выйдет из alternate screen и покажет курсор.

Главная ловушка: ставить обработчик process.on('SIGINT', cleanup). Он работает в простых случаях, но плохо комбинируется с библиотечным кодом, дублируется при перезапуске, висит после остановки и портит тесты. У Effect для этого есть штатный механизм: Scope и Effect.addFinalizer.

Скелет команды watch

// pulse-<nick>/src/cli/commands/watch.ts
import { Effect, Schedule, Terminal } from 'effect';
import { Command } from 'effect/unstable/cli';

import { Storage } from '../../services/storage.ts';
import { Sla } from '../../concurrency/sla-state.ts';

const enterAlternateScreen = (terminal: Terminal.Terminal) =>
  Effect.gen(function* () {
    yield* terminal.display('\x1b[?1049h\x1b[?25l'); // alternate screen + скрыть курсор
    yield* Effect.addFinalizer(() =>
      terminal.display('\x1b[?25h\x1b[?1049l').pipe(Effect.ignore),
    );
  });

const renderFrame = (terminal: Terminal.Terminal) =>
  Effect.gen(function* () {
    const storage = yield* Storage;
    const sla = yield* Sla;

    const monitors = yield* storage.list;
    const slaState = yield* sla.snapshot;

    yield* terminal.display('\x1b[H\x1b[2J'); // курсор в (1,1), очистка экрана
    yield* terminal.display(`Pulse · ${monitors.length} мониторов · активный ${slaState.active}\n\n`);
    for (const m of monitors) {
      const colour = m.lastStatus && m.lastStatus < 400 ? '\x1b[32m' : '\x1b[31m';
      yield* terminal.display(
        `${colour}● \x1b[0m${m.url.padEnd(40)} ${String(m.lastStatus ?? '---').padStart(3)}  ${m.lastLatencyMs ?? '---'} ms\n`,
      );
    }
  });

export const watchCommand = Command.make('watch', {}, () =>
  Effect.gen(function* () {
    const terminal = yield* Terminal.Terminal;
    yield* enterAlternateScreen(terminal);
    yield* renderFrame(terminal).pipe(
      Effect.repeat(Schedule.spaced('500 millis')),
    );
  }).pipe(Effect.scoped),
);

Что тут важно:

  • Terminal.Terminal это сервис прямо из ядра effect, реализацию под Node даёт NodeServices.layer. Никакого process.stdout.write в нашем коде.
  • enterAlternateScreen это функция, которая делает две вещи: переходит в alternate screen и регистрирует finalizer на возврат. Finalizer выполнится в момент закрытия scope, неважно как scope закрывается (нормально, по ошибке, по прерыванию).
  • Effect.scoped оборачивает программу команды в свой scope. Когда NodeRuntime прервёт корневой файбер по SIGINT, scope watchCommand тоже закроется, finalizer-ы развернут терминал в исходное.
  • renderFrame.pipe(Effect.repeat(Schedule.spaced('500 millis'))) это та самая логика “перерисовывать каждые полсекунды” из урока 11. Schedule делает паузу между запусками, файбер крутится в NodeRuntime.

Где SIGINT попадает в программу

В уроке 06 разбирался механизм Fiber.interrupt. У NodeRuntime.runMain есть встроенный обработчик: на SIGINT/SIGTERM рантайм вызывает interrupt на корневом файбере. Корневой файбер это тот, что запустил твою программу. Прерывание распространяется по дереву файберов вниз и закрывает scope-ы.

В нашем случае:

NodeRuntime root fiber
└─ watchCommand handler fiber
   └─ scope of Effect.scoped
      ├─ finalizer: вернуть raw mode, выйти из alternate screen
      └─ renderFrame loop fiber
         └─ Schedule.spaced timer

Ctrl+C SIGINT попадает в root файбер, interrupt каскадно проходит вниз, на закрытии scope finalizer возвращает терминал. На экране у пользователя моргнёт его prompt, никакой “битой” последовательности типа \x1b[?25l без парного \x1b[?25h не останется.

Что runMain делает под капотом

Тут полезно остановиться и посмотреть, как именно runMain ловит сигнал. Это не магия рантайма, а примерно полсотни строк поверх двух простых вещей: Effect.runFork (запустить эффект в файбере и получить на него ссылку) и обычный process.on('SIGINT') из Node. Вот его суть, очищенная от мелочей:

// упрощённый @effect/platform-node-shared
const runMain = (effect) => {
  // 1. запускаем программу в корневом fiber-е и держим ссылку на него
  const fiber = Effect.runFork(effect);

  // 2. фиктивный таймер не даёт Node закрыть процесс раньше времени
  const keepAlive = setInterval(() => {}, 2 ** 31 - 1);

  let receivedSignal = false;

  // 4. наблюдатель срабатывает, когда корневой fiber завершился любым Exit
  fiber.addObserver((exit) => {
    if (!receivedSignal) {
      process.removeListener('SIGINT', onSignal);
      process.removeListener('SIGTERM', onSignal);
    }
    clearInterval(keepAlive);
    teardown(exit, (code) => {
      if (receivedSignal || code !== 0) process.exit(code);
    });
  });

  // 3. сигнал прерывает корневой fiber и снимает обработчики
  function onSignal() {
    receivedSignal = true;
    process.removeListener('SIGINT', onSignal);
    process.removeListener('SIGTERM', onSignal);
    fiber.interruptUnsafe();
  }

  process.on('SIGINT', onSignal);
  process.on('SIGTERM', onSignal);
};

Каждая строчка тут закрывает конкретную задачу, которую ты бы иначе решал руками.

1 · Effect.runFork отдаёт ссылку на корневой файбер. Это тот же fork из урока 06, только на самом верху: программу нужно не просто запустить, а сохранить дескриптор, чтобы потом было что прерывать. Всё дерево твоих мониторов, scope-ов и таймеров растёт из этого одного файбера.

2 · keepAlive держит event loop живым. Node закрывает процесс, как только в event loop не осталось запланированной работы. Обычно ожидание Effect.sleep, таймера расписания или сетевого ответа само по себе считается такой работой, но setInterval(() => {}, 2 ** 31 - 1) это страховка: даже если в моменте все эффекты висят на внутренних промисах рантайма, фиктивный таймер не даёт Node решить, что делать нечего, и выйти преждевременно. На завершении его обязательно гасит clearInterval, иначе он сам держал бы процесс вечно.

3 · process.on('SIGINT') это единственное место, где Effect трогает сырой Node API. Сигнал приходит в обычный callback, вне всякого файбера, поэтому прервать корневой файбер отсюда можно только небезопасным методом fiber.interruptUnsafe(). Разберём название:

  • Unsafe, потому что вызов идёт из императивного Node-callback, а не из Effect. Это та же граница, что и runCallback из урока 06. Приставка стоит суффиксом, и это в v4 сквозная конвенция: offerUnsafe, openUnsafe, completeUnsafe, interruptUnsafe.
  • Метод не ждёт завершения: он ставит файберу флаг прерывания и сразу возвращается. Обработчик сигнала обязан отработать мгновенно, а само сворачивание дерева пойдёт уже на корневом файбере, на его ближайшем suspension point. Дождаться Exit из Effect-кода это отдельная функция, Fiber.await(fiber).
  • Необязательный аргумент это идентификатор того, кто прерывает. Опустили, значит прерывание атрибутируется самому файберу: в Cause это будет выглядеть как “файбер прервал сам себя”, а не “пришёл внешний interrupt от чужого файбера”.

Сразу за этим обработчик снимает оба слушателя. Это сознательный предохранитель: первый Ctrl+C запускает мягкую остановку, а так как слушателя больше нет, второй Ctrl+C уже падает в дефолтное поведение Node, то есть в жёсткий немедленный kill. Если какой-нибудь finalizer завис (например, медленный server.close() ждёт долгие соединения), у оператора всегда есть способ добить процесс вторым нажатием.

4 · addObserver это финальный аккорд. Когда корневой файбер получил флаг прерывания, дальше работает механика из урока 06: на ближайшем suspension point interrupt каскадом идёт вниз по дереву, дочерние файберы прерываются, scope-ы закрываются, finalizer-ы отрабатывают в порядке LIFO. Всё это занимает время (ровно столько, сколько живут твои finalizer-ы), и только когда дерево полностью свернулось, корневой файбер выдаёт свой Exit. Вот тут и срабатывает наблюдатель: гасит keepAlive и считает exit-код через teardown. Заметь красивую симметрию с уроком 06: там Fiber.interrupt(fiber) дожидался Exit синхронно внутри Effect; здесь обработчик сигнала ждать не может, поэтому ожидание вынесено в addObserver.

teardown и exit-код. Дефолтный teardown короткий:

const defaultTeardown = (exit, onExit) => {
  onExit(Exit.isFailure(exit) && !Cause.hasInterruptsOnly(exit.cause) ? 1 : 0);
};

Логика простая: если программа упала настоящей ошибкой (failure, в котором есть не только прерывание), это код 1. Если завершилась успехом или была прервана начисто (а именно так выглядит остановка по SIGTERM), это код 0. То есть graceful shutdown по сигналу выходит с нулём, а не с ошибкой, и это правильно: тебя попросили остановиться, ты корректно остановился, сбоя не было. Здесь и нужны предикаты Cause из урока 03: они отличают “чистое прерывание” от “ошибка плюс прерывание”. В v4 Cause плоский, внутри лежит массив причин, и предикаты поэтому во множественном числе: hasFails, hasDies, hasInterrupts.

Последняя деталь, строчка if (receivedSignal || code !== 0) process.exit(code). Принудительный выход вызывается в двух случаях: пришёл сигнал или код ненулевой. А при штатном добровольном завершении (код 0, сигнала не было) runMain process.exit не зовёт вовсе: keepAlive уже погашен, и процесс выходит сам, дав event loop спокойно дослить остатки (последний flush лога, буфер stdout). На сигнале же выход форсируется даже с кодом 0, чтобы процесс не повис, если что-то внешнее ещё удерживает event loop.

Отсюда и вывод, который ты уже видел в разделе про SIGINT: руками process.on('SIGINT', cleanup) писать не нужно. Всё, что ты бы там делал (поймать сигнал, прервать работу, дождаться чистки, выставить exit-код), runMain уже делает корректно и в правильном порядке. Твоя задача только вешать finalizer-ы на scope.

Чтение клавиш

Если хочется реагировать на q и r руками, поднимаешь дополнительный файбер, который читает stdin:

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

const inputs = Effect.gen(function* () {
  const terminal = yield* Terminal.Terminal;
  const queue = yield* terminal.readInput;
  return Stream.fromQueue(queue).pipe(Stream.map((event) => event.key.name));
});

Обрати внимание на форму terminal.readInput: это не “прочитать одно нажатие”, а “открыть очередь нажатий”. Эффект скоупный, и подписка на stdin снимется вместе со скоупом. Дальше очередь превращается в стрим обычным Stream.fromQueue (09 · Stream).

Внутри handler ты подписываешься на этот стрим и в Stream.runForEach решаешь, что делать с каждой клавишей:

const inputLoop = (signal: Deferred.Deferred<void>) =>
  Effect.gen(function* () {
    const stream = yield* inputs;
    yield* stream.pipe(
      Stream.runForEach((key) =>
        key === 'q' ? Deferred.succeed(signal, undefined) : Effect.void,
      ),
    );
  });

q сигналит через Deferred (см. урок 07), и главный файбер делает Deferred.await параллельно с renderFrame.repeat. Тот, кто первый отработает (либо пользователь нажал q, либо рантайм получил SIGINT), завершит программу через Effect.race.

В реальной команде watch это выглядит примерно так:

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

export const watchCommand = Command.make('watch', {}, () =>
  Effect.gen(function* () {
    const terminal = yield* Terminal.Terminal;
    const stop = yield* Deferred.make<void>();

    yield* enterAlternateScreen(terminal);

    yield* Effect.race(
      renderFrame(terminal).pipe(Effect.repeat(Schedule.spaced('500 millis'))),
      Effect.race(
        inputLoop(stop),
        Deferred.await(stop),
      ),
    );
  }).pipe(Effect.scoped),
);

Effect.race это структурный механизм: выигравший продолжает, проигравшие прерываются. Когда Deferred.await(stop) отрабатывает (пользователь нажал q), две другие гонки прерываются, их файберы каскадно interrupt-ятся, scope закрывается, terminal возвращается в исходное.

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

  • Terminal.Terminal это платформенный сервис, отдаёт display и readInput через типизированные эффекты.
  • Alternate screen, raw mode, скрытие курсора это ANSI-последовательности, отправляются через terminal.display.
  • Корректный выход это finalizer в Scope, не process.on('SIGINT'). NodeRuntime сам ловит сигналы и каскадно interrupt-ит файберы.
  • Effect.scoped создаёт локальный scope для команды, закрытие scope гарантированно вернёт терминал.
  • Чтение клавиш через Stream, остановка через Deferred, всё это race-ится с основным рендером.

Раздел 4 · Prompt, интерактивный setup

Зачем интерактив

Команда pulse init это тот момент, когда пользователь первый раз запустил программу: ему надо задать пару базовых вопросов и сохранить ответы в конфиг. Можно было бы заставить его руками редактировать pulse.json, но это плохой UX. Лучше провести через три вопроса в терминале:

? Имя проекта: my-pulse
? Интервал по умолчанию: (используй стрелки)
  ❯ 30 секунд
    1 минута
    5 минут
? Записывать события в JSONL? (y/N): y

В effect/unstable/cli для этого есть модуль Prompt с тремя основными типами: Prompt.text (свободный текст), Prompt.select (одно из множества), Prompt.toggle (yes/no). Рядом с ними лежат confirm, password, multiSelect, autoComplete, list, date и file.

Три вопроса

// pulse-<nick>/src/cli/commands/init.ts
import { Effect } from 'effect';
import { Command, Prompt } from 'effect/unstable/cli';

import { Config } from '../../services/config.ts';

const askName = Prompt.text({
  message: 'Имя проекта',
  default: 'pulse',
  validate: (input) =>
    input.length >= 2
      ? Effect.succeed(input)
      : Effect.fail('Имя должно быть не короче двух символов'),
});

const askInterval = Prompt.select({
  message: 'Интервал по умолчанию',
  choices: [
    { title: '30 секунд', value: '30s' },
    { title: '1 минута', value: '1m' },
    { title: '5 минут', value: '5m' },
  ],
});

const askJsonl = Prompt.toggle({
  message: 'Записывать события в JSONL?',
  initial: true,
  active: 'да',
  inactive: 'нет',
});

export const initCommand = Command.make('init', {}, () =>
  Effect.gen(function* () {
    const name = yield* askName;
    const interval = yield* askInterval;
    const jsonl = yield* askJsonl;

    const config = yield* Config;
    yield* config.write({ name, defaultInterval: interval, jsonl });

    yield* Effect.log(`сохранил конфиг для проекта ${name}`);
  }),
);

Что здесь важно:

  • Prompt.text({ validate }) принимает функцию string -> Effect<string, string, never>, которая решает, ок ли ввод. Если вернуть Effect.fail, пользователю показывают сообщение и спрашивают заново. Это та же дисциплина “валидация на границе”, что и у Schema в Раздел 1, только асинхронная и интерактивная.
  • Prompt.select отдаёт ровно тип value из выбранного choice: в нашем случае '30s' | '1m' | '5m', не string. Опечататься по значению нельзя, оно вычисляется из массива.
  • Prompt.toggle отдаёт boolean. active/inactive это подписи рядом с курсором.
  • Effect.gen собирает три prompt-а в линейный сценарий: сначала имя, потом интервал, потом флажок. Если на любом шаге пользователь нажал Ctrl+C, Prompt.X отдаёт interrupt, finalizer scope чисто закрывает терминал.

Композиция и условные ветви

Prompt это Effect, поэтому композиция работает как обычно. “Если включил jsonl, спроси, в какой файл писать”:

const init = Effect.gen(function* () {
  const name = yield* askName;
  const interval = yield* askInterval;
  const jsonl = yield* askJsonl;
  const jsonlPath = jsonl
    ? yield* Prompt.text({ message: 'Путь к JSONL-файлу', default: './pulse.jsonl' })
    : null;
  return { name, interval, jsonl, jsonlPath };
});

Никаких машин состояний, никаких “промежуточных стейтов между вопросами”. Это просто последовательность yield* в Effect.gen, и компилятор знает все типы.

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

  • Prompt.text, Prompt.select, Prompt.toggle из effect/unstable/cli это базовая тройка, на ней закрывается большинство интерактивных сценариев.
  • validate это Effect<string, string, never>, ошибка показывается пользователю, успех возвращает прошедшее значение.
  • Prompt-ы это Effect, поэтому композируются в Effect.gen как любые другие шаги.
  • Ctrl+C на любом prompt-е чисто прерывает программу, никаких висящих обработчиков.

Раздел 5 · HttpApi и HttpApiBuilder, декларативный сервер

Зачем декларативный HTTP

В классическом express-стиле endpoint это функция (req, res) -> void, которая руками вытаскивает body, валидирует, считает ответ, сериализует, пишет в res. Типы по дороге теряются на каждом шаге, валидация дублируется, OpenAPI-спека пишется отдельно и расходится с кодом через две недели.

HttpApi из effect/unstable/httpapi это другой подход. Сначала ты описываешь API как структуру данных: пути, методы, схемы запросов и ответов. Потом отдельно реализуешь handler-ы: тип запроса и ответа компилятор знает из схемы. Из той же структуры автоматически генерируются OpenAPI-спека, клиент с типизированными вызовами, тесты.

Описание API

// pulse-<nick>/src/http/api.ts
import { Schema } from 'effect';
import { HttpApi, HttpApiEndpoint, HttpApiGroup, HttpApiSchema } from 'effect/unstable/httpapi';

const Monitor = Schema.Struct({
  id: Schema.String,
  url: Schema.String,
  intervalMs: Schema.Number,
  lastStatus: Schema.NullOr(Schema.Number),
  lastLatencyMs: Schema.NullOr(Schema.Number),
});

const NotFound = Schema.TaggedStruct('NotFound', { id: Schema.String }).pipe(
  HttpApiSchema.status(404),
);

const monitorsGroup = HttpApiGroup.make('monitors').add(
  HttpApiEndpoint.get('list', '/api/monitors', {
    success: Schema.Array(Monitor),
  }),
  HttpApiEndpoint.get('byId', '/api/monitors/:id', {
    params: { id: Schema.String },
    success: Monitor,
    error: NotFound,
  }),
  HttpApiEndpoint.post('create', '/api/monitors', {
    payload: { url: Schema.String, intervalMs: Schema.Number },
    success: Monitor,
  }),
);

const Status = Schema.Struct({
  ok: Schema.Boolean,
  sla: Schema.Struct({
    active: Schema.Literals(['primary', 'fallback']),
    consecutiveFailures: Schema.Number,
  }),
  monitors: Schema.Number,
});

const statusGroup = HttpApiGroup.make('status').add(
  HttpApiEndpoint.get('snapshot', '/status', { success: Status }),
);

export const PulseApi = HttpApi.make('pulse').add(monitorsGroup, statusGroup);

Что здесь происходит:

  • HttpApiGroup.make(name) это контейнер для связанных endpoint-ов. .add(...) принимает сразу несколько, отдельный вызов на каждый не нужен. Пока две группы (monitors и status); в Разделе 7 добавим третью для SSE.
  • HttpApiEndpoint.get(id, path, options) описывает один endpoint. Всё, что раньше навешивалось цепочкой сеттеров, теперь лежит в объекте опций: params, query, headers, payload, success, error.
  • Параметры пути пишутся прямо в строке пути как :id, а их схемы идут в params. Раньше путь собирался через шаблонный литерал с подстановкой HttpApiSchema.param(...); эта конструкция была изобретательной, но читалась тяжело, и её убрали.
  • HTTP-код ошибки прицепляется к самой схеме через HttpApiSchema.status(404), а не к вызову .addError. Логика простая: код это свойство ответа, значит место ему на схеме ответа.

Уже на этом уровне у тебя есть источник истины для API: схемы запросов, схемы ответов, ошибки. Сменишь поле в Monitor это сменится везде: и в handler-е, и у клиента, и в OpenAPI-спеке.

Реализация handler-ов

// pulse-<nick>/src/http/handlers.ts
import { Effect, Layer } from 'effect';
import { HttpApiBuilder } from 'effect/unstable/httpapi';

import { Sla } from '../concurrency/sla-state.ts';
import { Storage } from '../services/storage.ts';
import { PulseApi } from './api.ts';

const MonitorsLive = HttpApiBuilder.group(PulseApi, 'monitors', (handlers) =>
  handlers
    .handle('list', () =>
      Effect.gen(function* () {
        const storage = yield* Storage;
        return yield* storage.list;
      }),
    )
    .handle('byId', ({ params: { id } }) =>
      Effect.gen(function* () {
        const storage = yield* Storage;
        const monitor = yield* storage.get(id);
        if (monitor === null) {
          return yield* Effect.fail({ _tag: 'NotFound' as const, id });
        }
        return monitor;
      }),
    )
    .handle('create', ({ payload }) =>
      Effect.gen(function* () {
        const storage = yield* Storage;
        return yield* storage.add(payload);
      }),
    ),
);

const StatusLive = HttpApiBuilder.group(PulseApi, 'status', (handlers) =>
  handlers.handle('snapshot', () =>
    Effect.gen(function* () {
      const storage = yield* Storage;
      const sla = yield* Sla;
      const monitors = yield* storage.list;
      const slaState = yield* sla.snapshot;
      return {
        ok: slaState.active === 'primary',
        sla: slaState,
        monitors: monitors.length,
      };
    }),
  ),
);

export const PulseApiLive = HttpApiBuilder.layer(PulseApi).pipe(
  Layer.provide(MonitorsLive),
  Layer.provide(StatusLive),
);

Каждый handler это Effect, получающий уже разобранные параметры (params, payload, query). Тип параметров и возврата компилятор берёт из PulseApi: если ты в list вернёшь не массив Monitor-ов, компилятор ругнётся, не на запуске, а на этапе типчека.

Финальная сборка это HttpApiBuilder.layer(PulseApi), куда через Layer.provide подаются реализации групп. Имя layer вместо прежнего api тут не случайно: результат это обычный Layer, который дальше складывается с остальными точно так же, как любой другой.

Сервисы (Storage, Sla) приходят из канала R. То же MainLive, что использует CLI и тесты, обслуживает HTTP-сервер.

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

  • HttpApi это описание API: пути, методы, схемы запросов и ответов. Живёт в effect/unstable/httpapi.
  • Endpoint описывается объектом опций: HttpApiEndpoint.get(id, path, { params, query, payload, success, error }).
  • HTTP-код ошибки навешивается на схему через HttpApiSchema.status(code).
  • HttpApiBuilder.group это реализация: для каждого endpoint-а свой Effect-handler с типизированным запросом и ответом; HttpApiBuilder.layer собирает их вместе.
  • Источник истины один (HttpApi-объект), отсюда выводятся handler-ы, клиент, OpenAPI, тесты.
  • Сервисы handler-ов приходят из R, тот же MainLive обслуживает CLI и HTTP без дубликации.

Раздел 6 · NodeHttpServer и провязка с MainLive

Роутер это сервис, а маршрут это Layer

Прежде чем поднимать сервер, надо усвоить главную идею HTTP-слоя в Effect 4, и она контринтуитивная. Роутер больше не значение, в которое ты складываешь маршруты. Роутер это сервис, а каждый маршрут это Layer, который сам себя в этом сервисе регистрирует.

Сравни две формы. Раньше маршруты набирались в пустой роутер через pipe:

// как это выглядело раньше
export const pulseRouter = HttpRouter.empty.pipe(
  HttpRouter.get('/status', statusHandler),
  HttpRouter.get('/events', eventsHandler),
);
export const pulseHttpApp = HttpServer.serve(pulseRouter);

Теперь маршруты это слои, и складываются они как слои:

// pulse-<nick>/src/http/server.ts
import { Layer } from 'effect';
import { HttpRouter } from 'effect/unstable/http';

export const PulseRoutes = Layer.mergeAll(
  HttpRouter.add('GET', '/status', statusHandler),
  HttpRouter.add('GET', '/api/monitors', monitorsHandler),
  HttpRouter.add('GET', '/events', eventsHandler),
);

export const pulseHttpServer = HttpRouter.serve(PulseRoutes);

Что изменилось построчно. Метод стал явным первым аргументом: вместо HttpRouter.get(path, h) пишется HttpRouter.add('GET', path, h). HttpRouter.empty исчез, потому что складывать больше не во что. Запуск идёт от Layer-а приложения: HttpRouter.serve(PulseRoutes) вместо HttpServer.serve(router).

Что это даёт. Маршрут теперь честно объявляет свои зависимости в типе Layer-а, и они всплывают в требованиях serve. То есть если eventsHandler тянет MonitorEventsPubSub, это видно в типе всего сервера, и подать сервис можно обычным Layer.provide в корне композиции. Есть и HttpRouter.provideRequest(layer) для зависимостей, которые обязаны создаваться на каждый запрос отдельно.

Серверный Layer

HttpRouter.serve(routes) отдаёт Layer, которому не хватает одного: собственно сервера, слушающего порт. Его даёт платформенный пакет @effect/platform-node:

// pulse-<nick>/src/cli/commands/serve.ts
import { NodeHttpServer } from '@effect/platform-node';
import { Effect, Layer } from 'effect';
import { Command, Flag } from 'effect/unstable/cli';
import { HttpRouter } from 'effect/unstable/http';
import { createServer } from 'node:http';

import { PulseRoutes } from '../../http/server.ts';

const port = Flag.integer('port').pipe(Flag.withDefault(8080));

const serverLayer = (port: number) =>
  HttpRouter.serve(PulseRoutes).pipe(
    Layer.provide(NodeHttpServer.layer(createServer, { port })),
  );

export const serveCommand = Command.make('serve', { port }, ({ port }) =>
  Effect.gen(function* () {
    yield* Effect.log(`pulse слушает на :${port}`);
    yield* Layer.launch(serverLayer(port));
  }),
);

Что здесь происходит:

  • HttpRouter.serve(appLayer) берёт Layer со всеми зарегистрированными маршрутами и просит у контекста HttpServer.
  • NodeHttpServer.layer(createServer, { port }) создаёт HttpServer-сервис поверх ноудовского http.createServer. port приходит из CLI-флага.
  • Layer.launch запускает Layer и возвращает Effect, который никогда не завершается сам (сервер слушает, пока живёт scope). На SIGINT/SIGTERM scope закроется, finalizer-ы сервера остановят listen, освободят сокеты.
  • Логирование запросов включено по умолчанию. Выключается через HttpRouter.disableLogger, а не включается через middleware.

Если ты собирал API декларативно (Раздел 5), вместо PulseRoutes в serve идёт PulseApiLive: HttpApiBuilder.layer регистрирует свои маршруты в том же самом сервисе-роутере, так что дальше по цепочке разницы нет.

Где здесь MainLive

Маршруты внутри тянут Storage, Sla, MonitorEventsPubSub. Эти сервисы у нас в MainLive из урока 04. Provide делается уровнем выше, в cli/index.ts:

cli.pipe(
  Effect.provide(MainLive),
  Effect.provide(NodeServices.layer),
  NodeRuntime.runMain,
);

MainLive стоит до запуска CLI, поэтому все handler-ы (включая serveCommand) получают Storage, Sla, остальные сервисы. Сервер не знает, откуда они: его задача собрать HTTP-уровень и отдать управление handler-ам.

Жизненный цикл сервера

NodeRuntime root fiber
└─ serveCommand handler fiber
   └─ Layer.launch scope
      ├─ NodeHttpServer.layer: создаёт http.createServer().listen(port)
      │  └─ finalizer: server.close()
      └─ HttpRouter.serve: регистрирует request handler на сервере
         └─ finalizer: снимает регистрацию

На SIGINT root файбер interrupt-ит handler, тот закрывает scope Layer.launch. Finalizer NodeHttpServer вызывает server.close(), который перестаёт принимать новые соединения и ждёт окончания текущих. Когда они закроются (или истечёт grace period), процесс выходит с кодом 0.

Это и есть graceful shutdown бесплатно: ты не пишешь process.on('SIGINT'), не считаешь активные соединения руками, не вызываешь server.close() явно. Всё делает scope.

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

  • Роутер это сервис, маршрут это Layer: HttpRouter.add(method, path, handler) плюс Layer.mergeAll.
  • NodeHttpServer.layer(createServer, { port }) оборачивает ноудовский HTTP-сервер в Layer.
  • HttpRouter.serve(appLayer) подключает маршруты к сервису HttpServer, отдаёт Layer. Для тестов есть HttpRouter.toWebHandler(appLayer), обычный fetch-стиль без сокетов.
  • Layer.launch(layer) запускает Layer и возвращает “вечный” Effect, который живёт, пока живёт scope.
  • MainLive provide-ится до запуска CLI, и HTTP-handler-ы получают доступ к тем же сервисам, что и CLI.
  • Graceful shutdown это finalizer scope, не отдельная инфраструктура.

Раздел 7 · SSE из PubSub, эндпоинт /events

Что такое SSE

Server-Sent Events это простой протокол: сервер отвечает с Content-Type: text/event-stream, держит соединение открытым, и шлёт по нему сообщения в формате:

event: probe-success
data: {"targetId":"github","status":200,"latencyMs":142}

event: probe-failure
data: {"targetId":"example","reason":"timeout"}

Клиент в браузере подписывается через new EventSource('/events') и получает события. Никакого WebSocket, никаких heartbeat-ов, ничего сложного. Идеально подходит для случая “сервер шлёт обновления, клиенту от сервера ничего отвечать не надо”.

В Pulse уже есть MonitorEventsPubSub из урока 07: туда worker-ы публикуют события каждого пробинга. Нам остаётся подписаться на этот pubsub из endpoint-а и стримить события клиенту.

Stream.fromPubSub и HttpServerResponse.stream

// pulse-<nick>/src/http/events.ts
import { Effect, Stream } from 'effect';
import { HttpServerResponse } from 'effect/unstable/http';

import { MonitorEventsPubSub } from '../concurrency/coordination.ts';

const sseEncoder = new TextEncoder();

const encodeSse = (event: { _tag: string }) =>
  sseEncoder.encode(`event: ${event._tag}\ndata: ${JSON.stringify(event)}\n\n`);

export const eventsHandler = Effect.gen(function* () {
  const pubsub = yield* MonitorEventsPubSub;

  const stream = Stream.fromPubSub(pubsub).pipe(Stream.map(encodeSse));

  return HttpServerResponse.stream(stream, {
    contentType: 'text/event-stream',
    headers: {
      'cache-control': 'no-cache',
      connection: 'keep-alive',
    },
  });
}).pipe(Effect.orDie);

Effect.orDie в хвосте нужен потому, что HttpRouter.add ждёт handler без ожидаемых ошибок в канале: всё, что может пойти не так на уровне домена, надо разобрать до этой точки, а необработанное честно превращается в дефект.

Что здесь происходит:

  • Stream.fromPubSub(pubsub) подписывается на pubsub и отдаёт Stream событий. Подписка живёт ровно столько, сколько живёт scope этого Stream. Когда клиент отвалится, scope закроется, подписка снимется.
  • Stream.map(encodeSse) форматирует каждое событие в SSE-фрейм и кодирует в Uint8Array (HttpServerResponse.stream хочет байты).
  • HttpServerResponse.stream(stream, { contentType, headers }) отдаёт HTTP-ответ, тело которого это поток байтов. Сервер сам читает из стрима и пишет в socket, backpressure из урока 09 работает по всей цепочке.
  • Заголовок cache-control: no-cache обязателен для SSE, иначе прокси может задержать буферизацию. connection: keep-alive подсказка прокси не закрывать соединение по таймауту.

Подключение к маршрутам

Самый прямой путь это отдельный маршрут рядом с остальными:

// в src/http/server.ts
export const PulseRoutes = Layer.mergeAll(
  HttpRouter.add('GET', '/status', statusHandler),
  HttpRouter.add('GET', '/api/monitors', monitorsHandler),
  HttpRouter.add('GET', '/events', eventsHandler),
);

Если ты идёшь декларативным путём через HttpApi, у SSE есть свой тип успеха, HttpApiSchema.StreamSse, и endpoint описывается так же, как любой другой:

// в src/http/api.ts
import { MonitorEvent } from '../events.ts';

const eventsGroup = HttpApiGroup.make('events').add(
  HttpApiEndpoint.get('stream', '/events', {
    success: HttpApiSchema.StreamSse({ data: MonitorEvent }),
  }),
);

export const PulseApi = HttpApi.make('pulse').add(monitorsGroup, statusGroup, eventsGroup);

StreamSse({ data: MonitorEvent }) принимает обычную Schema: каждое событие потока кодируется в JSON и уезжает в поле data SSE-фрейма. Это та самая схема MonitorEvent из урока 02, никакой отдельный кодек заводить не нужно.

// в src/http/handlers.ts
const EventsLive = HttpApiBuilder.group(PulseApi, 'events', (handlers) =>
  handlers.handleRaw('stream', () => eventsHandler),
);

export const PulseApiLive = HttpApiBuilder.layer(PulseApi).pipe(
  Layer.provide(MonitorsLive),
  Layer.provide(StatusLive),
  Layer.provide(EventsLive),
);

handlers.handleRaw это вариант для endpoint-ов, которые сами возвращают HttpServerResponse, без оборачивания результата схемой. Стримы попадают сюда.

Закрытие подписки

Главная тонкость SSE: что происходит, когда клиент закрывает вкладку или отваливается по сети.

  • На уровне HTTP сервер видит, что сокет закрылся.
  • NodeHttpServer interrupt-ит файбер, который обслуживает запрос.
  • Stream.fromPubSub сидит в scope этого файбера; на interrupt scope закрывается.
  • На закрытии scope подписка на PubSub чисто снимается.

Внутри MonitorEventsPubSub не остаётся “забытых” подписчиков, которые продолжают копить события. Это та же дисциплина scope-managed ресурсов из урока 05, прокинутая через всю цепочку HTTP.

Тест “три события за 90 секунд”

ДЗ просит проверить, что EventSource видит три события при интервале 30 секунд. Это пишется так:

import { describe, it } from '@effect/vitest';
import { Effect, Fiber, Stream } from 'effect';
import { TestClock } from 'effect/testing';
import { expect } from 'vitest';

import { MainTest } from '../layers/main-test.ts'; // тестовый Layer
import { eventsHandler } from '../src/http/events.ts';

describe('events endpoint', () => {
  it.effect('streams three events in 90 seconds at 30s interval', () =>
    Effect.gen(function* () {
      const response = yield* eventsHandler;
      const body = Stream.decodeText(response.body); // Stream<string>
      const fiber = yield* Effect.forkChild(Stream.runCollect(body));

      yield* TestClock.adjust('90 seconds');

      const parts = yield* Fiber.join(fiber); // Array<string>
      const events = parts.join('').match(/^event: /gm) ?? [];
      expect(events.length).toBe(3);
    }).pipe(Effect.provide(MainTest), Effect.provide(TestClock.layer())),
  );
});

В тестовом MainTest Layer-е Schedule.spaced('30 seconds') управляется TestClock-ом из урока 13: ровно три тика за виртуальные 90 секунд, ровно три события в pubsub, ровно три SSE-фрейма в теле ответа. Никаких реальных таймеров.

Три мелочи, из-за которых старый код тут не скомпилируется. TestClock приезжает Layer-ом из effect/testing. Effect.fork называется Effect.forkChild. И Stream.runCollect отдаёт обычный массив, так что склейка это простой parts.join(''), без Chunk.join.

Есть и способ проверить endpoint целиком, вместе с роутингом и заголовками, не поднимая сокет: HttpRouter.toWebHandler(PulseRoutes) отдаёт функцию из Request в Response в стандартном fetch-стиле.

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

  • Stream.fromPubSub(pubsub) это типовая обёртка над подписчиком, живущая в scope.
  • HttpServerResponse.stream(byteStream, { contentType, headers }) отдаёт chunked-ответ, backpressure уважается.
  • SSE-фрейм это event: name\ndata: json\n\n плюс UTF-8 в байты.
  • Закрытие клиента -> interrupt файбер -> закрытие scope -> снятие подписки, ничего лишнего держать не надо.
  • Тестируется через TestClock и Stream.runCollect, без сетевого слоя.

Раздел 8 · Pulse · один MainLive, два entry-point-а

Граф финального проекта

К концу урока 11 у нас был MainLive Layer со всеми сервисами Pulse: HttpClient, Storage, Logger, Clock, MonitorEventsPubSub, SlaState, Probe, Schedule. Теперь поверх него стоят две поверхности:

                ┌─────────────────────────────────────────────┐
                │                  MainLive                   │
                │                                             │
                │  HttpClient · Storage · Logger · Clock      │
                │  MonitorEventsPubSub · SlaState (TxRef)     │
                │  Probe · Schedule · BatchedDnsResolver      │
                │                                             │
                └─────────────────────────────────────────────┘
                                    ▲
                       ┌────────────┴────────────┐
                       │                         │
            ┌──────────┴──────────┐   ┌──────────┴──────────┐
            │  src/cli/index.ts   │   │  src/http/server.ts │
            │                     │   │                     │
            │  Command.run(...)   │   │  HttpRouter.serve   │
            │  + NodeServices     │   │  + NodeHttpServer   │
            │  + NodeRuntime      │   │  + NodeRuntime      │
            └─────────────────────┘   └─────────────────────┘
                       │                         │
                  pulse watch              :8080/status
                  pulse monitor add        :8080/api/monitors
                  pulse init               :8080/events  (SSE)

Обе ветки Effect.provide(MainLive) подсасывают одно и то же. Различаются только тем, что они подключают поверх Main: CLI тащит NodeServices.layer ради Terminal и Stdio, HTTP-ветка тащит NodeHttpServer.layer ради сервера. Сам Main не знает, какая поверхность его дернула.

Два процесса или один

В простом случае pulse serve и pulse watch это два отдельных процесса: оператор запускает watch в терминале, серверный процесс висит как systemd-юнит. Они не делят память, у каждого свой MainLive Layer со своими Storage и MonitorEventsPubSub.

Когда хочется делить состояние между ними, у тебя два пути:

  • Внешнее хранилище. Storage пишет в Postgres/SQLite/файл, pubsub заменяется на Redis Streams. Тогда два процесса работают над общим стейтом.
  • Один процесс, две поверхности. Серверный процесс запускает и worker, и HTTP-сервер одновременно. pulse serve это команда из CLI, но её handler внутри стартует и Probe-шедулер, и HttpApi-сервер. Тогда terminal-команды pulse monitor add идут через HTTP к этому процессу, а не через локальный Storage.

В нашей серии мы выбираем второй путь: pulse serve запускает полный сервис (worker + HTTP), pulse watch подсоединяется по SSE и рисует таблицу, pulse monitor add шлёт POST на /api/monitors. CLI становится тонким клиентом к HTTP, не “вторым процессом, который независимо живёт”.

Один Runtime для тестов

В тестах это особенно полезно. У тебя одна программа: создаёшь Runtime с MainTest Layer, прогоняешь и CLI-команды, и HTTP-запросы через тот же runtime:

import { describe, it } from '@effect/vitest';
import { Effect, ManagedRuntime } from 'effect';

import { MainTest } from '../layers/main-test.ts';
import { addCommand } from '../src/cli/commands.ts';
import { eventsHandler } from '../src/http/events.ts';

describe('end-to-end', () => {
  it.effect('add via CLI, observe via SSE', () =>
    Effect.gen(function* () {
      const runtime = ManagedRuntime.make(MainTest);
      // ...
    }),
  );
});

Тут стоит обратить внимание на it.effect. Раньше для тестов, которым нужен скоуп, был отдельный it.scoped; теперь любой it.effect уже скоупный, и отдельного варианта не существует. Для теста, которому нужно настоящее время вместо виртуального, есть it.live.

Один Layer, один runtime, два API. В уроке 13 · Тестирование и code style этот паттерн раскрыт детальнее.

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

  • MainLive это общий граф зависимостей, не привязанный ни к CLI, ни к HTTP.
  • Поверх Main стоят два entry-point-а, оба Effect.provide(MainLive), и каждый добавляет свои платформенные Layer-ы.
  • Делить состояние между процессами это отдельная задача (внешнее хранилище), а делить между поверхностями одного процесса бесплатно.
  • В тестах один ManagedRuntime обслуживает и CLI-команды, и HTTP-запросы; TestClock подменяется на уровне Layer.

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

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

ДЗ

Дальше

Следующий урок · 13. Тестирование и code style. Финал серии: TestClock для проверки Schedule и SSE без реальных таймеров, Layer.test для http-handler-ов, dual API, Data, branded types, и чек-лист “когда тащить Effect в проект, а когда хватит neverthrow”.

Контекст из соседних уроков. Layer-граф собирается в 04 · Services и Layer: MainLive это композиция StorageLive, SlaLive, MonitorEventsPubSubLive и остальных. Stream и backpressure в 09 · Stream: SSE это Stream поверх PubSub. Schedule.spaced и retry в 11 · Runtime и Schedule: на нём строится и pulse watch (Schedule.spaced('500 millis')), и сам worker. Транзакционный Sla из 08 · Транзакции: handler /status отдаёт sla.snapshot атомарно.