Перейти к содержимому

Практическое руководство Chronacta

Этот документ сочетает концептуальное введение (Event Sourcing, DDD, подписки, проекции, регистр схем), сквозной пример заказа, паттерны интеграции и практическое руководство по установке, эксплуатации и интеграции Chronacta v1.0.


Chronacta — серверная база на Go с хранением событий (event sourcing). Она ведёт неизменяемый журнал: события группируются в потоки, к каждому потоку можно записывать, читать, подписываться, строить проекции и управлять сервером через сетевой API.

Chronacta хорошо подходит для предметных областей, где важны:

  • полная история изменений (заказы, платежи, счета, бизнес-процессы);
  • аудит и воспроизведение («что произошло и в каком порядке»);
  • интеграция между сервисами через общий журнал событий;
  • восстановление состояния из событий после сбоев или при миграции.
Это не… Почему
Реляционная БД нет SQL, JOIN, произвольных UPDATE по строкам
Универсальная очередь события — главный и постоянный источник данных, а не одноразовые сообщения в очереди
Замена Kafka/NATS «как есть» Chronacta — надёжный журнал событий с нумерацией по потокам; NATS/Kafka — шина для доставки сообщений во внешние системы
CRUD-хранилище текущего состояния текущее состояние выводится из истории событий
  • хранилище «только добавление»: события не перезаписываются; на диске — журнал упреждающей записи (WAL) и файловые сегменты;
  • потоки, порядковые номера событий внутри потока и сквозная нумерация по всему серверу;
  • защита от одновременной записи в один поток и ключи идемпотентности;
  • чтение одного потока и общего журнала $all;
  • подписки на время сессии и постоянные подписки с сохранением позиции;
  • регистр JSON-схем для формата данных событий;
  • проекции — готовые сводки по событиям (встроенные, CEL, Starlark);
  • резервное копирование, экспорт и импорт (JSONL и бинарный формат);
  • основной API — gRPC; REST и WebSocket — дополнительные шлюзы;
  • CLI, веб-админка, TLS, роли и права, OIDC, мультитенантность, метрики Prometheus/OpenTelemetry;
  • опциональный HA-кластер: один узел принимает записи, остальные реплицируют.

Полный набор операций — через gRPC и Go SDK. REST-клиенты для TypeScript и Python покрывают основные сценарии, но не все методы gRPC.


CRUD-модель (классическая БД):

Таблица orders:
id=42, status=Paid, total=1500 ← хранится только «сейчас»

При изменении строка перезаписывается. Старое значение теряется, если отдельно не ведётся журнал аудита.

Event Sourcing:

Поток order-42:
v1 OrderCreated {"order_id":"42","total":1500}
v2 OrderItemAdded {"sku":"A1","qty":2}
v3 OrderPaid {"payment_id":"pay-9"}

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

Свойство Что это даёт
Неизменяемость после фиксации событие не меняется — надёжный аудит
Только добавление новые данные — только новым событием; исправление = ещё одно событие
Запросы ко времени можно восстановить состояние агрегата на любой момент истории потока
Повторная обработка тот же набор событий можно прогнать заново (новая проекция, новый потребитель)
Разделение записи и чтения запись идёт в журнал; чтение — через проекции или прямое чтение потока
Команда в вашем сервисе
→ проверка бизнес-правил
→ формирование доменного события
→ запись в Chronacta (указываете поток и сколько событий уже было в нём)
→ подписчики, проекции и интеграции обрабатывают событие асинхронно

Chronacta хранит и доставляет события. Бизнес-правила агрегата задаются в коде вашего сервиса — Chronacta не знает, что такое «Order» или «Invoice», он знает stream_id, event_type и data.

Хорошие сценарии:

  • домен с богатой историей (финансы, заказы, соответствие требованиям);
  • несколько представлений для чтения из одного журнала событий (отчёты, поиск, аналитика);
  • интеграционные события между командами/сервисами;
  • отладка инцидентов («покажи все события по сущности X»).

Слабые сценарии:

  • простой CRUD без истории (справочник из 50 строк);
  • частые «забыть всё и перезаписать» без доменного смысла;
  • сверхбыстрое key-value-хранилище с миллионами мелких ключей, где не нужна история событий.

CQRS (разделение команд и запросов) разделяет:

  • Сторона команд — принимает команды и меняет состояние через события;
  • Сторона запросов — отдаёт данные для интерфейса и API, часто из проекций.

Chronacta — это хранилище событий + механизмы чтения и реакции:

CQRS-роль В Chronacta
Хранилище событий (запись) запись в поток, $all
Модель для чтения результат проекции или ваш сервис читает поток и держит кэш
Интеграция постоянная подписка, ваш фоновый обработчик

CQRS не обязателен для event sourcing, но на практике почти всегда появляется хотя бы одна модель для чтения или потребитель событий.


Domain-Driven Design (DDD) — способ моделировать сложный бизнес-домен. Event Sourcing часто используют вместе с DDD, но это разные идеи: DDD — про язык и границы домена, ES — про способ хранения.

DDD-понятие Смысл Как выражается в Chronacta
Domain Event факт, который уже произошёл в домене запись в поток: event_type + JSON data
Aggregate (агрегат) группа связанных объектов с общими правилами целостности обычно один поток на экземпляр агрегата, напр. order-42
Aggregate Root точка входа для изменений агрегата ваш сервис записывает события только в «свой» поток
Bounded Context (ограниченный контекст) граница модели и предметного языка префиксы имён потоков, отдельные сервисы, мультитенантность t/{tenant}/…
Ubiquitous Language общие термины команды и кода имена event_type: OrderCreated, не EvtType1
Repository (репозиторий) загрузка и сохранение агрегата ваш код: читает поток и сворачивает события; сохраняет — записывает новые
Integration Event событие для другого контекста отдельный поток или общий журнал $all + подписка

