Files
dzentra_bot/docs/migrations/build_060_18_architecture.md

64 KiB
Raw Permalink Blame History

Build 060.18 — Trade Stream Consistency Controller

Статус: Architecture Specification
Build: 060.18
Ветка: Trades Feed (Time & Sales)
Документ: build_060_18_architecture.md
Связанные документы: build_060_18.md — план реализации данного Build.


Назначение документа

Настоящий документ является официальной архитектурной спецификацией Build 060.18 и определяет построение подсистемы обеспечения согласованности потока сделок (Trade Stream Consistency).

Документ фиксирует все архитектурные решения, принятые до начала реализации, и служит единственным источником истины (Single Source of Truth) при разработке данного Build.

Все решения, описанные ниже, считаются утвержденными до начала реализации и не должны изменяться в процессе написания кода без подготовки нового ADR.


Статус Build

Build 060.18 является продолжением серии Build 060, посвящённой построению новой подсистемы получения рыночных данных (Market Data Acquisition).

К моменту начала данного Build в проекте уже существуют:

  • каноническая модель Trade;
  • транспортные Parser;
  • Mapper;
  • Value Validation;
  • REST и WebSocket интеграции;
  • Trades Feed;
  • Trades Handler.

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

Однако на текущем этапе отсутствует механизм, обеспечивающий согласованность самого потока сделок.

Build 060.18 закрывает именно эту архитектурную задачу.


Контекст

После завершения Build 060.17 система умеет получать сделки одновременно из различных транспортных источников.

Например:

  • REST Backfill;
  • WebSocket Trade Stream;
  • будущие источники исторических данных.

Все эти источники после прохождения Parser, Validation и Mapper возвращают одинаковый объект:

Trade

Однако наличие канонической модели ещё не означает существование канонического потока.

Например, поток может содержать:

100
101
100

или

100
101
99

или

500
501
500 (с другой ценой)

Каждый из этих случаев нарушает различные инварианты системы.

Следовательно, между построением объекта Trade и публикацией сделки потребителю должен существовать отдельный уровень контроля согласованности.

Именно этот уровень реализуется Build 060.18.


Предпосылки

Настоящий Build опирается на архитектурные решения, принятые ранее.

Build 057

Определены базовые принципы Market Data Acquisition.

Разделены транспортный и доменный уровни.


Build 060.1

Построена каноническая модель Trade.

Все источники данных приводятся к единому объекту.


Build 060.2

Завершена унификация транспортных моделей.


Build 060.4

Завершено разделение Parser и Mapper.

Parser отвечает только за транспорт.

Mapper отвечает только за построение Domain Model.


Build 060.9

Построена система Value Validation.

Корректность отдельных значений гарантируется до появления объекта Trade.


Build 060.11

Определены транспортные адаптеры REST и WebSocket.


Build 060.17

Построен Trades Feed.

Получение сделок полностью функционирует.

На этом этапе поток состоит из независимых объектов Trade.

Build 060.17 сознательно не занимается проверкой согласованности последовательности сделок.


Проблема

Build 060.17 гарантирует корректность каждой отдельной сделки.

Однако система ещё не гарантирует корректность последовательности сделок.

Например:

Trade
Trade
Trade

не обязательно образуют корректный поток.

В частности отсутствуют гарантии:

  • монотонности;
  • отсутствия повторов;
  • отсутствия конфликтующих дублей;
  • целостности опубликованного потока.

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

Build 060.18 устраняет именно эту проблему.


Основная идея Build

Главная идея Build заключается в разделении двух понятий.

До настоящего момента существовала только каноническая модель Trade.

После завершения Build появляются две независимые сущности.

Первая:

Canonical Trade

Вторая:

Canonical Trade Stream

Это принципиально разные понятия.

Trade представляет собой корректный объект предметной области.

Trade Stream представляет собой последовательность объектов Trade, удовлетворяющую дополнительным инвариантам.

Следовательно, появление объекта Trade ещё не означает его публикацию в поток.

Перед публикацией Trade обязан пройти проверку согласованности.


Цель Build

Build обязан обеспечить существование единственного канонического потока сделок.

После завершения Build система должна гарантировать следующие свойства.

Каждая опубликованная сделка:

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

Все источники данных должны использовать одинаковый механизм проверки согласованности.

Источник происхождения сделки не должен влиять на работу системы.


Что НЕ входит в Scope Build

Настоящий Build сознательно НЕ реализует:

  • восстановление пропущенных сделок;
  • обнаружение разрывов последовательности;
  • повторную синхронизацию через REST;
  • сохранение состояния между перезапусками процесса;
  • журналирование истории потока;
  • долговременное хранение дедупликационного окна;
  • диагностику качества соединения;
  • обработку потери WebSocket.

