Раздел 25 · Effect-TS

Транзакции: TxRef, TxQueue, Effect.tx, ABA, тесты транзакций

senior~120 мин

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

Транзакции: TxRef, TxQueue, Effect.tx, ABA, тесты транзакций

Сцена · два кассира и один счёт

В банке один счёт, на нём 1000. Два кассира одновременно принимают перевод: Алиса снимает 700, Боб снимает 500. Каждый по отдельности проверяет: “хватает ли денег”. Каждый видит “да, хватает”. Оба списывают. Баланс уходит в минус 200, банк должен инвестору 200, и никто не понимает, кто виноват.

Это не “ошибка многопоточности”. Это композиция операций над общим состоянием без атомарности границы. У Ref.update атомарность есть. У последовательности read, проверить, decide, write атомарности нет, и вклинивание чужого файбера между шагами это не теоретический риск, а гарантированный сценарий под нагрузкой.

В этом уроке разбираем транзакционную память в Effect: TxRef, TxQueue, TxHashMap, Effect.tx, Effect.txRetry, и почему “перевод денег” пишется одной короткой транзакцией.

Сцена · один язык вместо двух

Если ты читал про транзакционную память в Haskell или в Clojure, ты знаешь классическую конструкцию: отдельная монада STM, свой набор комбинаторов, и мост atomically, через который значение переходит в мир эффектов. В теле транзакции ты пишешь не на том же языке, что в остальной программе: у STM свой gen, свой fail, свои коллекции, и ничего из мира Effect внутрь не проходит.

Effect третьей версии жил ровно так. Перевод сотни со счёта на счёт выглядел вот так:

// Effect v3, отдельный мир STM
const program = Effect.gen(function* () {
  const alice = yield* TRef.make(500);
  const bob = yield* TRef.make(500);

  const transfer = STM.gen(function* () {
    const aliceBalance = yield* TRef.get(alice);
    if (aliceBalance < 100) yield* STM.retry;

    const bobBalance = yield* TRef.get(bob);
    yield* TRef.set(alice, aliceBalance - 100);
    yield* TRef.set(bob, bobBalance + 100);
  });

  yield* STM.commit(transfer); // мост из мира STM в мир Effect
});

STM<A, E, R> был отдельным типом вычислений, синхронным и чистым. Прерваться посередине, как Effect, он не мог, и именно поэтому рантайм имел право перезапускать его сколько угодно раз. Журнал вёлся тот же, что и сейчас: ссылка на TRef, значение на входе в транзакцию, текущее значение внутри неё. На STM.commit рантайм сверял: исходные значения прочитанных TRef всё ещё те же? Если да, записи применялись, если нет, журнал выбрасывался и тело прокручивалось с нуля, и так до успешного коммита.

Плата за такую чистоту одна, и она болезненная: цвет. Внутрь STM нельзя было позвать ни один Effect, даже безобидный Effect.log для отладки. Свой хелпер над Effect внутрь тоже не проходил, для транзакций приходилось писать второй.

В v4 модель поменяли: модули Tx* это тот же STM, только встроенный в Effect. Транзакция теперь обычный Effect.gen, а Effect.tx(...) вокруг него объявляет границу коммита:

// Effect v4, тот же перевод
const program = Effect.gen(function* () {
  const alice = yield* TxRef.make(500);
  const bob = yield* TxRef.make(500);

  const transfer = Effect.tx(
    Effect.gen(function* () {
      const aliceBalance = yield* TxRef.get(alice);
      if (aliceBalance < 100) yield* Effect.txRetry;

      const bobBalance = yield* TxRef.get(bob);
      yield* TxRef.set(alice, aliceBalance - 100);
      yield* TxRef.set(bob, bobBalance + 100);
    }),
  );

  yield* transfer; // границу задал Effect.tx, отдельного commit нет
});

Суть та же, механика под капотом та же, исчез отдельный тип и мост наружу.

Что это меняет на практике. Учить второй набор комбинаторов не нужно: Effect.gen, Effect.forEach, Effect.all, твои собственные хелперы работают внутри транзакции ровно так же, как снаружи. Транзакционные ячейки (TxRef и родня) отдают обычные Effect-ы, поэтому они свободно сочетаются с любым другим кодом.

Плата за это одна, и о ней стоит знать сразу: типы больше не запрещают положить в транзакцию сетевой вызов или запись в лог. Раньше это ловил компилятор, теперь ловишь ты. К этому вернёмся в Разделе 2.

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

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

  1. Раздел 1, Ref и его потолок: атомарна одна операция, не композиция.
  2. Раздел 2, TxRef плюс Effect.tx. Граница транзакции и что в неё нельзя класть.
  3. Раздел 3, как транзакции работают изнутри: журнал в сервисе Effect.Transaction, микротранзакции, оптимистический commit, retry на конфликте.
  4. Раздел 4, Effect.txRetry, condition variable без condition variable.
  5. Раздел 5, транзакционные коллекции: TxQueue, TxChunk, TxHashMap. Композиция “взять из одной, положить в другую” в одной транзакции.
  6. Раздел 6, ABA-проблема и почему транзакции от неё защищены по построению.
  7. Раздел 7, когда транзакция избыточна: Ref.modify хватает, когда хватает.
  8. Раздел 8, Pulse: TxRef<SlaState> для атомарного circuit breaker-а, который переключает основной URL на резервный после трёх провалов подряд.

К концу урока ты пишешь транзакцию из 3 шагов и понимаешь, чем она отличается от 3 операций над Ref, объясняешь, почему сетевой вызов внутри транзакции это баг, и используешь Effect.txRetry как замену опросу в цикле.

Раздел 1 · Ref и где он перестаёт спасать

Что Ref уже даёт

В 06 · Файберы и concurrency и 07 · Координация ты уже видел Ref<A>. Это типобезопасная мутируемая ячейка, и его операции по одной атомарны:

import { Effect, Ref } from 'effect';