Идентификатор потока: order-42 (один агрегат = один поток).

События домена:

{"event_type":"OrderCreated","data":{"customer_id":"c1","currency":"RUB"}}
{"event_type":"OrderLineAdded","data":{"sku":"BOOK-1","qty":1,"price":900}}
{"event_type":"OrderPaid","data":{"payment_ref":"pay-771"}}

Правила целостности (в коде вашего сервиса, не в Chronacta):

  • нельзя добавить строку после OrderCancelled;
  • OrderPaid только если сумма строк совпадает с ожидаемой;
  • повторная оплата отклоняется.

Конкурентная запись. Два запроса одновременно меняют один заказ. Оба прочитали поток и видят три события. Оба пытаются добавить четвёртое, сообщая серверу: «сейчас в потоке должно быть ровно три события» (в CLI: -expected-version 3). Первый запрос успевает — в потоке становится четыре события. Второй получает ошибку «версия не совпадает», заново читает поток и решает, выполнять ли команду ещё раз.

Ограниченный контекст и именование потоков

Заголовок раздела «Ограниченный контекст и именование потоков»

Плохо:

stream: data
stream: events
stream: temp

Хорошо:

order-42 ← контекст «Sales / Order»
invoice-9001 ← контекст «Billing / Invoice»
t-acme/order-42 ← мультитенантность: арендатор acme, локальное имя order-42

Имя потока — часть соглашения между командами. Зафиксируйте правила в своём репозитории (не обязательно в документации Chronacta).

Команда Событие
Время намерение сделать факт, что случилось
Пример PlaceOrder OrderPlaced
Где живёт HTTP/gRPC-обработчик вашего сервиса поток Chronacta
Отклонение «нельзя — нет товара» событие не пишется

Chronacta принимает только события (запись в поток). Валидация команды — до записи.

1. Клиент вызывает POST /orders { ... }
2. OrderService читает все события потока order-{id} и восстанавливает заказ
3. Order.Place(...) проверяет бизнес-правила
4. OrderService записывает OrderCreated и OrderLineAdded, указывая текущее число событий в потоке
5. Подписка «shipping» получает OrderCreated и резервирует товар на складе
6. Проекция orders-by-status обновляет список «New / Paid / Shipped»
7. Интерфейс показывает данные из проекции или вашего API для чтения

Поток — упорядоченная последовательность событий с монотонной версией в потоке (stream_version: 1, 2, 3, …).

  • один поток = один логический «журнал» (часто один агрегат);
  • версии плотные в рамках потока (без пропусков после успешной записи);
  • параллельные записи в разные потоки не конфликтуют.

Каждое событие содержит (назначает сервер):

Поле Назначение
event_id UUID, уникальный идентификатор
event_type строка доменного типа
data JSON-данные события
stream_id поток
stream_version позиция в потоке
global_position позиция в общем журнале $all
created_at время фиксации на сервере
schema_name / schema_version необязательно; связь с регистром схем

Помимо версии внутри потока, каждое событие получает global_position — порядок фиксации на всём кластере.

read-all / $all нужен когда:

  • проекция слушает все потоки (сквозной обзор);
  • интеграция строит общую хронологию;
  • нужен строгий общий порядок (в HA у реплик чтение может отставать).

Ожидаемая версия (защита от одновременных изменений)

Заголовок раздела «Ожидаемая версия (защита от одновременных изменений)»

При записи вы сообщаете серверу, сколько событий уже есть в потоке. Если за это время кто-то успел записать раньше, число не совпадёт — запись отклонится, и вы перечитаете поток заново. Так не теряются изменения при параллельной работе.

Значение -expected-version Что означает
-2 поток ещё не существует (создаём новый)
-1 не проверять версию (осторожно: возможны гонки)
N ≥ 0 в потоке сейчас должно быть ровно N событий

Проще говоря: «я видел N событий и на их основе добавляю следующее».

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


После успешной записи событие не редактируется и не удаляется через обычный API. Исправление — новое событие (OrderAddressCorrected), компенсация — OrderRefunded. Физическое удаление старых записей — только операторская процедура очистки (scavenge) с подтверждением и проверкой, что потребители не пострадают.

Несколько событий в одном запросе проходят единую цепочку фиксации:

проверка → ожидаемая версия → номера событий → WAL → сегмент на диске → fsync (по политике) → индексы → уведомление подписчиков и проекций

Параметр CHRONACTA_SYNC_WRITES задаёт, ждёт ли клиент принудительной записи на диск перед подтверждением.

В режиме mode=cluster только лидер принимает записи. Реплики копируют журнал WAL. Подтверждение клиенту — после записи на лидере и согласования с кворумом; иначе при падении лидера до репликации данные могут потеряться.

Механизм Вопрос, на который отвечает
Чтение потока «покажи историю агрегата X»
Чтение $all «покажи общий хронологический журнал»
Проекция «какое готовое состояние или сводку мы уже посчитали из журнала?»
Подписка «уведоми меня, когда появится новое событие (и запомни позицию)»

Доставка «как минимум один раз» и идемпотентность

Заголовок раздела «Доставка «как минимум один раз» и идемпотентность»

Постоянная подписка гарантирует доставку как минимум один раз: после nack или сбоя до ack событие придёт снова. Ваш обработчик обязан быть идемпотентным (дедупликация по event_id, идемпотентность внешних вызовов или таблица «уже обработано»).