Все перечисленные задачи относятся к следующим Build серии 060.

Build 060.18 отвечает исключительно за согласованность уже поступающих сделок.


Архитектурные принципы

При реализации Build используются следующие фундаментальные принципы.

1. Canonical First

Любая транспортная информация должна быть преобразована в каноническую модель до начала проверки согласованности.

Контроллер никогда не работает с транспортными структурами.


2. Single Domain Model

Во всей системе существует единственная модель Trade.

REST и WebSocket не имеют собственных моделей после этапа Mapper.


3. Stream Before Publication

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


4. Immutable Domain Objects

Объект Trade никогда не изменяется после создания.

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


5. One Source Of Truth

Единственным владельцем состояния потока является компонент Stream Consistency.

Никакие другие части Acquisition Pipeline не должны хранить собственную историю уже обработанных сделок.


6. Stateless Acquisition

Parser, Mapper, Validation и Handler остаются полностью stateless.

Первой stateful-компонентой Acquisition становится только Stream Consistency Controller.


7. Separation Of Responsibilities

Trade описывает предметную область.

TradeStreamState хранит состояние.

TradeStreamConsistencyController принимает доменные решения.

Ни один из этих компонентов не должен брать на себя ответственность другого.



8. Unique File Naming

Во всём проекте Dzentra не допускается существование нескольких файлов с одинаковым именем независимо от их расположения в каталогах.

Единственным исключением является служебный файл:

__init__.py

Имя каждого файла должно однозначно отражать его назначение и позволять определить его содержимое без открытия файла.

Например:

trade_stream_consistency_controller.py

trade_stream_state.py

trade_stream_protocol.py

trade_stream_exceptions.py

Не допускается использование общих имён файлов, таких как:

controller.py

state.py

protocol.py

exceptions.py

Canonical Trade Stream

Введение

Главной архитектурной задачей Build 060.18 является переход от понятия отдельной канонической сделки к понятию канонического потока сделок.

До настоящего Build система гарантировала корректность каждого объекта Trade независимо.

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

Именно поток, а не отдельная сделка, становится новой доменной сущностью Acquisition Layer.


Определение Canonical Trade Stream

Canonical Trade Stream — это последовательность канонических объектов Trade, удовлетворяющая всем инвариантам согласованности потока.

Именно эта последовательность считается единственным допустимым источником сделок для всех последующих компонентов системы.

Все последующие модули Dzentra должны считать, что получаемый ими поток уже является согласованным.

Повторная проверка порядка, дедупликации или целостности не допускается.


Отличие Canonical Trade от Canonical Trade Stream

Canonical Trade отвечает на вопрос:

Является ли данный объект корректной сделкой?

Canonical Trade Stream отвечает на вопрос:

Может ли данная сделка стать частью уже существующего потока?

Это принципиально разные уровни проверки.

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

Например:

Trade #100
Trade #101
Trade #99

Все три объекта являются корректными Trade.

Однако третья сделка нарушает согласованность потока.

Следовательно, она никогда не становится частью Canonical Trade Stream.


Граница формирования потока

До прохождения проверки согласованности существует только набор независимых объектов Trade.

REST
        │
WebSocket
        │
        ▼
Parser
        ▼
Validation
        ▼
Mapper
        ▼
Trade

После прохождения Stream Consistency появляется новая сущность.

Trade
        │
        ▼
TradeStreamConsistencyController
        ▼
Canonical Trade Stream
        ▼
Consumer

Именно TradeStreamConsistencyController создаёт канонический поток.

Никакой другой компонент системы не имеет такого права.


Инварианты Canonical Trade Stream

Любой поток, публикуемый системой, обязан удовлетворять следующим инвариантам.


Инвариант №1

Поток состоит исключительно из объектов Trade.

Транспортные модели никогда не публикуются.


Инвариант №2

Каждая опубликованная сделка успешно прошла Value Validation.


Инвариант №3

Каждая опубликованная сделка успешно прошла проверку Stream Consistency.


Инвариант №4

Для каждого символа порядок trade_id никогда не уменьшается.


Инвариант №5

Каждая биржевая сделка публикуется не более одного раза.


Инвариант №6

Конфликтующие повторы никогда не публикуются.


Инвариант №7

После публикации сделка считается неизменяемой.


Инвариант №8

Каждый символ имеет полностью независимую историю потока.


Инвариант №9

Проверка согласованности выполняется до публикации сделки.


Инвариант №10

После публикации сделка никогда повторно не проверяется.


Identity Model

Назначение

Для обеспечения дедупликации необходимо определить каноническую идентичность сделки.