const program = Effect.gen(function* () {
  const counter = yield* Ref.make(0);

  // 100 параллельных инкрементов, итог гарантированно 100
  yield* Effect.all(
    Array.from({ length: 100 }, () => Ref.update(counter, (n) => n + 1)),
    { concurrency: 'unbounded' },
  );

  return yield* Ref.get(counter);
});
// program: Effect<number, never, never>, всегда возвращает 100

Ref.update(ref, fn) это атомарный read-modify-write над одной ячейкой. Вклиниться между чтением и записью другому файберу нельзя, рантайм гарантирует.

Ref.modify(ref, fn) сильнее: атомарно читает, вызывает fn(current) => [returnValue, nextState], и записывает nextState. В 07 · Координация этим способом мы атомарно создавали per-domain semaphore, защищаясь от двух одновременных создателей.

Где это перестаёт работать

Перевод денег между двумя счетами это операция над двумя ячейками. Атомарность по одной не помогает, потому что между обновлениями двух разных Ref-ов чужой файбер успевает прочитать промежуточное состояние:

import { Effect, Ref } from 'effect';

const transfer = (from: Ref.Ref<number>, to: Ref.Ref<number>, amount: number) =>
  Effect.gen(function* () {
    const balance = yield* Ref.get(from);
    if (balance < amount) return yield* Effect.fail('insufficient');
    yield* Ref.update(from, (b) => b - amount);
    // *** опасное окно: счёт `from` уже списан, `to` ещё не пополнен ***
    yield* Ref.update(to, (b) => b + amount);
  });

Если в этот момент кто-то третий читает оба счёта и считает их сумму, он увидит “минус amount”. Это не безобидно: на этом строятся отчёты регулятору, проверки баланса в тестах, инварианты, на которые опирается весь остальной код.

Сценарий с двумя кассирами на одном счёте ещё хуже:

// fiber A
const balance = yield* Ref.get(account); // прочитал 1000
if (balance >= 700) {
  // fiber B вклинился: тоже прочитал 1000, тоже увидел "хватает"
  yield* Ref.update(account, (b) => b - 700); // 1000 - 700 = 300
}
// fiber B
yield* Ref.update(account, (b) => b - 500); // 300 - 500 = -200, банк должен

Корень проблемы в том, что read и последующий write разорваны другим эффектом. Атомарности Ref.update хватает только если проверку удаётся свернуть в саму функцию обновления.

Что можно вытащить из Ref.modify

Бывает, что одиночная транзакция всё-таки сводится к одной ячейке: счёт один, from и to это поля одного объекта, пара Map<accountId, balance> хранится целиком в одной Ref<Map>. Тогда Ref.modify спасает:

import { Effect, HashMap, Ref } from 'effect';

type Accounts = HashMap.HashMap<string, number>;

const transferOnRef = (
  ref: Ref.Ref<Accounts>,
  fromId: string,
  toId: string,
  amount: number,
) =>
  Ref.modify(ref, (accounts) => {
    const from = HashMap.get(accounts, fromId);
    const to = HashMap.get(accounts, toId);
    if (from._tag === 'None' || to._tag === 'None') return [{ ok: false as const }, accounts];
    if (from.value < amount) return [{ ok: false as const }, accounts];

    const next = HashMap.set(
      HashMap.set(accounts, fromId, from.value - amount),
      toId,
      to.value + amount,
    );
    return [{ ok: true as const }, next];
  });

Это работает, ровно потому что вся логика проверки и обновления уехала внутрь чистой функции, которую Ref.modify атомарно прокручивает через одну ячейку. У этого подхода два честных недостатка:

  • всё состояние сидит в одной ячейке, и любая транзакция держит её целиком; параллелизма по непересекающимся подмножествам нет;
  • как только понадобится наблюдать изменение (ждать, пока на счёте появится 700), вернёшься к опросу Ref.get в цикле или будешь платить за Deferred на каждое ожидание.

Когда состояние шире одной ячейки или нужно “ждать условия”, приходят транзакции.

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

  • Ref.update и Ref.modify атомарны по одной ячейке.
  • Цепочка из нескольких операций над разными Ref-ами не атомарна, между ними можно вклиниться.
  • Ref.modify спасает, если состояние удаётся собрать в одну ячейку и уложить логику в чистую функцию.
  • Когда состояние нужно держать в нескольких ячейках или ждать условия, Ref упирается в потолок.

Раздел 2 · TxRef и Effect.tx

Сама идея

Транзакционная память переносит идею транзакций из баз данных в память процесса. Программа описывает, что хочет сделать с общим состоянием, рантайм запоминает, какие ячейки прочитаны и какие записаны, в конце атомарно фиксирует изменения, если ни одна из прочитанных ячеек не изменилась снаружи. Если изменилась, транзакция перезапускается с нуля, и так до победы.

Отсюда следует главное правило: внутри транзакции не должно быть ничего, что нельзя откатить. Никакого console.log, никакого fetch, никакой записи в файл. Тело транзакции должно быть расчётом над набором транзакционных ячеек, и весь его результат либо целиком применяется, либо целиком отбрасывается.

TxRef это Ref для транзакций

import { Effect, TxRef } from 'effect';

const program = Effect.gen(function* () {
  const counter = yield* TxRef.make(0);
  // counter: TxRef<number>

  yield* Effect.tx(TxRef.update(counter, (n) => n + 1));

  return yield* TxRef.get(counter);
});

Смотри на типы: TxRef.make, TxRef.get, TxRef.set, TxRef.update, TxRef.modify возвращают обычные Effect-ы. Отдельного транзакционного мира нет, поэтому и мост из него не нужен.

Заметь ещё, что последняя строка обошлась без Effect.tx. Одиночная операция атомарна сама по себе, ровно как Ref.get: заворачивать её в транзакцию не во что. Effect.tx нужен там, где шагов несколько и они должны стать одним.

Effect.tx, граница коммита

Effect.tx берёт обычный эффект и делает из него транзакцию:

import { Effect, TxRef } from 'effect';

const transfer = (from: TxRef.TxRef<number>, to: TxRef.TxRef<number>, amount: number) =>
  Effect.gen(function* () {
    const balance = yield* TxRef.get(from);
    if (balance < amount) return yield* Effect.fail('insufficient' as const);
    yield* TxRef.update(from, (b) => b - amount);
    yield* TxRef.update(to, (b) => b + amount);
  });