6. Подписки, проекции и регистр схем: зачем и когда что использовать

Заголовок раздела «6. Подписки, проекции и регистр схем: зачем и когда что использовать»

Это три разные подсистемы. Их часто путают, потому что все «реагируют на события».

Регистр схем Подписка Проекция
Задача соглашение о формате данных доставка событий обработчику формирование модели для чтения на сервере
Когда срабатывает при записи (если указана схема) выборка/отправка после фиксации после фиксации, фоновая свёртка
Состояние схемы JSON позиция подписки позиция и результат проекции
Кто потребляет сервер (валидация) ваш обработчик / сервис среда выполнения Chronacta (builtin/CEL/Starlark)
Типичный клиент все, кто пишет события интеграции, процессы, побочные эффекты дашборды, агрегаты, счётчики

Что это: каталог JSON Schema для типов событий. Хранится на сервере (data/schemas/).

Зачем:

  1. Соглашение о формате между тем, кто пишет события, и тем, кто их читает: все знают структуру OrderCreated.
  2. Валидация при записи: запись с -schema-name / -schema-version проверяется до фиксации.
  3. Эволюция схем: режимы backward / forward / full задают, как новая версия схемы совместима со старой (обратная / прямая / полная совместимость).

Чего регистр не делает:

  • не меняет уже записанные события;
  • не заменяет доменную валидацию (можно запретить бизнес-операцию, даже если JSON валиден);
  • не является «базой данных схем» для произвольных API — только данные событий.

Пример по шагам:

Окно терминала
# 1. Регистрация контракта
./bin/chronacta schema register -name order-created -version 1 \
-file order-created.schema.json -compatibility backward
# 2. Запись с привязкой к схеме
./bin/chronacta append -stream order-42 -type OrderCreated \
-data '{"order_id":"42","currency":"RUB"}' \
-schema-name order-created -schema-version 1 \
-expected-version -2
# 3. Проверка данных до отправки (в CI или сервисе)
./bin/chronacta schema validate -name order-created -version 1 \
-data '{"order_id":"42","currency":"RUB"}'

Эволюция: версия 2 схемы может добавить необязательное поле. Старые события не меняются; новые записи могут использовать v2, если режим совместимости позволяет.

Что это: механизм получения событий после фиксации — для вашего кода, другого сервиса или фонового обработчика.

Окно терминала
./bin/chronacta subscribe -stream orders -from 1
  • живёт пока открыто соединение / сессия CLI;
  • догоняющее чтение + события в реальном времени;
  • нет сохранённой позиции на сервере для этого потребителя;
  • подходит: отладка, просмотр «хвоста» журнала, локальный скрипт.
Окно терминала
./bin/chronacta subscription create -id orders-worker \
-stream orders -consumer worker-pool -start 1
./bin/chronacta subscription pull -id orders-worker
./bin/chronacta subscription ack -id orders-worker -position 42
  • позиция чтения хранится на сервере (data/subscriptions/);
  • pull / ack / nack — явная выборка, подтверждение и отказ в обработке;
  • доставка как минимум один раз — повторы нормальны;
  • несколько обработчиков могут читать одну подписку (конкурирующие потребители);
  • после лимита повторов — событие в dead letter (очередь необработанных).

Зачем постоянная подписка:

Сценарий Пример
Побочный эффект отправить email после OrderPaid
Интеграция синхронизировать склад при StockReserved
Saga / длинный процесс следующий шаг процесса
Исходящие (outbox) сервис пишет в Chronacta, подписка забирает «свои» события

Подписка vs просто чтение потока: при чтении потока вы сами опрашиваете журнал и храните смещение. Постоянная подписка — сервер помнит позицию, аренду, повторы и необработанные события.

Подписка vs проекция: подписка отдаёт событие наружу вашему коду. Проекция считает внутри Chronacta и хранит результат.

Что это: фоновый обработчик общего журнала $all (или отдельного потока), который сворачивает события в состояние или результат по программе.

Зачем (для быстрого чтения):

Задача Программа проекции
Счётчик заказов builtin count
Сумма по полю builtin sum_field
Группировка builtin group_count
Кастомная логика CEL или Starlark reduce(event, state)

Пример:

Окно терминала
./bin/chronacta projection create -name orders-count -runtime builtin -program count \
-source-stream order-42 \
-filter 'event.event_type == "OrderCreated" || event.event_type == "OrderPaid"'
./bin/chronacta projection rebuild -name orders-count
./bin/chronacta projection result -name orders-count

Жизненный цикл:

создание → (пересборка, если история уже есть) → работает → позиция сдвигается при успехе
↓ ошибка
повтор → dead letter (журнал событий не меняется)

Когда проекция на Chronacta, когда свой сервис:

Проекция на Chronacta Свой обработчик + база данных
простые агрегаты, счётчики сложные JOIN, полнотекстовый поиск
операционные дашборды тяжёлая аналитика
быстрый старт без второй БД уже есть PostgreSQL для чтения

Проекции не заменяют доменный агрегат в вашем коде: агрегат для команд по-прежнему загружается и сворачивается из потока в вашем коде (или через отдельную проекцию только для чтения).

┌─────────────────────┐
Ваш сервис ──►│ запись OrderCreated │
│ + схема v1 │
└──────────┬──────────┘
│
┌─────────────────────┼─────────────────────┐
▼ ▼ ▼
Регистр схем Проекция Постоянная подписка
(проверка JSON) orders-by-status email-worker
│ │ │
│ ▼ ▼
│ результат: {Paid: 12} SendGrid API
└─ отклонение невалидных данных
  1. Писать текущее состояние вместо событий — {"status":"Paid"} каждый раз вместо OrderPaid.
  2. Один гигантский поток на всё — теряется изоляция агрегатов и защита от гонок.
  3. Ждать доставку ровно один раз от подписки — нужна идемпотентность на стороне потребителя.
  4. Дублировать проекцию и подписку с одной логикой — выберите одно место для свёртки.
  5. Регистр схем без дисциплины — схемы есть, но при записи не указывают -schema-name.