Без формального определения идентичности невозможно определить:

  • повтор;
  • конфликтующий повтор;
  • новую сделку.

Рассматриваемые варианты

Во время проектирования были рассмотрены несколько вариантов.


Вариант 1

trade_id

Отклонён.

Причина:

Архитектура не должна предполагать глобальную уникальность trade_id.

Разные торговые инструменты могут использовать одинаковые идентификаторы.


Вариант 2

(symbol, trade_id)

Принят.

Данная пара однозначно определяет одну сделку внутри канонического потока.


Вариант 3

(symbol, trade_id, source)

Отклонён.

После Mapper источник происхождения сделки перестаёт иметь доменное значение.

REST и WebSocket обязаны порождать одну и ту же каноническую сущность.


Вариант 4

(symbol, trade_id, timestamp)

Отклонён.

Время исполнения является атрибутом сделки, а не частью её идентичности.


Официальная модель идентичности

Во всей системе Build 060.18 официальной идентичностью сделки считается:

(symbol, trade_id)

Никакие другие поля не участвуют в определении идентичности.


Ordering Model

Назначение

После определения идентичности необходимо определить правило упорядочивания потока.


Рассматриваемые варианты


Ordering по timestamp

Отклонён.

Причины:

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

Ordering по executed_at

Отклонён.

executed_at является характеристикой сделки, но не механизмом восстановления порядка.


Ordering по (timestamp, trade_id)

Отклонён.

Избыточно.

Усложняет систему без появления дополнительных гарантий.


Ordering по trade_id

Принят.

Биржа уже определяет последовательность исполнения сделок.

Следовательно, система должна использовать именно её.


Официальное правило Ordering

Для каждого символа поток обязан удовлетворять условию:

trade_id(new) >= trade_id(last)

При этом допускаются разрывы последовательности.

Например:

100
101
103
120

является корректным потоком.

Build 060.18 не занимается анализом пропущенных идентификаторов.

Эта задача относится к Build 060.19.


Нарушение порядка

Если новая сделка имеет меньший trade_id, чем последняя опубликованная сделка данного символа, поток считается нарушенным.

Пример:

100
101
98

В этом случае сделка отвергается.

Контроллер генерирует TradeOrderingError.

Состояние потока при этом не изменяется.


Deduplication Model

Назначение

После проверки порядка необходимо определить правила обработки повторных сделок.


Определение дубликата

Дубликатом считается сделка, имеющая ту же идентичность:

(symbol, trade_id)

что и уже зарегистрированная сделка.


Полный дубликат

Если все канонические поля совпадают, повтор считается корректным.

Например:

BTC
trade_id = 150
price = 100
quantity = 5

и повтор:

BTC
trade_id = 150
price = 100
quantity = 5

представляют одну и ту же сделку.

Контроллер не публикует её повторно.

Ошибки не возникает.


Конфликтующий дубликат

Если идентичность совпадает, но хотя бы одно бизнес-поле отличается:

  • price;
  • quantity;
  • executed_at;
  • aggressor_side;

поток считается противоречивым.

Например:

trade_id = 500

price = 100

позже:

trade_id = 500

price = 101

Такая ситуация невозможна внутри корректного канонического потока.

Контроллер обязан немедленно завершить обработку ошибкой.

Публикация сделки запрещается.


Что не считается дубликатом

Следующие сделки никогда не считаются повтором:

trade_id = 500

trade_id = 501

Даже если совпадают:

  • цена;
  • объём;
  • время исполнения.

Идентичность определяется исключительно парой:

(symbol, trade_id)

Итоговая модель обработки сделки

Каждая поступающая сделка проходит последовательность проверок.

Trade
    │
    ▼
Определение символа
    │
    ▼
Получение состояния потока
    │
    ▼
Проверка существования trade_id
    │
    ├─────────────── Да ───────────────┐
    │                                  │
    ▼                                  ▼
Проверка совпадения              Полный дубликат
бизнес-полей                           │
    │                                  ▼
    │                           Не публиковать
    │
    ▼
Конфликт?
    │
    ├──── Да ───► TradeConsistencyError
    │
    ▼
Нет
    │
    ▼
Проверка порядка
    │
    ├──── Нарушение ─► TradeOrderingError
    │
    ▼
Регистрация сделки
    │
    ▼
Публикация

Архитектурный результат

После завершения Build 060.18 в системе появляется новый уровень доменной модели.

До Build:

Trade

После Build:

Trade
        │
        ▼
Canonical Trade Stream

Именно канонический поток становится единственным допустимым источником данных для всех последующих компонентов Market Intelligence Pipeline.


Архитектура компонентов