const run = Effect.gen(function* () {
  const alice = yield* TxRef.make(1000);
  const bob = yield* TxRef.make(0);

  yield* Effect.tx(transfer(alice, bob, 700));
  // граница: после этой строки изменения применены атомарно
});

Сравни с примером на Ref из Раздела 1. Код почти такой же, за одним отличием: его нельзя разорвать. Вся последовательность из четырёх шагов (get, проверка, два update) видна внешнему миру как одна точка во времени. Любой другой файбер, который смотрит на alice и bob, увидит либо состояние “до перевода”, либо “после”, промежуточное “минус 700, ещё не плюс 700” не существует.

Провал тоже транзакционный: Effect.fail('insufficient') внутри тела откатывает всё, что транзакция успела записать, и выносит ошибку наружу обычным каналом E.

Композиция транзакций

Главное преимущество перед Ref.modify это композиция. Транзакции собираются как лего, и сборка остаётся атомарной:

const transferTwice = Effect.tx(
  Effect.gen(function* () {
    yield* transfer(alice, bob, 100);
    yield* transfer(bob, charlie, 100);
  }),
);
// одна транзакция из шести шагов, все атомарны вместе

Ref.modify так не умеет: три отдельных Ref.modify над тремя ячейками это три отдельные точки атомарности, между ними другой файбер увидит промежуточное состояние.

Отдельно стоит запомнить поведение вложенных Effect.tx. Они не создают вложенную транзакцию: внутренний вызов присоединяется к уже идущей и пользуется её журналом, а границу задаёт самый внешний (механику разберём в Разделе 3). Поэтому оборачивать каждую мелкую операцию в tx безопасно, лишних коммитов не появится:

const inner = Effect.tx(TxRef.set(ref2, 20));

yield* Effect.tx(
  Effect.gen(function* () {
    yield* TxRef.set(ref1, 10);
    yield* inner; // присоединился к внешней транзакции, отдельного коммита нет
  }),
);
// либо записались обе ячейки, либо ни одной

Что не стоит класть внутрь транзакции

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

const broken = Effect.tx(
  Effect.gen(function* () {
    const balance = yield* TxRef.get(account);
    yield* Effect.log(`balance is ${balance}`); // *** так не надо ***
    // если транзакция конфликтует и ретраится, лог напечатается дважды
    yield* TxRef.update(account, (b) => b - 100);
  }),
);

Типы это разрешают, а рантайм повторит тело столько раз, сколько случилось конфликтов. Значит, повторится и лог, и отправленное письмо, и списание с карты. Правило простое: внутри Effect.tx только работа с транзакционными ячейками и чистые вычисления. Всё остальное выносится за границу:

const balance = yield* Effect.tx(transfer(alice, bob, 700));
yield* Effect.log('transfer committed'); // здесь можно, транзакция уже закоммичена

Если держать это правило в голове тяжело, помогает договорённость по именам: функции, предназначенные для тела транзакции, называть с приставкой (txTransfer, txRecordFailure) и не звать из них ничего, кроме TxRef и родни.

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

  • TxRef<A> это Ref<A> для транзакций. Его операции обычные Effect-ы, отдельного транзакционного мира нет.
  • Effect.tx(effect) объявляет границу: тело прокручивается и атомарно коммитится.
  • Одиночная операция над TxRef атомарна и без Effect.tx, оборачивать нужно только композицию.
  • Вложенные Effect.tx присоединяются к внешней транзакции, а не создают свою.
  • Сайд-эффекты внутри транзакции типами не запрещены, но повторятся на каждом ретрае. Держи внутри только ячейки и чистые расчёты.

Раздел 3 · Как это работает изнутри

Снапшот и оптимистический commit

Под капотом всё устроено по классической схеме (статья Harris/Marlow/Peyton Jones/Herlihy 2005, ссылка в resources, идея с тех пор не поменялась):

  1. На входе в транзакцию рантайм заводит журнал: набор прочитанных ячеек и набор записанных, поначалу пустые.
  2. Каждый TxRef.get(ref) записывает версию ref в набор чтений и возвращает текущее значение из набора записей, если оно там есть, иначе из самой ячейки.
  3. Каждый TxRef.set или TxRef.update пишет новое значение в набор записей, не трогая саму ячейку.
  4. На границе Effect.tx рантайм берёт короткий глобальный лок и проверяет: для каждой прочитанной ячейки её текущая версия равна той, что была на чтении? Если да, записи атомарно сливаются в реальные ячейки, лок отпускается. Если нет, журнал выбрасывается, транзакция запускается с нуля.

Никаких блокировок на самих ячейках, никаких очередей ожидания, никаких deadlock. Цена: при высокой конкуренции за одни и те же ячейки транзакции могут много раз ретраиться, и время до коммита растёт. Но обычно такой давки нет, потому что транзакции короткие и наборы пересекающихся ячеек небольшие.

Где лежит журнал · Effect.Transaction

Журнал это не спрятанная механика рантайма, а обычный сервис в Context. Вот он целиком, из исходников effect:

class Transaction extends Context.Service<
  Transaction,
  {
    retry: boolean;
    readonly journal: Map<TxRef<any>, { readonly version: number; value: any }>;
  }
>()('effect/Effect/Transaction') {}

Внутри ровно две вещи: журнал (ячейка, её версия на момент чтения, текущее значение внутри транзакции) и флаг “тело попросило ретрай”.

Effect.tx для транзакции это то же, что Effect.scoped для ресурсов из 05 · Ресурсы. Он создаёт состояние, подставляет сервис в Context на время тела и вычитает требование из R-канала: в сигнатуре это буквально Effect<A, E, Exclude<R, Transaction>>. Scope в R означает “этому эффекту нужна область жизни”, Transaction в R означает “этому эффекту нужна идущая транзакция”.

Микротранзакции и почему Effect.tx не вкладывается

Тогда почему TxRef.get работает и вне всякого Effect.tx, откуда он берёт сервис? Оттуда, что каждая операция над Tx*-ячейкой сама завёрнута в Effect.tx. Вот как устроен TxRef.modify, через который выражены и get, и set, и update:

// исходники effect, TxRef.ts, сокращено
export const modify = (self, f) =>
  Effect.Transaction.pipe(
    Effect.flatMap((state) => Effect.sync(() => { /* работа с журналом */ })),
    Effect.tx, // <- вот она, граница
  );

То есть одиночный TxRef.set это микротранзакция: свой журнальчик на одну запись, мгновенный коммит. Значит вот такой код это не одна транзакция, а четыре подряд:

const transfer = Effect.gen(function* () {
  const aliceBalance = yield* TxRef.get(alice); // микротранзакция 1
  const bobBalance = yield* TxRef.get(bob); // микротранзакция 2

  yield* TxRef.set(alice, aliceBalance - 100); // микротранзакция 3
  yield* TxRef.set(bob, bobBalance + 100); // микротранзакция 4
});

Между любыми двумя из них чужой файбер спокойно вклинивается, и получается ровно та же дыра, что у цепочки Ref.update в Разделе 1. Атомарность есть у каждого шага, у последовательности её нет.

Теперь смотри, что делает внешний Effect.tx. Первым делом он заглядывает в Context текущего файбера:

// исходники effect, Effect.ts, сокращено
export const tx = (effect) =>
  withFiber((fiber) => {
    const state = Context.getOrUndefined(fiber.context, Transaction);
    if (state) return effect; // транзакция уже идёт, присоединяемся к ней
    // иначе заводим журнал и крутим тело до консистентного коммита
  });

Если транзакция уже идёт, tx не создаёт ничего нового и возвращает тело как есть. Журнал заводится только на самом внешнем вызове. Поэтому обёртка снаружи не добавляет пятую транзакцию, а схлопывает четыре микротранзакции в одну: каждая операция найдёт в Context готовый журнал и напишет в него, а коммит случится один, на границе.

Отсюда же и правило про вложенные Effect.tx из Раздела 2: они не вкладываются, они присоединяются.

Что значит “оптимистический”

“Оптимистический” значит “сначала делаем, потом проверяем”. В мире локов было бы наоборот: возьми локи на все участвующие ячейки в правильном порядке, сделай работу, отпусти. Это пессимистично: ты предполагаешь, что конфликт будет, и заранее ставишь барьер.

Транзакции ставят обратное предположение: конфликта обычно нет. Под этим предположением два не пересекающихся transfer (Алиса в Боба, Карл в Диму) идут параллельно без какой-либо координации. Только если они касаются одних и тех же ячеек, проигравший ретраится. На реальных нагрузках это сильно дешевле, чем явные локи.

Почему ретраи безопасны

Ровно настолько, насколько чисто тело транзакции. Если внутри только TxRef и вычисления, у транзакции нет наблюдаемых эффектов до коммита: ретрай повторяет тот же расчёт, и результат не отличить от однократного выполнения. Как только внутрь попал Effect.log или fetch, гарантия ломается, и ломается тихо.

Это та же дисциплина, что у Effect.retry в 03 · Tagged-ошибки: повторять безопасно то, что идемпотентно. Разница в том, что раньше здесь помогала система типов, а теперь помогаешь ты сам.

Чем платим за неявность

Модель v4 удобнее, но у неё есть счета к оплате, и лучше знать о них заранее.

Границы транзакции не видно в типах. Забыть Effect.scoped компилятор не даст: Scope торчит в R и требует, чтобы его кто-то поставил. С транзакциями наоборот, Transaction из R вычищен почти везде, потому что каждая операция сама себе микротранзакция. Единственный эффект, который его показывает, это Effect.txRetry (Раздел 4). У функции над TxRef сигнатура не отличит “это кусок транзакции” от “это самостоятельная операция”, и следить за границами приходится глазами и именами.

Асинхронщина растягивает транзакцию во времени. В v3 это запрещал тип: STM был синхронный, засунуть в него sleep или сетевой вызов было нечем. Теперь yield* Effect.sleep('1 second') внутри Effect.tx компилируется. Секунду тело будет прокручиваться, всё это время журнал держит версии прочитанных ячеек, и шанс, что кто-то запишет в них раньше коммита, растёт вместе с длительностью. Чем длиннее транзакция, тем чаще она проигрывает гонку и уходит на ретрай с нуля.

То же и с планировщиком. Тело Effect.tx крутится в непрерываемой области (в исходниках это uninterruptibleMask), то есть посреди транзакции файбер не прервут. Но непрерываемость это не непрерывность: на каждом асинхронном шаге рантайм вправе отложить продолжение и заняться другими файберами, а когда он вернётся, заранее не известно. Короткая чистая транзакция такого повода не даёт, длинная даёт на каждом шаге. Получается, с чем боролись, на то и напоролись: свободу звать Effect внутри транзакции оплачиваешь ростом числа конфликтов.

Имя Transaction не про STM. В проекте, где рядом живут транзакции базы, оно читается двусмысленно. Заводя свои хелперы, ставь приставку tx в имена (txTransfer, txRecordFailure), это единственная разметка границ, которая у тебя осталась.

Другая сторона счёта: бойлерплейта стало заметно меньше. Второго набора комбинаторов учить не надо, свой txCheck пишется одной строкой (Раздел 4), Effect.forEach и твои собственные хелперы работают внутри транзакции без переходников. Это общее направление, куда движется DX Effect в v4: меньше отдельных миров, больше обычного Effect.

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

  • Транзакция работает оптимистически: журнал чтений и записей, проверка версий на коммите, ретрай при конфликте.
  • Журнал живёт в сервисе Effect.Transaction в Context, а Effect.tx это его Effect.scoped: подставляет сервис и вычитает Transaction из R.
  • Каждая операция над Tx*-ячейкой это микротранзакция; внешний Effect.tx схлопывает их в одну, потому что находит журнал в Context и переиспользует его.
  • Никаких локов на ячейках, никаких deadlock, ретраи дешёвые ровно пока тело чистое и короткое.
  • Цена неявности: границу не видно в типах, а асинхронный шаг внутри транзакции растягивает её и повышает шанс конфликта.