7. Сквозной пример: заказ от команды до проекции и подписки

Заголовок раздела «7. Сквозной пример: заказ от команды до проекции и подписки»

Ниже — полный сценарий домена Sales / Order на одном потоке order-42. Предполагается запущенный сервер (CHRONACTA_DATA_DIR=./data ./bin/chronacta-server).

Элемент Значение
Ограниченный контекст Sales (оформление заказа)
Агрегат Order #42
ID потока order-42
Интеграция Shipping (отдельный обработчик через подписку)
Модель для чтения проекция «сколько событий в заказе» + интерфейс читает поток

Создайте файлы схем и зарегистрируйте их до первой записи:

order-created.schema.json:

{
"$schema": "https://json-schema.org/draft/2020-12/schema",
"type": "object",
"required": ["order_id", "customer_id", "currency"],
"properties": {
"order_id": { "type": "string" },
"customer_id": { "type": "string" },
"currency": { "type": "string", "minLength": 3, "maxLength": 3 }
},
"additionalProperties": false
}

order-paid.schema.json:

{
"$schema": "https://json-schema.org/draft/2020-12/schema",
"type": "object",
"required": ["payment_ref", "amount_cents"],
"properties": {
"payment_ref": { "type": "string" },
"amount_cents": { "type": "integer", "minimum": 0 }
},
"additionalProperties": false
}
Окно терминала
./bin/chronacta schema register -name order-created -version 1 \
-file order-created.schema.json -compatibility backward
./bin/chronacta schema register -name order-paid -version 1 \
-file order-paid.schema.json -compatibility backward

Шаг 1 — создать заказ (команда PlaceOrder, в потоке ещё нет событий):

Окно терминала
./bin/chronacta append -stream order-42 -type OrderCreated \
-data '{"order_id":"42","customer_id":"c1","currency":"RUB"}' \
-schema-name order-created -schema-version 1 \
-expected-version -2 -idempotency-key place-order-42

Шаг 2 — добавить строку (команда AddLine; в потоке уже одно событие, поэтому -expected-version 1):

Окно терминала
./bin/chronacta append -stream order-42 -type OrderLineAdded \
-data '{"sku":"BOOK-1","qty":1,"unit_price_cents":90000}' \
-expected-version 1 -idempotency-key add-line-42-1

Шаг 3 — оплатить (команда PayOrder; в потоке два события, поэтому -expected-version 2):

Окно терминала
./bin/chronacta append -stream order-42 -type OrderPaid \
-data '{"payment_ref":"pay-771","amount_cents":90000}' \
-schema-name order-paid -schema-version 1 \
-expected-version 2 -idempotency-key pay-42

Проверка истории агрегата:

Окно терминала
./bin/chronacta read -stream order-42 -from 1
./bin/chronacta stream-info -stream order-42

Ожидаемый результат: три события, current_version=3.

Chronacta не хранит «текущий объект Order» — ваш сервис читает все события потока и собирает состояние:

func LoadOrder(events []*event.Event) Order {
var o Order
for _, e := range events {
switch e.EventType {
case "OrderCreated":
o = Order{ID: jsonString(e.Data, "order_id"), Status: "New"}
case "OrderLineAdded":
o.Lines = append(o.Lines, parseLine(e.Data))
case "OrderPaid":
o.Status = "Paid"
case "OrderCancelled":
o.Status = "Cancelled"
}
}
return o
}

Команда PayOrder сначала читает поток и собирает заказ (LoadOrder), проверяет, что статус «New» и сумма строк верна, затем записывает OrderPaid, указывая текущее число событий в потоке (в коде — len(events)).

Создаём проекцию — счётчик событий по потоку заказа (операционная модель для чтения):

Окно терминала
./bin/chronacta projection create -name order-42-event-count \
-runtime builtin -program count -source-stream order-42
./bin/chronacta projection rebuild -name order-42-event-count
./bin/chronacta projection result -name order-42-event-count

Ожидаемый результат после трёх записей: count = 3.

Для отчёта по всем потокам «все оплаченные заказы» используйте проекцию на $all с CEL-фильтром по event.event_type == "OrderPaid" — см. projections.md.

Сервис доставки реагирует на оплату и резервирует товар вне Chronacta:

Окно терминала
./bin/chronacta subscription create -id shipping-order-42 \
-stream order-42 -consumer shipping-worker -start 1

Обработчик (псевдокод):

цикл:
пакет = PullEvents("shipping-order-42")
для каждого события в пакете:
если event.type == "OrderPaid" и ещё не обработано(event.id):
reserveWarehouse(event.data)
Ack(position=event.global_position)
иначе:
Ack(...) # пропуск нерелевантных после проверки идемпотентности

Проверка вручную:

Окно терминала
./bin/chronacta subscription pull -id shipping-order-42
# обработали OrderCreated, OrderLineAdded — подтверждение
./bin/chronacta subscription ack -id shipping-order-42 -position <global_position>
# на OrderPaid — побочный эффект + подтверждение
HTTP POST /orders/42/pay
│
▼
OrderService читает поток order-42 в Chronacta
│
▼
Запись OrderPaid (ожидаем 2 события в потоке, схема, ключ идемпотентности)
│
├─► Регистр схем проверяет JSON
├─► Проекция order-42-event-count: result = 3
└─► Подписка shipping-order-42 вызывает API склада
Шаг Команда / симптом
Конфликт версии указали неверное число событий в потоке — перечитайте поток и повторите
Дубль после таймаута повтор с тем же idempotency-key
Проекция отстаёт projection status, detect-stalled
Подписка застряла subscription dead-letters, неподтверждённые позиции
Невалидный JSON запись отклонена до фиксации