После определения модели Canonical Trade Stream необходимо определить архитектурные компоненты, обеспечивающие его существование.

Build 060.18 вводит в Acquisition Layer первую stateful-подсистему.

До настоящего момента все компоненты Acquisition являлись stateless.

Появление Stream Consistency является первой точкой, где система начинает хранить собственное состояние.

Именно поэтому данный Build вводит строго определённые границы ответственности.


Архитектура подсистемы

Подсистема Stream Consistency состоит из двух компонентов.

TradeStreamConsistencyController
                │
                ▼
        TradeStreamState

Оба компонента являются внутренними компонентами Acquisition Layer.

Никакие другие части системы не имеют права изменять их состояние.


TradeStreamConsistencyController

Назначение

TradeStreamConsistencyController является единственной точкой формирования Canonical Trade Stream.

Никакой другой компонент системы не имеет права принимать решение о публикации сделки.

Контроллер определяет:

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

Ответственность Controller

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

К ним относятся:

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

Controller НЕ отвечает

Контроллер сознательно не отвечает за:

  • хранение данных;
  • структуру дедупликационного окна;
  • алгоритм удаления старых записей;
  • реализацию FIFO;
  • внутреннее устройство состояния.

Все перечисленные задачи принадлежат исключительно TradeStreamState.


Главный принцип Controller

Controller принимает решения.

State хранит данные.

Это фундаментальное архитектурное правило Build 060.18.


TradeStreamState

Назначение

TradeStreamState представляет собой внутреннюю модель состояния одного потока сделок.

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

Например:

BTCUSDT
        │
        ▼
TradeStreamState

или

ETHUSDT
        │
        ▼
TradeStreamState

Состояние различных символов никогда не смешивается.


Почему состояние существует отдельно

Во время проектирования рассматривалась возможность хранения всех данных непосредственно внутри Controller.

Этот вариант был отклонён.

Причина:

Controller должен выражать бизнес-правила.

State должен выражать состояние предметной области.

Разделение этих двух ролей существенно упрощает дальнейшее развитие системы.


Инварианты TradeStreamState

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


Инвариант №1

Состояние обслуживает ровно один symbol.


Инвариант №2

Состояние никогда не принимает Trade другого symbol.

Даже если Controller ошибётся.

Это является дополнительной защитой целостности системы.


Инвариант №3

Состояние хранит только уже опубликованные сделки.

Непринятые сделки никогда не изменяют состояние.


Инвариант №4

Состояние никогда самостоятельно не принимает бизнес-решения.

Оно лишь предоставляет операции хранения.


Инвариант №5

Состояние полностью скрывает собственную реализацию.

Controller не знает:

  • используется ли OrderedDict;
  • используется ли deque;
  • используется ли Ring Buffer;
  • используется ли специализированная структура данных.

Внутреннее состояние

Минимальная модель состояния включает:

symbol

last_trade_id

recent_trades

Этого достаточно для реализации всех инвариантов Build 060.18.


symbol

Каждое состояние хранит собственный символ.

Например:

BTCUSDT

Это необходимо по двум причинам.

Во-первых,

для проверки принадлежности входящей сделки.

Во-вторых,

для полноценной диагностики ошибок.


last_trade_id

last_trade_id означает:

максимальный trade_id, успешно опубликованный системой для данного символа.

Это очень важно.

last_trade_id НЕ означает:

последний увиденный trade_id.

Например:

100

101

95

После возникновения OrderingError состояние остаётся:

last_trade_id = 101

Никаких изменений не происходит.


recent_trades

recent_trades представляет собой окно уже опубликованных сделок.

Окно необходимо исключительно для проверки повторов.

После выхода сделки из окна

она перестаёт участвовать в дедупликации.


Почему используется окно

Полная история сделок потенциально бесконечна.

Следовательно,

невозможно хранить все Trade.

Используется ограниченное окно фиксированного размера.

Размер окна определяется конфигурацией.


Требования к окну

Окно обязано обеспечивать:

поиск

O(1)

вставку

O(1)

удаление самой старой записи

O(1)

Эти требования являются обязательными.


FIFO Window

Во время проектирования рассматривались различные структуры хранения.


Полная история

Отклонена.

Причина:

неограниченный рост памяти.


HashSet

Отклонён.

Причина:

невозможно проверить конфликтующий повтор.


LRU Cache

Отклонён.

Причина:

LRU ориентирован на обращения.

Trade Stream ориентирован на порядок появления.


FIFO Window

Принят.

FIFO полностью соответствует природе потока сделок.

Самые старые сделки постепенно забываются.

Новые добавляются в конец окна.


Внутренний ключ