Раздел 4 · Effect.txRetry, ожидание условия

Условное ожидание без polling

Сценарий: воркер хочет снять 700 со счёта Алисы, на счёте сейчас 200. Можно вернуть 'insufficient' и заставить вызывающего ретраить с задержкой. А можно сказать: “подожди, пока на счёте станет хватать”. В классическом мире это condition variable: ждущий файбер засыпает, кто-то делает notify, ждущий просыпается и проверяет условие.

Тут этого не нужно. Достаточно сказать “если условие не выполнено, ретрай”:

import { Effect, TxRef } from 'effect';

const withdrawWhenAvailable = (account: TxRef.TxRef<number>, amount: number) =>
  Effect.tx(
    Effect.gen(function* () {
      const balance = yield* TxRef.get(account);
      if (balance < amount) return yield* Effect.txRetry;
      yield* TxRef.update(account, (b) => b - amount);
    }),
  );

Effect.txRetry это особенное “ничего не получилось, повторите попытку”. Рантайм видит, что транзакция упёрлась в него, и подвешивает её до изменения какой-нибудь из прочитанных ячеек. Как только изменилась (кто-то пополнил счёт), ждущая транзакция автоматически перезапускается с нуля и снова проверяет условие.

Никаких таймеров, никаких опросов в цикле, никаких ручных уведомлений. Producer пишет TxRef.update(account, (b) => b + 1000) и не знает ничего про ждущих, рантайм будит их сам, потому что видел, что они читали account.

Обрати внимание на тип: Effect.txRetry это Effect<never, never, Transaction>. Transaction в требованиях означает “меня можно позвать только внутри транзакции”, а Effect.tx это требование вычитает. Позвать txRetry из обычного кода компилятор не даст, и это единственное, что здесь всё ещё стережёт система типов.

Свой txCheck, сахар над “условие или retry”

Шаблон “если не condition, то retry” встречается часто, и на него удобно завести собственный однострочник:

import { Effect } from 'effect';

const txCheck = (condition: boolean) => (condition ? Effect.void : Effect.txRetry);

Дальше он читается естественно, “проверь, что баланса хватает, иначе подожди”:

const withdrawWhenAvailable = (account: TxRef.TxRef<number>, amount: number) =>
  Effect.tx(
    Effect.gen(function* () {
      const balance = yield* TxRef.get(account);
      yield* txCheck(balance >= amount);
      yield* TxRef.update(account, (b) => b - amount);
    }),
  );

Это как раз пример того, ради чего “один язык вместо двух” и затевался: свой комбинатор для транзакций пишется обычной функцией над Effect, никакого отдельного набора примитивов заводить не пришлось.

Цепочки условий

Композиция работает и здесь. Можно написать “перевод между двумя счетами с ожиданием обоих условий”:

const transferWhenReady = (
  from: TxRef.TxRef<number>,
  to: TxRef.TxRef<number>,
  amount: number,
  ceiling: number,
) =>
  Effect.tx(
    Effect.gen(function* () {
      const fromBalance = yield* TxRef.get(from);
      const toBalance = yield* TxRef.get(to);
      yield* txCheck(fromBalance >= amount); // ждём, пока хватает у from
      yield* txCheck(toBalance + amount <= ceiling); // и пока to не переполнен
      yield* TxRef.update(from, (b) => b - amount);
      yield* TxRef.update(to, (b) => b + amount);
    }),
  );

Если хоть одно условие не выполнено, транзакция спит. На любом изменении from или to (рантайм знает, потому что они прочитаны) транзакция перезапускается, оба условия проверяются заново. Состояние “одно условие выполнено, второе нет, висим в полусостоянии” невозможно.

Прерывание ждущей транзакции

Граница транзакции это точка приостановки. Если файбер, который ждёт условия, получает Fiber.interrupt, рантайм аккуратно его сворачивает: журнал отбрасывается, ничего не остаётся в подвешенном состоянии. Это та же дисциплина, что и у Deferred.await или Queue.take (см. 07 · Координация): транзакция атомарно либо случилась, либо нет, никаких “почти случилась”.

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

  • Effect.txRetry ставит транзакцию в режим “жду изменения прочитанных ячеек”; рантайм будит её сам.
  • Его тип требует Transaction, поэтому вне Effect.tx он не скомпилируется.
  • Свой txCheck(cond) пишется одной строкой, потому что транзакция это обычный Effect.
  • Polling и condition variable не нужны: пишешь условие как часть транзакции, дальше работа рантайма.
  • Прерывание ждущей транзакции чистое, рантайм знает, что записи ещё не применены.

Раздел 5 · TxQueue, TxChunk, TxHashMap, транзакционные коллекции

Зачем транзакционные коллекции

Чтобы атомарно “взять из одной структуры и положить в другую”. На обычных Queue (см. 07 · Координация) последовательность Queue.take(in) плюс Queue.offer(out) это два эффекта. Между ними в out может уже что-то прилететь, или in снова заполнится, или сам файбер может быть прерван между шагами.

В транзакции это пишется одним блоком:

import { Effect, TxQueue } from 'effect';

const forward = (input: TxQueue.TxQueue<number>, output: TxQueue.TxQueue<number>) =>
  Effect.tx(
    Effect.gen(function* () {
      const value = yield* TxQueue.take(input); // ждём элемента в input
      yield* TxQueue.offer(output, value);
    }),
  );

// в основном коде
const program = Effect.gen(function* () {
  const a = yield* TxQueue.bounded<number>(8);
  const b = yield* TxQueue.bounded<number>(8);

  yield* Effect.forkChild(forward(a, b).pipe(Effect.forever));
});

TxQueue.take на пустой очереди уходит в ретрай (помнишь Раздел 4: take на пустом это и есть “ждать условия”). Когда в input приходит элемент, транзакция просыпается, берёт его и атомарно кладёт в output. Если в этот момент output полон, весь forward ретраится, элемент возвращается в input. Программа никогда не увидит “взял из input и потерял, не успев положить в output”.

Заметь: конструкторы коллекций тоже обычные эффекты, yield* TxQueue.bounded(8) без всякого моста.