Правило Пример
Прошедшее время, факт OrderCreated, PaymentCaptured
Доменный язык OrderLineAdded, не AddLineOk
Стабильность не переименовывать тип без стратегии миграции
Версия схемы отдельно OrderCreated + schema v2, не OrderCreatedV2 как тип

Мелкие (OrderLineAdded на каждую строку):

  • проще частичная повторная обработка и аудит;
  • больше записей в потоке.

Крупные (OrderPlaced со всеми строками):

  • меньше записей;
  • сложнее менять одну строку без полного переписывания смысла.

Практика: один факт — одно событие; пакет строк в одном событии допустим, если так принято в домене.

data metadata
Содержимое доменные поля correlation_id, causation_id, user_id, trace_id
Контракт регистр схем обычно без схемы, соглашение команды
Повторная обработка участвует в свёртке часто игнорируется при восстановлении состояния

Пример metadata при записи (через SDK/API):

{"correlation_id":"req-abc","causation_id":"cmd-pay-42","actor":"user:u1"}
Доменное событие Интеграционное событие
Поток поток агрегата order-42 отдельный поток integration-order-paid или потребитель $all
Аудитория тот же ограниченный контекст другой контекст / внешняя система
Форма полные доменные данные часто упрощённые данные для внешнего API

Не обязательно дублировать: многие команды слушают $all или подписку на доменный поток с фильтром по event_type.

  1. Зарегистрировать схему v2 с обратной совместимостью (backward).
  2. Сервисы, которые записывают события, постепенно переходят на v2.
  3. Потребители и модели для чтения учитывают обе версии при свёртке или фильтрации.
  4. Старые события в потоке не переписываются.

9. Паттерны интеграции и длинных процессов

Заголовок раздела «9. Паттерны интеграции и длинных процессов»
1. Ваш сервис в одной транзакции БД: бизнес-строка + «ожидает публикации»
2. Запись в Chronacta с ключом идемпотентности = id outbox
3. Постоянная подписка или проекция подтверждает доставку во внешние системы
4. Пометить запись outbox как обработанную (в вашей БД)

Chronacta не участвует в двухфазной фиксации с PostgreSQL — таблица outbox остаётся в вашей БД; Chronacta — надёжный журнал после успешной записи.

Длинный процесс «заказ → оплата → отгрузка → возврат»:

  • оркестратор хранит состояние saga в своём потоке saga-checkout-42 или в БД;
  • каждый шаг — запись доменного события + реакция подписки;
  • компенсации — события ShipmentCancelled, PaymentRefunded, а не удаление записей.

Подписки Chronacta дают доставку как минимум один раз; обработчик шага saga идемпотентен по (saga_id, step, event_id).

9.3. Хореография (без центрального оркестратора)

Заголовок раздела «9.3. Хореография (без центрального оркестратора)»

Несколько сервисов без центрального оркестратора:

OrderPaid → подписка A (склад)
→ подписка B (email)
→ проекция C (дашборд)

Каждый потребитель независим, порядок между сервисами не гарантирован — закладывайте идемпотентность и согласованность со временем.

Для публикации зафиксированных событий во внешнюю шину:

Окно терминала
CHRONACTA_NATS_URL=nats://127.0.0.1:4222 ./bin/chronacta-server

Chronacta остаётся главным источником данных; NATS — транспорт с доставкой «как минимум один раз». См. nats-bridge.md.

Задача Механизм
Побочный эффект в другом сервисе постоянная подписка
Дашборд / агрегат на сервере проекция
Контракт данных регистр схем
Рассылка наружу NATS-мост + идемпотентные потребители
Длинный многошаговый процесс поток saga + подписки

10. Тестирование, повторная обработка и локальная разработка

Заголовок раздела «10. Тестирование, повторная обработка и локальная разработка»

Unit-тест функции свёртки на наборе событий:

events := []*event.Event{
{EventType: "OrderCreated", Data: []byte(`{"order_id":"42"}`)},
{EventType: "OrderPaid", Data: []byte(`{"payment_ref":"p1"}`)},
}
o := LoadOrder(events)
if o.Status != "Paid" { t.Fatal(...) }

Паттерн из репозитория: поднять сервер или использовать examples/sdk-test — запись, чтение, проверка версии.

Окно терминала
make build
CHRONACTA_DATA_DIR=./data-test ./bin/chronacta-server &
./bin/chronacta-example-sdk

После изменения логики модели для чтения:

Окно терминала
./bin/chronacta projection rebuild -name order-42-event-count
./bin/chronacta projection result -name order-42-event-count

Исходный поток не меняется — повторная обработка безопасна для аудита.

10.4. Повторная обработка подписки (осторожно)

Заголовок раздела «10.4. Повторная обработка подписки (осторожно)»

Постоянная подписка не «перематывается» сама при смене обработчика. Для повторной обработки:

  • создать новую подписку с -start 1 и новым -id;
  • или операторский replay / сброс позиции (если поддерживается вашей версией CLI);
  • не пересобирайте production-подписки без плана идемпотентности.
  • побочные эффекты должны быть идемпотентны.
Окно терминала
./bin/chronacta export stream -stream order-42 -output fixtures/order-42.jsonl -format jsonl
./bin/chronacta import -file fixtures/order-42.jsonl -format auto \
-duplicate-policy skip_by_event_id