Поскольку один экземпляр TradeStreamState обслуживает только один symbol,

ключом окна становится исключительно:

trade_id

Полная идентичность

(symbol, trade_id)

используется только на уровне Controller.

Это уменьшает объём памяти

и упрощает внутреннюю структуру состояния.


Жизненный цикл состояния

Build 060.18 определяет простой жизненный цикл.


Создание

Состояние создаётся лениво.

Первое появление сделки данного символа приводит к созданию нового TradeStreamState.


Использование

После создания состояние используется всеми последующими сделками данного символа.


Уничтожение

Build 060.18 не удаляет состояния.

Они существуют до завершения процесса.

Управление жизненным циклом относится к Runtime Layer и будет рассматриваться отдельно.


Публичный контракт Controller

Controller предоставляет единственную публичную операцию.

accept(trade: Trade) -> Trade | None

Других публичных методов Build 060.18 не вводит.


Семантика accept()

Если сделка успешно прошла проверку,

Controller возвращает исходный объект Trade.

Если поступил полный повтор,

возвращается:

None

Если обнаружено нарушение инвариантов,

генерируется соответствующее исключение.


Почему возвращается Trade

Во время проектирования рассматривались альтернативы.


bool

Отклонён.

Причина:

вызывающий код вынужден хранить исходный объект отдельно.


Result Object

Отклонён.

Причина:

избыточен для Build 060.18.


Исключение для дубликатов

Отклонено.

Повтор является нормальной ситуацией.

Он не считается ошибкой.


Trade | None

Принят.

Контракт минимален,

естественен

и легко расширяется в будущем.


Исключения Controller

Build вводит только два новых доменных исключения.


TradeOrderingError

Возникает,

если сделка нарушает монотонность потока.

Например:

100

101

95

TradeConsistencyError

Возникает,

если найден конфликтующий повтор.

Например:

trade_id = 500

price = 100

позже

trade_id = 500

price = 101

Такой поток считается внутренне противоречивым.


Взаимодействие Controller и State

Важнейшим архитектурным принципом Build является инкапсуляция состояния.

Controller никогда не обращается к внутренним структурам данных напрямую.

Вместо этого он взаимодействует со State исключительно через его операции.

Концептуально взаимодействие выглядит следующим образом.

Controller

        │

        ▼

TradeStreamState

        │

        ├── получить последнюю сделку

        ├── получить зарегистрированную сделку

        ├── зарегистрировать новую сделку

        └── поддерживать размер окна

Таким образом:

  • Controller ничего не знает о реализации хранения;
  • State ничего не знает о бизнес-правилах проверки согласованности.

Именно это разделение делает подсистему устойчивой к дальнейшему развитию.


Алгоритм работы TradeStreamConsistencyController

После определения архитектуры компонентов необходимо формально определить алгоритм обработки каждой сделки.

Build 060.18 рассматривает Controller как детерминированный автомат (Deterministic State Machine).

Это означает, что результат обработки полностью определяется двумя величинами:

Текущее состояние

+

Входящая Trade

↓

Новое состояние

+

Результат обработки

Никакие внешние факторы не влияют на принятие решения.


Почему выбран детерминированный автомат

Во время проектирования рассматривались несколько моделей.


Процедурный алгоритм

Обычная последовательность условий.

if

if

if

if

Работает.

Но по мере развития Build начинает быстро усложняться.


Таблица правил

Возможна.

Однако становится плохо читаемой.


Deterministic State Machine

Принята.

Причины:

  • полностью предсказуемое поведение;
  • простое тестирование;
  • возможность восстановления состояния;
  • одинаковое поведение REST и WebSocket;
  • естественное расширение для Recovery Build.

Формальная модель

Для каждой входящей сделки существует единственный возможный результат.

TradeStreamState

+

Trade

↓

TradeStreamState'

+

Result

где

Result представляет собой одно из следующих состояний:

Accepted

Duplicate

TradeOrderingError

TradeConsistencyError

Других исходов Build 060.18 не предусматривает.


Последовательность обработки

Каждая входящая сделка проходит одинаковую последовательность шагов.

Trade

↓

Получить состояние символа

↓

Поиск trade_id

↓

Определение дубликата

↓

Проверка порядка

↓

Регистрация

↓

Публикация

Порядок этапов является обязательным.

Изменение последовательности проверок запрещается.


Почему сначала проверяется дубликат

Во время проектирования рассматривался альтернативный вариант.

Сначала Ordering.

Потом Deduplication.

Например:

100

101

100

Ordering немедленно сообщил бы,

что сделка старая.

Но на самом деле это корректный повтор.

Следовательно,

Ordering не должен выполняться первым.