TxQueue API

import { TxQueue } from 'effect';

const make = TxQueue.bounded<string>(64); // Effect<TxQueue<string>>
// также: TxQueue.unbounded(), TxQueue.sliding(N), TxQueue.dropping(N)

TxQueue.offer(q, 'hello'); // на полной bounded уходит в ретрай
TxQueue.take(q);           // на пустой уходит в ретрай
TxQueue.size(q);
TxQueue.peek(q);           // посмотреть без снятия

Семантика стратегий та же, что у обычной Queue. Есть и знакомые по 07 · Координация операции завершения: TxQueue.end, TxQueue.fail, TxQueue.awaitCompletion. Ключевое отличие от обычной очереди в том, что все эти операции можно собирать в составные транзакции, которые либо целиком случаются, либо целиком нет.

TxChunk, транзакционная последовательность

import { Chunk, Effect, TxChunk } from 'effect';

const program = Effect.gen(function* () {
  const chunk = yield* TxChunk.make(Chunk.fromIterable([10, 20, 30]));
  const first = yield* TxChunk.get(chunk, 0); // 10
  yield* TxChunk.append(chunk, 40);
  const size = yield* TxChunk.size(chunk); // 4
});

TxChunk<A> держит последовательность в транзакционной ячейке и даёт над ней привычный набор операций. Транзакции, которые его не трогают, с ним и не конфликтуют.

TxHashMap, словарь

import { Effect, TxHashMap } from 'effect';

const program = Effect.gen(function* () {
  const map = yield* TxHashMap.empty<string, number>();
  yield* TxHashMap.set(map, 'apples', 10);
  yield* TxHashMap.set(map, 'oranges', 5);
  const apples = yield* TxHashMap.get(map, 'apples'); // Option<10>
  yield* TxHashMap.remove(map, 'oranges');
  const size = yield* TxHashMap.size(map); // 1
});

На уровне идиом это закрывает большинство сценариев “общий словарь между файберами”: кеш, сессия пользователя, индекс по id. Отдельно полезен TxHashMap.modifyAt(map, key, fn), атомарное “прочитать по ключу, посчитать, записать обратно”, тот самый паттерн, ради которого в 07 · Координация приходилось изворачиваться через Ref.modify.

Остальное семейство

Tx-модулей больше, и все они устроены одинаково: обычные эффекты, атомарность на границе Effect.tx.

МодульЧто это
TxRefодна ячейка
TxQueueочередь
TxPubSubшина с fan-out
TxHashMap, TxHashSetсловарь и множество
TxChunkпоследовательность
TxSemaphoreсчётчик разрешений
TxDeferredодноразовый канал
TxSubscriptionRefячейка, на изменения которой можно подписаться потоком

Правило выбора простое: если значение участвует в составных атомарных операциях, бери Tx-версию; если нет, обычной хватит.

Композиция: TxQueue плюс TxHashMap в одной транзакции

import { Effect, TxHashMap, TxQueue } from 'effect';

type Job = { id: string; payload: string };

const dequeueAndIndex = (
  jobs: TxQueue.TxQueue<Job>,
  index: TxHashMap.TxHashMap<string, Job>,
) =>
  Effect.tx(
    Effect.gen(function* () {
      const job = yield* TxQueue.take(jobs);
      yield* TxHashMap.set(index, job.id, job);
    }),
  );

“Возьми job из очереди и положи в индекс”. Атомарно. Если кто-то параллельно делает TxHashMap.get(index, id), он увидит либо состояние “до взятия из очереди” (job-а в индексе нет, в очереди есть), либо “после взятия” (job-а в очереди нет, в индексе есть). Промежуточное состояние не существует.

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

  • TxQueue, TxChunk, TxHashMap: транзакционные версии очереди, последовательности, словаря. Есть и остальные, от TxSemaphore до TxSubscriptionRef.
  • Их операции это обычные Effect-ы, включая конструкторы.
  • TxQueue.take на пустой и TxQueue.offer на полной уходят в ретрай, рантайм будит на изменении.
  • Атомарность распространяется на композицию: “take из одной, set в другую” внутри одного Effect.tx.

Раздел 6 · ABA-проблема и почему транзакции от неё защищены

Что такое ABA

ABA-проблема это классическая ловушка lock-free структур данных, основанных на CAS (compare-and-swap). Поток T1 читает в переменной значение A, идёт что-то делать, готовится сделать CAS(A -> C). За это время поток T2 успевает переписать значение: A -> B -> A. Когда T1 наконец вызывает CAS(A -> C), операция успешно проходит: текущее значение действительно A, и оно меняется на C. С точки зрения T1 ничего не происходило. На самом деле произошло два события, которые могли сломать инварианты, на которые T1 опирался.

Классический пример: lock-free стек на единственном указателе на голову. T1 читает голову (узел A). Хочет вытолкнуть A из стека и вернуть его. За это время T2 вытолкнул A, вытолкнул B, и положил A обратно. T1 делает CAS(head, A, A.next), где A.next это указатель, прочитанный T1 в самом начале. CAS проходит. Но A.next уже не имеет смысла: теперь он указывает на узлы старой подцепи, которой давно нет в стеке. Структура повреждена.

Почему транзакции это не задевает

Здесь нет CAS на отдельной ячейке. Есть версия каждой прочитанной ячейки в журнале транзакции и проверка всех этих версий на коммите.

Применительно к нашему примеру: T1 в своей транзакции делает TxRef.get(head) и получает A. В журнал рантайм записывает: “я прочитал head версии 17”. Дальше T2 в своих транзакциях делает head: A -> B, потом head: B -> A, каждая со своим коммитом. После каждого коммита версия head инкрементируется: была 17, стала 18, потом 19. Когда T1 пытается закоммититься, рантайм проверяет: “версия head сейчас 19, у меня в журнале 17, конфликт”. Транзакция T1 ретраится с нуля, всё хорошо.

ABA здесь невозможна, потому что сравнивается не значение ячейки, а её версия. Любая запись инкрементирует версию, даже если значение совпало с прежним. У T1 нет шанса “не заметить” две операции, потому что она следит за версиями, не за значениями.