Полезно для тестовых наборов в CI и переноса между средами разработки и staging.

  • Соглашение об именах потоков задокументировано в README сервиса.
  • Каждый event_type имеет схему или явное «без схемы».
  • При записи указывается ожидаемое число событий в потоке и ключ идемпотентности на границах API.
  • Обработчик подписки покрыт тестом на повторную доставку.
  • Пересборка проекции проверена после смены программы.

Client / CLI / SDK
|
gRPC :2113 <── REST / WebSocket / Admin BFF
|
auth, RBAC, limits, tracing
|
фиксация -> WAL -> сегменты и индексы
|
subscriptions / projections / backup / replication

Основной контракт API — api/proto/stream/v1/stream.proto. На диске обычно каталоги segments/, wal/, schemas/, subscriptions/, projections/, idempotency/, auth/ и в HA — cluster/.

В HA лидер — единственный, кто принимает записи. Реплики применяют кадры WAL и могут отдавать чтение с задержкой. Запись подтверждается после фиксации на лидере и согласования с кворумом.


Окно терминала
git clone https://gitverse.ru/AndreyI/chronacta.git
cd chronacta
git checkout v1.0.0
make proto-tools # один раз при отсутствии protoc plugins
make proto
make build
./bin/chronacta-server --version

Проверки перед production:

Окно терминала
make proto test test-race vet build
make openapi-validate ops-validate sdk-smoke release-validate
Окно терминала
CHRONACTA_DATA_DIR=./data ./bin/chronacta-server

В другом терминале — сквозной пример «создать заказ»:

Окно терминала
./bin/chronacta health
# Создание агрегата (новый поток)
./bin/chronacta append -stream order-42 -type OrderCreated \
-data '{"order_id":"42","customer_id":"c1"}' -expected-version -2
# Развитие агрегата
./bin/chronacta append -stream order-42 -type OrderLineAdded \
-data '{"sku":"A1","qty":2}' -expected-version 1
./bin/chronacta read -stream order-42 -from 1
./bin/chronacta read-all -from 1
./bin/chronacta verify

Сервер конфигурируется переменными CHRONACTA_*. config.yaml — справочный список настроек, а не основной автоматически загружаемый файл.


Окно терминала
sudo useradd -r -s /bin/false chronacta
sudo mkdir -p /var/lib/chronacta/data /var/lib/chronacta/backups /etc/chronacta/tls
sudo chown -R chronacta:chronacta /var/lib/chronacta

Минимальный /etc/chronacta/env:

Окно терминала
CHRONACTA_PROFILE=production
CHRONACTA_DATA_DIR=/var/lib/chronacta/data
CHRONACTA_BACKUP_DIR=/var/lib/chronacta/backups
CHRONACTA_GRPC_PORT=2113
CHRONACTA_METRICS_ENABLED=true
CHRONACTA_METRICS_ADDRESS=127.0.0.1:9090
CHRONACTA_AUTH_ENABLED=true
CHRONACTA_AUTH_DEVELOPMENT_MODE=false
CHRONACTA_AUTH_STORE=/var/lib/chronacta/data/auth/users.json
CHRONACTA_TLS_ENABLED=true
CHRONACTA_TLS_CERT_FILE=/etc/chronacta/tls/server.crt
CHRONACTA_TLS_KEY_FILE=/etc/chronacta/tls/server.key
CHRONACTA_CLUSTER_REPLICATION_TOKEN=<длинный случайный токен>

Профиль production требует auth, TLS, выключенный режим разработки и replication token. При вебхуке нужен секрет. Если WebSocket слушает не только localhost, задайте CHRONACTA_WS_ALLOWED_ORIGINS; передачу токена в query-параметре оставляйте выключенной.

Пример запуска через systemd:

[Service]
User=chronacta
EnvironmentFile=/etc/chronacta/env
WorkingDirectory=/opt/chronacta
ExecStart=/opt/chronacta/bin/chronacta-server
Restart=on-failure
RestartSec=5
LimitNOFILE=65535
Окно терминала
sudo systemctl daemon-reload
sudo systemctl enable --now chronacta
sudo systemctl status chronacta

Окно терминала
./bin/chronacta auth bootstrap -username admin -password '<initial-secret>'
./bin/chronacta auth login -username admin -password '<secret>'

Пользователи и роли управляются командами auth create-user, auth create-role, auth grant, auth revoke, auth list-users, auth list-roles, auth audit.

Основные права: stream.read, stream.append, schema.manage, projection.read, projection.manage, backup, verify, tenant.read, tenant.admin, lifecycle.execute.

Окно терминала
./bin/chronacta health \
-server chronacta.example.com:2113 \
-tls -tls-ca /etc/chronacta/tls/ca.crt \
-token "$CHRONACTA_TOKEN"

OIDC использует фактическое имя CHRONACTA_OIDC_ISSUER:

Окно терминала
CHRONACTA_OIDC_ENABLED=true
CHRONACTA_OIDC_ISSUER=https://idp.example.com/
CHRONACTA_OIDC_CLIENT_ID=chronacta
CHRONACTA_OIDC_DEFAULT_ROLES=reader
CHRONACTA_OIDC_USERNAME_CLAIM=email
CHRONACTA_OIDC_TENANT_CLAIM=tenant_id

Окно терминала
./bin/chronacta stream list
./bin/chronacta stream-info -stream order-42
./bin/chronacta read -stream order-42 -from 1 -count 100
./bin/chronacta read-all -from 1 -count 100
./bin/chronacta storage-status
./bin/chronacta verify
./bin/chronacta verify-db