Официальная последовательность

Контроллер всегда выполняет проверки в следующем порядке.


Шаг 1

Получение состояния символа.

Если состояние отсутствует,

оно создаётся.


Шаг 2

Поиск существующей сделки.

Если trade_id найден,

Controller обязан сравнить всю каноническую модель.


Шаг 3

Определение типа повтора.

Если совпадают все бизнес-поля,

сделка считается полным дубликатом.

Возвращается:

None

Если найдено хотя бы одно различие,

генерируется

TradeConsistencyError

После этого обработка завершается.


Шаг 4

Если сделка новая,

выполняется проверка порядка.

Если

trade_id < last_trade_id

генерируется

TradeOrderingError

Состояние не изменяется.


Шаг 5

Если все проверки успешно завершены,

сделка регистрируется.


Шаг 6

После регистрации

Trade публикуется вызывающему компоненту.


Почему регистрация выполняется перед публикацией

Это принципиальное решение Build.

Рассмотрим последовательность.

Плохой вариант.

Trade

↓

Publication

↓

Registration

Если между публикацией и регистрацией произойдёт исключение,

система потеряет согласованность.

Потребитель уже увидел сделку,

а внутреннее состояние ещё нет.


Правильный вариант.

Trade

↓

Registration

↓

Publication

После регистрации

Controller уже находится в согласованном состоянии.


Неизменяемость состояния при ошибках

Любая ошибка обязана обладать свойством атомарности.

Если возникает:

  • TradeOrderingError;
  • TradeConsistencyError;

состояние Controller не изменяется.

Это означает:

  • last_trade_id остаётся прежним;
  • окно дедупликации не изменяется;
  • история публикаций остаётся неизменной.

Архитектурная диаграмма

Полный путь обработки сделки выглядит следующим образом.

REST / WebSocket

        │

        ▼

Parser

        │

        ▼

Value Validation

        │

        ▼

Mapper

        │

        ▼

Trade

        │

        ▼

TradeStreamConsistencyController

        │

        ▼

Canonical Trade Stream

        │

        ▼

TradesFeed Consumer

После Build 060.18 именно эта схема становится официальной архитектурой Trades Feed.


Интеграция в Acquisition Pipeline

Основной принцип

Build 060.18 не изменяет существующий Acquisition Pipeline.

Он лишь добавляет новый этап.

До Build:

Handler

↓

Consumer

После Build:

Handler

↓

TradeStreamConsistencyController

↓

Consumer

Никакие существующие компоненты не меняют собственную ответственность.


Integration Boundary

TradeStreamConsistencyController располагается между Handler и Consumer.

Он получает уже канонический Trade.

И публикует только Canonical Trade Stream.

Это официальная граница между:

Canonical Trade

↓

Canonical Trade Stream

Dependency Injection

Controller создаётся в Composition Root.

Никакой компонент не создаёт Controller самостоятельно.

Итоговая схема выглядит следующим образом.

Composition Root

        │

        ├──────────────┐

        ▼              ▼

TradesHandler   TradeStreamConsistencyController

        │              │

        └──────┬───────┘

               ▼

          TradesFeed

Таким образом:

  • Feed ничего не знает о реализации Controller;
  • Controller ничего не знает о Feed;
  • Handler ничего не знает о Controller.

Каждый компонент получает только необходимые зависимости.


План изменений файлов

Build 060.18 вводит новую внутреннюю подсистему.

Предлагаемая структура.

market_data/
└── acquisition/
    └── consistency/
        ├── __init__.py
        ├── trade_stream_consistency_controller.py
        ├── trade_stream_state.py
        ├── trade_stream_protocol.py
        └── trade_stream_exceptions.py

trade_stream_consistency_controller.py

Содержит:

TradeStreamConsistencyController

Отвечает исключительно за доменные решения.


trade_stream_state.py

Содержит:

TradeStreamState

Отвечает исключительно за хранение состояния.


trade_stream_exceptions.py

Содержит:

TradeOrderingError

TradeConsistencyError

trade_stream_protocol.py

Содержит публичный контракт подсистемы Stream Consistency.

Все внешние компоненты должны зависеть от данного Protocol, а не от конкретной реализации Controller.


Соглашение об именовании файлов

Build 060.18 закрепляет правило уникальности имён файлов в проекте Dzentra.

Все новые файлы должны иметь уникальные имена в пределах всего репозитория.

Имя файла должно отражать его содержимое, а не только роль внутри текущего каталога.

Например:

trade_stream_consistency_controller.py

trade_stream_state.py

trade_stream_protocol.py

trade_stream_exceptions.py

Данное соглашение является обязательным архитектурным стандартом для всех последующих Build.