Это не “защита, которая работает в большинстве случаев”. Это базовая семантика транзакционной памяти, выводимая из её определения. Тем самым она избавляет тебя от целого класса нетривиальных багов, на которые написан отдельный том литературы по lock-free программированию.

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

  • ABA-проблема возникает там, где CAS сравнивает значения и не различает “значение совпало случайно”.
  • Транзакция сравнивает версии прочитанных ячеек, не значения, поэтому ABA в ней невозможна по построению.
  • Это снимает целый класс багов, без необходимости вручную управлять памятью снятых узлов и согласовывать их освобождение между потоками.

Раздел 7 · Когда транзакция, когда Ref хватает

У транзакции есть цена

  • транзакции под высокой конкуренцией за ячейки могут много раз ретраиться, тогда время до коммита растёт линейно;
  • журнал чтений и записей это аллокации в куче на каждую транзакцию, для очень коротких операций это заметный оверхед;
  • сайд-эффект, случайно попавший внутрь, повторится на ретрае, и типы об этом не предупредят.

Под “очень короткими операциями” имеется в виду что-то вроде counter += 1. Один атомарный инкремент через Ref.update это одна-две машинных инструкции под капотом, через транзакцию это аллокация журнала и проверка версий. Для миллиона инкрементов в секунду разница ощутима.

Эвристика выбора

Бери Ref, когда:

  • состояние помещается в одну ячейку;
  • операции это короткие read-modify-write, без условного ожидания;
  • Ref.modify решает задачу за одну атомарную функцию.

Бери транзакцию, когда:

  • состояние шире одной ячейки и нужна атомарность по их пересечению;
  • нужно ожидание условия (Effect.txRetry);
  • нужна композиция атомарных операций (несколько transfer в одной транзакции, “взять из очереди и положить в карту”);
  • код на Ref.modify начинает превращаться в кашу из вложенных проверок.

В Pulse мы дальше используем оба: счётчики времени до следующей проверки сидят в Ref, а SlaState (несколько связанных полей плюс правило “после трёх провалов”) в TxRef.

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

  • Транзакция это не замена Ref, это другой инструмент.
  • Под высокой конкуренцией за одни ячейки транзакция может ретраиться много раз, аллокации не нулевые.
  • Эвристика: одна ячейка и простой modify это Ref; несколько ячеек, ожидание условия, композиция это транзакция.

Раздел 8 · Pulse · TxRef<SlaState> и атомарный circuit breaker

Сцена · переключение URL после трёх провалов

В Pulse у каждого таргета есть основной URL и резервный. Правило: если основной URL подряд три раза вернул 5xx или таймаут, переключаемся на резервный. На любом успехе счётчик провалов обнуляется. Любой файбер, который смотрит “какой сейчас активный URL”, должен видеть согласованное состояние: либо “всё ок, основной”, либо “переключились на резервный, провалов 0”.

Это идеальная задача для транзакции: одно состояние с инвариантом, несколько файберов читают и пишут параллельно, правило перехода зависит от текущего значения.

Шаг 1 · SlaState и TxRef

// pulse-<nick>/src/services/sla.ts
export type SlaState = {
  readonly active: 'primary' | 'fallback';
  readonly consecutiveFailures: number;
};

const FAILURE_THRESHOLD = 3;

active это какой URL сейчас в работе, consecutiveFailures это счётчик подряд идущих провалов на активном URL.

Шаг 2 · сервис Sla на транзакциях

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

export class Sla extends Context.Service<Sla>()('Pulse/Sla', {
  make: Effect.gen(function* () {
    const ref = yield* TxRef.make<SlaState>({ active: 'primary', consecutiveFailures: 0 });

    const recordFailure = Effect.gen(function* () {
      const current = yield* TxRef.get(ref);
      const failures = current.consecutiveFailures + 1;
      const next: SlaState =
        failures >= FAILURE_THRESHOLD && current.active === 'primary'
          ? { active: 'fallback', consecutiveFailures: 0 }
          : { ...current, consecutiveFailures: failures };
      yield* TxRef.set(ref, next);
      return next;
    });

    const recordSuccess = Effect.gen(function* () {
      const current = yield* TxRef.get(ref);
      if (current.consecutiveFailures === 0) return current;
      const next: SlaState = { ...current, consecutiveFailures: 0 };
      yield* TxRef.set(ref, next);
      return next;
    });

    return {
      recordFailure: Effect.tx(recordFailure),
      recordSuccess: Effect.tx(recordSuccess),
      snapshot: TxRef.get(ref),
    };
  }),
}) {
  static readonly layer = Layer.effect(this, this.make);
}

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

  • TxRef<SlaState>, инициализированный “основной, 0 провалов”, лежит в scope сервиса.
  • recordFailure это тело транзакции, обычный Effect.gen: читает текущее состояние, считает новый счётчик, решает, переключать активный URL или увеличить счётчик. На переключении сразу же сбрасывает счётчик. На выходе отдаёт новое состояние, чтобы вызывающий мог его залогировать снаружи.
  • recordSuccess сбрасывает счётчик. Маленькая оптимизация: если уже 0, ничего не записываем (чтение всё равно попадёт в журнал, но набор записей остаётся пустым, и коммит отрабатывает быстро).
  • Наружу сервис отдаёт эффекты, уже завёрнутые в Effect.tx. Вызывающий не знает, что внутри транзакция, он просто зовёт sla.recordFailure.
  • А snapshot это голый TxRef.get(ref), без Effect.tx: одиночное чтение атомарно само по себе, заворачивать нечего.

Шаг 3 · использование в worker-е

import { Effect, Result } from 'effect';

import { MonitorEvents } from './monitor-events.ts';
import { HttpService } from './http.ts';
import { Sla } from './sla.ts';