Надёжная запись с retry:

Окно терминала
./bin/chronacta append \
-stream order-42 -type OrderCreated \
-data '{"order_id":"42"}' \
-expected-version -2 -idempotency-key checkout-42

Подробнее о назначении — разделы 6 и 7. Здесь команды.

Окно терминала
./bin/chronacta schema register \
-name order-created -version 1 \
-file order-created.schema.json -compatibility backward
./bin/chronacta schema get -name order-created -version 1
./bin/chronacta schema list
./bin/chronacta schema validate -name order-created -version 1 \
-data '{"order_id":"42","currency":"RUB"}'

Запись с привязкой к схеме:

Окно терминала
./bin/chronacta append -stream order-42 -type OrderCreated \
-data '{"order_id":"42","currency":"RUB"}' \
-schema-name order-created -schema-version 1 \
-expected-version -2

Режимы совместимости: none, backward, forward, full. Новая версия схемы не изменяет уже записанные события.


Подписка на сессию — отладка и хвост журнала

Заголовок раздела «Подписка на сессию — отладка и хвост журнала»
Окно терминала
./bin/chronacta subscribe -stream order-42 -from 1
Окно терминала
./bin/chronacta subscription create -id order-notifications \
-stream order-42 -consumer email-sender -start 1
./bin/chronacta subscription pull -id order-notifications
./bin/chronacta subscription ack -id order-notifications -position 42
./bin/chronacta subscription nack -id order-notifications -position 42 -reason retry
./bin/chronacta subscription pause -id order-notifications
./bin/chronacta subscription resume -id order-notifications
./bin/chronacta subscription dead-letters -id order-notifications

Правила эксплуатации:

  • ack только после успешного побочного эффекта (email отправлен, запись на складе создана);
  • обработчик идемпотентен по event_id;
  • необработанные события (dead letters) — ручной разбор, не игнорировать.

Fan-out consumer groups (масштабирование обработки):

Окно терминала
./bin/chronacta consumer-group create -id workers -stream order-42
./bin/chronacta consumer-group join -id workers -member node-a
./bin/chronacta consumer-group join -id workers -member node-b

Каждый member получает события, где global_position % N == member_index.

Transient catch-up: медленный подписчик не отключается — сервер переводит поток в режим catch-up (см. SubscribeResponse.control). Legacy: CHRONACTA_SUBSCRIPTION_CATCHUP_MODE=false.

Connector: sidecar chronacta-connector для managed durable jobs — см. connector.md.

См. subscriptions.md и consumer-groups.md.

chronacta-ingress — отдельный коммерческий компонент. Он не входит в бесплатный chronacta-server и запускается отдельно. Его задача — принимать только выбранные события из NATS JetStream, записывать их в Chronacta и подтверждать сообщение в NATS только после успешной записи.

Не смешивайте события для хранения и operational-сообщения:

  • agents.events.persist — события, которые могут попасть в Chronacta;
  • agents.events.telemetry — ping, heartbeat и состояние агента; их обрабатывает отдельный NATS-only consumer;
  • chronacta.events — downstream subject для разработчиков после persistence;
  • agents.events.persist.dlq — invalid/rejected/poison-сообщения.

Ingress должен быть единственным production consumer’ом agents.events.persist. Разработчикам запрещается subscribe на этот raw subject — они получают только chronacta.events.>.

agent → agents.events.persist → chronacta-ingress → Chronacta → chronacta.events → developer
agent → agents.events.telemetry ───────────────────────────────→ health consumer
{
"version": 1,
"message_id": "01J...",
"agent_id": "agent-42",
"event_type": "ProcessStarted",
"stream_id": "agent-42",
"data": {"pid": 123},
"metadata": {"hostname": "pc-42"}
}
  1. Создайте в JetStream stream, включающий persist subject, и durable pull consumer chronacta-ingress с filter agents.events.persist.
  2. Настройте ACL: ingress может consume/ack raw persist subject и publish в downstream/DLQ; developer не может subscribe на raw subject.
  3. Создайте entitlement-файл:
{"entitlements":{"nats_ingress_enabled":true}}
  1. Заполните examples/ingress-test/ingress.yaml: URL NATS, stream, durable consumer, allowlist allow_event_types, Chronacta address и dlq_subject.
  2. Запустите шлюз отдельно:
Окно терминала
make ingress
CHRONACTA_INGRESS_CONFIG=./examples/ingress-test/ingress.yaml \
CHRONACTA_INGRESS_METRICS_ADDR=:9091 \
./bin/chronacta-ingress
  1. Только после проверки /health подключайте developer consumers к chronacta.events.>.

Allowlist работает как default-deny: в Chronacta попадут только явно разрешённые event_type. Не отправляйте ping в persist subject. При недоступности Chronacta ingress не делает ACK, поэтому JetStream повторит доставку. После append шлюз использует idempotency key nats:<message_id>, затем публикует downstream и подтверждает исходное сообщение. Доставка — at-least-once, поэтому downstream consumer должен быть идемпотентным.

Полный runbook: NATS ingress. Граница лицензии и rollback также описаны там.


Окно терминала
./bin/chronacta projection create -name orders-count -runtime builtin -program count \
-source-stream order-42
./bin/chronacta projection rebuild -name orders-count
./bin/chronacta projection list
./bin/chronacta projection status -name orders-count
./bin/chronacta projection result -name orders-count
./bin/chronacta projection errors -name orders-count
./bin/chronacta projection detect-stalled

Runtimes: builtin, CEL, Starlark. JavaScript не поддерживается.

После создания, если в потоке уже есть история — обязательно rebuild. Позиция сдвигается только после успешной свёртки.