Изменяемые файлы

Build 060.18 должен минимально затронуть существующий код.

Изменения предполагаются только в точках интеграции:

  • Composition Root (создание Controller);
  • TradesFeed (внедрение зависимости и последовательность вызовов).

Все остальные изменения должны быть локализованы внутри новой подсистемы consistency.


Architectural Decision Records

ADR-060.18-01 — Canonical Trade Identity

Решение

Идентичность сделки определяется парой:

(symbol, trade_id)

Другие поля в идентичность не входят.


ADR-060.18-02 — Stream State Ownership

Решение

Единственным владельцем состояния потока является TradeStreamConsistencyController, который управляет набором TradeStreamState по символам.


ADR-060.18-03 — Deterministic Stream State Machine

Решение

Контроллер рассматривается как детерминированный автомат:

State + Trade → New State + Result

Поведение полностью определяется текущим состоянием и входящей сделкой.


ADR-060.18-04 — Encapsulated Stream State

Решение

TradeStreamState инкапсулирует все детали хранения.

Контроллер взаимодействует только через его публичные операции.


ADR-060.18-05 — Domain State Ownership

Решение

Разделение ответственности фиксируется следующим образом:

  • TradeStreamConsistencyController — бизнес-правила и доменные решения;
  • TradeStreamState — состояние и структурные инварианты.

ADR-060.18-06 — Stream Boundary

Решение

TradeStreamConsistencyController является единственной точкой формирования Canonical Trade Stream.

До него существует только Canonical Trade.

После него существует только Canonical Trade Stream.

Никакой другой компонент не имеет права изменять или повторно проверять согласованность потока.


Граница Build 060.18

После завершения данного Build система получает:

  • каноническую модель Trade;
  • канонический поток Trade;
  • гарантии порядка;
  • гарантии дедупликации;
  • гарантии неизменяемости опубликованного потока.

На этом ответственность Build заканчивается.

Следующие задачи сознательно оставлены за пределами Build:

  • обнаружение пропусков (Gap Detection);
  • восстановление последовательности (Recovery);
  • синхронизация через REST Backfill;
  • сохранение состояния между перезапусками процесса;
  • диагностика качества соединения;
  • телеметрия и метрики работы контроллера.

Все перечисленные функции относятся к следующим Build серии 060.


Стратегия тестирования

Build 060.18 вводит первый stateful-компонент Acquisition Layer.

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

Тестирование должно подтверждать корректность поведения системы независимо от источника данных.


Основные принципы тестирования

При проектировании тестов используются следующие принципы.


Проверяются инварианты

Основной объект тестирования —

не отдельные методы,

а инварианты Canonical Trade Stream.


Тестируются только публичные контракты

Unit-тесты не должны зависеть от внутренней реализации:

  • FIFO;
  • OrderedDict;
  • структуры хранения;
  • внутренних коллекций.

Все проверки выполняются исключительно через публичный API.


Один тест — одна причина отказа

Каждый негативный сценарий проверяет только одно нарушение инварианта.

Это существенно упрощает анализ ошибок.


Детерминированность

Каждый тест обязан быть полностью воспроизводимым.

Никакие случайные значения,

генераторы,

таймеры,

или текущее время

не должны влиять на результат.


Уровни тестирования

Build 060.18 вводит три уровня проверки.


Unit Tests

Проверяются:

  • TradeStreamState
  • TradeStreamConsistencyController

изоляционно.

Все внешние зависимости заменяются тестовыми объектами.


Integration Tests

Проверяется полный Pipeline.

Trade

↓

TradeStreamConsistencyController

↓

Consumer

Основная задача —

доказать,

что Controller корректно интегрирован в Acquisition.


Regression Tests

Проверяется,

что Build 060.18

не изменил поведение уже существующих Build.

В частности:

  • Parser;
  • Mapper;
  • Validation;
  • Handler.

Матрица тестирования

1. Создание состояния

Проверяется

Первое появление нового symbol.

Ожидаемый результат

Создаётся новый TradeStreamState.


2. Повторное использование состояния

Проверяется

Несколько сделок одного symbol.

Ожидаемый результат

Используется существующее состояние.

Новое не создаётся.


3. Независимость символов

Проверяется

BTC

ETH

BTC

ETH

Ожидаемый результат

Состояния полностью независимы.


4. Первая сделка

Проверяется

Пустой поток.

Ожидаемый результат

Trade публикуется.


5. Возрастающий trade_id

100

101

102

Ожидаемый результат

Все сделки публикуются.


6. Разрыв последовательности

100

105

150

Ожидаемый результат

Ошибок нет.


7. Уменьшение trade_id

100