const probeOnce = (target: { id: string; primaryUrl: string; fallbackUrl: string }) =>
  Effect.gen(function* () {
    const sla = yield* Sla;
    const http = yield* HttpService;
    const bus = yield* MonitorEvents;

    const state = yield* sla.snapshot;
    const url = state.active === 'primary' ? target.primaryUrl : target.fallbackUrl;

    const outcome = yield* http.get(url).pipe(Effect.result);

    if (Result.isSuccess(outcome) && outcome.success.status < 500) {
      const next = yield* sla.recordSuccess;
      yield* bus.publish({
        _tag: 'ProbeSuccess',
        targetId: target.id,
        url,
        status: outcome.success.status,
      });
      return next;
    }

    const next = yield* sla.recordFailure;
    yield* bus.publish({ _tag: 'ProbeFailure', targetId: target.id, url, sla: next });
    return next;
  });

Worker берёт snapshot (одно атомарное чтение), уходит в HTTP (наружу, за границей транзакции), и по результату вызывает либо recordSuccess, либо recordFailure. Каждое из них одна атомарная транзакция. Между ними другие worker-ы могут параллельно записывать свои результаты, и никто никого не забивает.

Обрати внимание, что сам HTTP-запрос стоит между транзакциями, а не внутри. Это ровно то правило из Раздела 2: типы бы его не остановили, но при ретрае запрос ушёл бы повторно.

Снэпшот тут осознанно “устаревший”: между snapshot и реальным запросом ещё может прийти событие, которое переключит URL. Это нормально: HTTP-запрос уже в полёте, отменять его по такому поводу незачем. На следующем тике worker возьмёт свежий snapshot и пойдёт по новому URL.

Шаг 4 · подписка на изменения через Effect.txRetry

Иногда хочется не “записать новое состояние и забыть”, а “ждать, пока активный URL поменяется”, чтобы залогировать переключение или поднять алерт. Это пишется через дополнительный файбер, который сидит в транзакции с проверкой:

import { Effect, TxRef } from 'effect';

const watchSwitchOver = (
  ref: TxRef.TxRef<SlaState>,
  publish: (event: { _tag: 'SlaSwitched'; from: 'primary' | 'fallback' }) => Effect.Effect<void>,
) =>
  Effect.gen(function* () {
    let lastSeen: 'primary' | 'fallback' = 'primary';
    while (true) {
      const next = yield* Effect.tx(
        Effect.gen(function* () {
          const state = yield* TxRef.get(ref);
          if (state.active === lastSeen) return yield* Effect.txRetry;
          return state.active;
        }),
      );
      yield* publish({ _tag: 'SlaSwitched', from: lastSeen });
      lastSeen = next;
    }
  });

Транзакция читает active, проверяет “отличается ли от последнего, что я видел”, и если нет, уходит в Effect.txRetry. Рантайм подвешивает её до изменения active. Никаких таймеров, никаких опросов в цикле, никаких отдельных каналов уведомлений: транзакция сама “слушает” интересную ей ячейку. Заметь и здесь границу: publish стоит после Effect.tx, а не внутри.

Если такой сценарий нужен часто, посмотри на TxSubscriptionRef: это ячейка, у которой изменения сразу доступны как поток, и цикл выше сворачивается в Stream.

Этот файбер запускается через Effect.forkScoped (см. 06 · Файберы и concurrency) в MainLive, и при runtime.dispose() чисто прерывается.

Шаг 5 · тест на правило “три провала”

Для Effect-ового кода используем @effect/vitest. Тест на правило выглядит так:

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

import { Sla } from '../src/services/sla.ts';

describe('Sla', () => {
  it.effect('switches to fallback after three consecutive failures', () =>
    Effect.gen(function* () {
      const sla = yield* Sla;

      const after1 = yield* sla.recordFailure;
      expect(after1.active).toBe('primary');
      expect(after1.consecutiveFailures).toBe(1);

      const after2 = yield* sla.recordFailure;
      expect(after2.consecutiveFailures).toBe(2);

      const after3 = yield* sla.recordFailure;
      expect(after3.active).toBe('fallback');
      expect(after3.consecutiveFailures).toBe(0);
    }).pipe(Effect.provide(Sla.layer)),
  );

  it.effect('success resets the counter', () =>
    Effect.gen(function* () {
      const sla = yield* Sla;
      yield* sla.recordFailure;
      yield* sla.recordFailure;
      const afterSuccess = yield* sla.recordSuccess;
      expect(afterSuccess.consecutiveFailures).toBe(0);
      expect(afterSuccess.active).toBe('primary');
    }).pipe(Effect.provide(Sla.layer)),
  );
});

Две мелочи по тестовому обвесу. expect теперь импортируется из самого vitest, @effect/vitest его больше не переэкспортирует. И it.effect уже даёт scope, отдельного it.scoped нет.

Такие тесты тривиальны, потому что транзакция атомарна и детерминирована. Никаких “запусти пять раз и проверь, что обычно срабатывает”: один запуск, один результат, одна гарантия.

В 13 · Testing собран паттерн на конкурентные тесты: десять параллельных recordFailure, проверка инварианта “счётчик равен числу провалов или на переключении начался с 0”.

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

  • TxRef<State> плюс пара методов recordX это полный circuit breaker на транзакциях, 30 строк.
  • Транзакция включает чтение, расчёт, запись, и решение “переключать ли активный URL”. Атомарно.
  • HTTP-запрос стоит между транзакциями, а не внутри: на ретрае он ушёл бы повторно.
  • Внешний наблюдатель через Effect.txRetry подписывается на изменения без опроса в цикле.
  • Тесты простые, потому что транзакция детерминирована: один запуск, один результат.
  • Граница транзакции спрятана внутри сервиса; вызывающий получает обычные Effect-ы.

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

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

ДЗ

Дальше

Следующий урок · 09 · Stream и Sink. После транзакционной памяти удобно взять стримы: атомарные шаги внутри транзакции, потом стрим из событий наружу. Стримы это та же дисциплина producer/consumer из 07 · Координация, только декларативная и с map/filter/groupBy/mapEffect поверх.

Контекст по Ref живёт в 07 · Координация и 06 · Файберы и concurrency. Layer-граф сервисов в 04 · Services и Layer: Sla встанет в Layer наряду с MonitorEvents, UrlQueue, DomainLimiter. Тесты на транзакции подробнее в 13 · Testing.