JSON-проекции — для ops. Для SQL/UI:

  1. Опишите таблицы в YAML (examples/postgres-projector/schema.yaml) — схему задаёт команда read model.
  2. Сгенерируйте SQL: ./bin/chronacta sqlgen generate ddl|dml ...
  3. Примените DDL в Postgres вручную.
  4. Запустите chronacta-connector с handler: go:PostgresReadModel.

См. postgres-projections.md.

См. projections.md.


Окно терминала
./bin/chronacta backup create -archive ./backups/node.tar.gz
./bin/chronacta backup inspect -archive ./backups/node.tar.gz
./bin/chronacta backup list -dir ./backups

Restore выполняется после остановки сервера в пустой каталог:

Окно терминала
./bin/chronacta-admin restore \
--archive ./backups/node.tar.gz \
--target ./restored-data --confirm

Перед изменениями используйте --dry-run. Архив содержит описание состава и контрольные суммы SHA-256.

Окно терминала
./bin/chronacta export stream -stream order-42 -output order-42.jsonl -format jsonl
./bin/chronacta import -file order-42.jsonl -format auto \
-duplicate-policy reject_existing_stream

При импорте назначаются новые версии в потоках и глобальные позиции. Политики: reject_existing_stream (по умолчанию), append, skip_by_event_id.


Окно терминала
./bin/chronacta cluster status -json
./bin/chronacta cluster add-node -node-id node-2 -addr host2:2113 -role follower
./bin/chronacta cluster remove-node -node-id node-2
./bin/chronacta cluster promote -node-id node-2
./bin/chronacta cluster snapshot push -target-node-id node-2

Перед постепенным обновлением сделайте backup, проверьте кворум и отставание реплик. Сначала обновляйте followers, затем бывшего leader.

REST и WebSocket включаются переменными CHRONACTA_GATEWAY_REST_* и CHRONACTA_GATEWAY_WS_*. REST покрывает часть API; OpenAPI — на /v1/openapi.json. WebSocket: /v1/ws/streams/{id} и /v1/ws/all; передавайте Authorization: Bearer ....

Admin UI — серверный интерфейс на /admin. В production привязывайте его к localhost или размещайте за обратным прокси с TLS.


Production-клиенты используют pkg/client/resilience для NATS-style reconnect:

rc, err := resilience.Connect(ctx, resilience.Config{
Config: client.Config{Address: "127.0.0.1:2113", Token: token},
})
if err != nil { log.Fatal(err) }
defer rc.Close()
// Append с idempotency key для безопасных retry
result, err := rc.AppendToStream(ctx, "order-42", 1,
[]*event.Event{{EventType: "OrderPaid", Data: []byte(`{"payment_ref":"pay-9"}`)}},
client.AppendOptions{IdempotencyKey: "pay-9"})
// Live tail с auto-resubscribe после обрыва
live, err := rc.SubscribeLive(ctx, "order-42", 0, func(_ context.Context, ev *event.Event) error {
return process(ev)
})

Durable worker: rc.RunDurableWorker. Примеры: examples/resilience-test, examples/subscription-test.

Низкоуровневый reconnect без resilience: examples/sdk-test (только для демонстрации).

cli, err := client.Dial(ctx, client.Config{Address: "127.0.0.1:2113", Token: token})

Для TLS передайте client.TLSConfig.

Окно терминала
npm install @chronacta/client
import { RestClient } from "@chronacta/client";
const client = new RestClient({
baseUrl: "https://chronacta.example.com:8081",
token: process.env.CHRONACTA_TOKEN,
});
await client.append("order-42", "OrderCreated", { order_id: "42" });
const events = await client.readStream("order-42");
Окно терминала
pip install chronacta
from chronacta import RestClient
import os
client = RestClient("https://chronacta.example.com:8081",
token=os.environ["CHRONACTA_TOKEN"])
client.append("order-42", "OrderCreated", {"order_id": "42"})
print(client.read_stream("order-42"))

TypeScript и Python — облегчённые REST-клиенты. Для постоянных подписок, проекций и полного Admin API используйте Go/gRPC или документированные REST endpoints.


Окно терминала
curl -s http://127.0.0.1:9090/healthz
curl -s http://127.0.0.1:9090/readyz
curl -s http://127.0.0.1:9090/metrics

Prometheus собирает задержки, ошибки, сбои авторизации, WAL, подписки/проекции и отставание репликации. JSON-логи содержат ID запроса, метод, код gRPC и длительность; секреты маскируются.

Симптом Что проверить
запись отклонена expected-version, token, quotas, disk
таймаут после записи повтор с тем же ключом идемпотентности
проекция отстаёт projection detect-stalled, ошибки, необработанные события
повторная доставка подписки норма «как минимум один раз»; идемпотентность потребителя
реплика отстаёт лидер, сеть, cluster status, отставание репликации

  • Зафиксированы тег релиза и контрольные суммы.
  • Data и backups находятся на отдельных защищенных путях.
  • Включены auth, TLS и выключен режим разработки.
  • Replication token не хранится в репозитории.
  • Admin UI, metrics, REST и WebSocket не выставлены напрямую в интернет.
  • Настроены Prometheus, логи, оповещения и мониторинг диска.
  • Выполнены резервное копирование, inspect, пробное восстановление и verify.
  • Для HA проверены кворум, отставание, переключение при сбое и постепенное обновление.
  • Модель домена документирована: id потока, типы событий, версии схем.
  • Клиенты передают ожидаемую версию потока и ключ идемпотентности.
  • Постоянные подписки и проекции готовы к повторной доставке и пересборке.