101

95

Проверяется

Ordering.

Ожидаемый результат

TradeOrderingError.


8. Полный дубликат

100

100

Проверяется

Deduplication.

Ожидаемый результат

Возвращается

None

9. Конфликтующий повтор

trade_id = 500

price = 100

позже

trade_id = 500

price = 101

Проверяется

Consistency.

Ожидаемый результат

TradeConsistencyError.


10. Неизменяемость состояния после ошибки

После возникновения исключения проверяется:

  • last_trade_id;
  • окно дедупликации;
  • количество опубликованных сделок.

Ничего не должно измениться.


11. Работа FIFO

Заполняется окно.

Добавляется ещё одна сделка.

Проверяется корректное удаление самой старой записи.


12. Проверка границы окна

После удаления старой записи

её повтор

не должен считаться известной сделкой.


13. Проверка symbol isolation

BTC

не должен влиять на

ETH.

Даже при совпадающих trade_id.


14. Проверка большого потока

Несколько тысяч последовательных сделок.

Проверяется:

  • отсутствие деградации;
  • корректность состояния;
  • отсутствие ошибок порядка.

Проверка производительности

Build 060.18 не вводит отдельного Benchmark Build.

Однако контроллер обязан удовлетворять следующим требованиям.


Поиск сделки

Средняя сложность

O(1)

Добавление сделки

Средняя сложность

O(1)

Удаление старой записи

Средняя сложность

O(1)

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

Размер памяти ограничивается размером окна дедупликации.

Не допускается бесконечный рост.


Матрица архитектурных инвариантов

Инвариант Проверка
Immutable Trade Unit
Canonical Identity Unit
Ordering Unit
Deduplication Unit
Conflicting Duplicate Unit
State Isolation Unit
FIFO Window Unit
Integration Pipeline Integration
Stateless Handler Integration
Controller Boundary Integration

Definition of Done

Build считается завершённым только после выполнения всех перечисленных условий.


Архитектура

  • Все ADR реализованы без отклонений.
  • Не нарушены границы ответственности компонентов.
  • Controller остаётся Domain Service.
  • State остаётся Domain State.

Код

  • Не изменена модель Trade.
  • Не изменены Parser.
  • Не изменены Mapper.
  • Не изменена Validation.
  • Handler остаётся stateless.

Функциональность

Поддерживаются:

  • Ordering;
  • Deduplication;
  • Conflict Detection;
  • Symbol Isolation.

Производительность

Все операции соответствуют требованиям:

O(1)

Тестирование

Все Unit Tests проходят.

Все Integration Tests проходят.

Regression Tests проходят без изменений.


Документация

Обновлены:

  • Build Documentation;
  • ADR;
  • Architecture Diagram;
  • File Plan.

Пошаговый план реализации

Реализация Build должна выполняться строго по следующему порядку.


Шаг 1

Создание структуры каталога

consistency/


Шаг 2

Создание файла trade_stream_exceptions.py.


Шаг 3

Создание файла trade_stream_state.py.


Шаг 4

Реализация FIFO Window.


Шаг 5

Создание файла trade_stream_consistency_controller.py.


Шаг 6

Создание файла trade_stream_protocol.py.


Шаг 7

Интеграция через Composition Root.


Шаг 8

Интеграция с TradesFeed.


Шаг 9

Unit Tests.


Шаг 10

Integration Tests.


Шаг 11

Regression Tests.


Build Boundary

После завершения Build 060.18 подсистема получения сделок считается архитектурно завершённой с точки зрения согласованности потока.

Следующим этапом развития становится Build 060.19.

Его задачами являются:

  • обнаружение пропусков последовательности;
  • восстановление пропущенных сделок;
  • интеграция REST Backfill;
  • повторная синхронизация потока;
  • подготовка к непрерывной работе Trade Feed.

Build 060.18 сознательно не реализует перечисленные функции.

Он создаёт архитектурный фундамент, на который будут последовательно опираться все последующие Build серии 060.


Заключение

Build 060.18 завершает формирование базовой архитектуры Trades Feed (Time & Sales).

Если предыдущие Build определяли, что такое каноническая сделка (Trade), то Build 060.18 впервые формализует, что такое канонический поток сделок (Canonical Trade Stream).

Ключевым результатом является введение TradeStreamConsistencyController — единственной точки формирования согласованного потока, и TradeStreamState — единственного владельца состояния этого потока.

Зафиксированные в документе инварианты, контракты и ADR обеспечивают:

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

Таким образом, Build 060.18 становится фундаментом для следующих этапов серии 060, включая Gap Detection, Recovery, REST Backfill и построение полностью отказоустойчивого Trade Feed.