diff --git a/docs/migrations/build_060_18.md b/docs/migrations/build_060_18.md new file mode 100644 index 0000000..c7aad47 --- /dev/null +++ b/docs/migrations/build_060_18.md @@ -0,0 +1,1492 @@ +# Build 060.18 — Trade Stream Consistency Controller + +**Engineering Migration Report** + +--- + +# Контроль документа + +| Свойство | Значение | +|----------|----------| +| Build | 060.18 | +| Название | Trade Stream Consistency Controller | +| Статус | Completed | +| Проект | Dzentra | +| Подсистема | Market Data Acquisition | +| Компонент | Trade Stream Consistency | +| Версия | 1.0 | + +--- + +# Связанные документы + +- build_060_18_architecture.md — архитектурная спецификация Build. +- build_060_17.md — Engineering Migration Report предыдущего Build. + +--- + +# Цель Build + +После завершения Build 060.17 подсистема Market Data Acquisition получила полноценную инфраструктуру формирования канонической модели сделки (`Trade`). + +К этому моменту архитектура уже обеспечивала: + +- получение транспортных сообщений из различных источников; +- преобразование транспортных моделей в канонический объект `Trade`; +- единый механизм Parser; +- Value Validation; +- Mapper; +- Trade Adapter; +- инфраструктуру Feed. + +Таким образом система уже умела получать отдельные корректные сделки независимо от источника их происхождения. + +Однако корректность отдельной сделки ещё не означает корректность последовательности сделок. + +Поток данных, поступающий от биржи, может содержать: + +- повторную доставку уже опубликованных сделок; +- нарушение порядка поступления сообщений; +- конфликтующие дубликаты; +- повторную передачу одной и той же сделки через различные транспортные каналы. + +До настоящего Build подобные ситуации никак не контролировались. + +Каждый компонент, получающий объект `Trade`, был вынужден самостоятельно предполагать, что поток уже является корректным. + +Подобная архитектура противоречила фундаментальному принципу Dzentra. + +Согласованность потока должна обеспечиваться централизованно. + +Потребители рыночных данных не должны повторно выполнять проверку порядка или дедупликацию. + +Главной задачей Build 060.18 становится построение специализированной подсистемы **Trade Stream Consistency**, обеспечивающей формирование единственного канонического потока сделок внутри Acquisition Layer. + +После завершения Build система получает: + +- специализированный `TradeStreamConsistencyProtocol`; +- специализированный `TradeStreamConsistencyController`; +- внутреннюю модель состояния `TradeStreamState`; +- специализированные доменные исключения; +- полноценное unit-тестирование новой подсистемы. + +При этом Build принципиально не затрагивает: + +- Parser; +- Mapper; +- Value Validation; +- Trade Adapter; +- Handler; +- Feed; +- Runtime; +- Subscription Layer; +- REST Backfill; +- Gap Recovery; +- Reconnect; +- интеграцию с Runtime. + +Все перечисленные задачи относятся к следующим этапам развития подсистемы Trades Feed. + +--- + +# Предпосылки + +К началу Build архитектура обработки сделок уже обеспечивала формирование единой канонической модели `Trade`. + +Полный конвейер обработки выглядел следующим образом. + +```text +Transport Message + │ + ▼ +Schema Validation + │ + ▼ +Parser + │ + ▼ +Value Validation + │ + ▼ +Mapper + │ + ▼ +Trade +``` + +Каждый уровень обладал строго определённой областью ответственности. + +Schema Validation отвечала за корректность транспортного документа. + +Parser извлекал необходимые поля. + +Value Validation проверяла корректность отдельных значений. + +Mapper строил каноническую модель предметной области. + +Полученный объект `Trade` уже являлся полностью независимым от транспортного формата и мог использоваться всеми последующими компонентами системы. + +Однако существовал один принципиальный архитектурный пробел. + +Система гарантировала корректность каждой отдельной сделки, но не гарантировала корректность последовательности этих сделок. + +Например, следующий поток состоял исключительно из корректных объектов `Trade`. + +```text +Trade #100 + +Trade #101 + +Trade #100 +``` + +Каждая сделка по отдельности являлась полностью корректной. + +Однако сам поток нарушал инвариант уникальности публикации. + +Аналогично поток + +```text +Trade #100 + +Trade #101 + +Trade #99 +``` + +также состоял из корректных объектов `Trade`, но нарушал инвариант порядка. + +Следовательно, между построением объекта `Trade` и его публикацией должен существовать самостоятельный уровень контроля согласованности потока. + +Именно эту архитектурную задачу решает Build 060.18. + +--- + +# Результаты архитектурного аудита + +Перед началом реализации Build был выполнен полный аудит существующей подсистемы **Market Data Acquisition**. + +Целью аудита являлась проверка соответствия фактической реализации архитектурной модели, утверждённой в `build_060_18_architecture.md`, а также определение точек интеграции новой подсистемы Stream Consistency. + +Особое внимание уделялось существующей инфраструктуре получения сделок. + +Первоначально предполагалось, что после реализации `TradeStreamConsistencyController` он будет интегрирован непосредственно в существующий `TradesFeed`. + +Именно такое решение рассматривалось на этапе архитектурного проектирования. + +Однако проведённый аудит показал, что фактическая структура подсистемы отличается от первоначальных предположений. + +--- + +## Анализ существующего Trades Feed + +В ходе проверки был полностью проанализирован существующий компонент: + +```text +feeds/trades_feed.py +``` + +Аудит показал, что данный компонент реализует исключительно сценарий получения **одиночной сделки** посредством REST API. + +Его публичный контракт имеет следующий вид. + +```text +load_trade(symbol) + │ + ▼ +Trade +``` + +Feed не содержит: + +- непрерывного потока сообщений; +- подписок WebSocket; +- обработки последовательности сделок; +- собственного внутреннего состояния; +- механизма публикации Trade Consumer; +- инфраструктуры обработки событий. + +Каждый вызов `load_trade()` полностью независим от предыдущих вызовов. + +Таким образом существующий `TradesFeed` представляет собой stateless-компонент, предназначенный исключительно для получения отдельных канонических объектов `Trade`. + +--- + +## Отсутствие потокового Feed + +Дополнительно был проведён аудит инфраструктуры WebSocket. + +Проверка показала, что к моменту реализации Build 060.18 в проекте отсутствует специализированный компонент, отвечающий за обработку непрерывного потока сделок. + +Иными словами, архитектура уже содержит: + +- WebSocket Runtime; +- Subscription Layer; +- Router; +- Handler; +- Adapter; +- каноническую модель `Trade`; + +но ещё не содержит специализированного **WebSocket Trades Feed**, который бы организовывал непрерывное получение и публикацию последовательности сделок. + +Это стало ключевым результатом архитектурного аудита. + +--- + +## Последствия для реализации Build + +Полученные результаты существенно повлияли на окончательную реализацию Build. + +Интеграция `TradeStreamConsistencyController` в существующий REST Feed привела бы к смешению двух различных архитектурных уровней. + +REST Feed отвечает за получение отдельной сделки. + +Trade Stream Consistency отвечает за обработку непрерывного потока сделок. + +Эти два сценария обладают различной природой и различными требованиями к состоянию системы. + +Попытка объединить их в рамках одного компонента нарушила бы принцип единственной ответственности и создала бы искусственную зависимость между REST и потоковой обработкой. + +Поэтому было принято решение отказаться от подобной интеграции. + +--- + +# Архитектурное решение + +По результатам проведённого аудита было принято решение оставить новую подсистему **Trade Stream Consistency** полностью самостоятельной. + +В рамках Build реализуются только компоненты, непосредственно отвечающие за проверку согласованности потока. + +Интеграция с инфраструктурой получения данных сознательно переносится на последующие этапы развития Trades Feed. + +После завершения Build архитектура принимает следующий вид. + +```text +Trade + │ + ▼ +TradeStreamConsistencyController + │ + ▼ +Canonical Trade Stream +``` + +Таким образом Build 060.18 завершает построение самостоятельного слоя согласованности потока, не изменяя существующую инфраструктуру получения сделок. + +Это решение обеспечивает слабую связанность компонентов и позволяет независимо развивать: + +- инфраструктуру получения данных; +- механизмы восстановления потока; +- обработку разрывов последовательности; +- механизмы повторной синхронизации. + +Все перечисленные возможности смогут использовать уже готовый `TradeStreamConsistencyController` без изменения его внутренней реализации. + +--- + +# Почему Controller не интегрирован в существующий Feed + +На этапе архитектурного проектирования предполагалось, что новой подсистеме потребуется непосредственная интеграция в `TradesFeed`. + +Однако инженерный аудит показал, что такая интеграция является преждевременной. + +Причина заключается в различии ответственности компонентов. + +`TradesFeed` в текущей реализации получает отдельную сделку. + +`TradeStreamConsistencyController` принимает решения исключительно относительно непрерывного потока сделок. + +Следовательно, контроллер не может эффективно использоваться до появления полноценного потокового Feed. + +В результате было принято следующее окончательное решение Build. + +Настоящий Build завершает реализацию самостоятельной подсистемы Stream Consistency, полностью готовой к использованию. + +Её интеграция будет выполнена после появления специализированного **WebSocket Trades Feed**, который станет источником непрерывного потока сделок. + +Подобное решение позволило сохранить архитектурную чистоту проекта и избежать появления технического долга на раннем этапе развития подсистемы. + +--- + +# Новая подсистема Trade Stream Consistency + +Главным результатом настоящего Build становится появление в подсистеме **Market Data Acquisition** нового архитектурного уровня — **Trade Stream Consistency**. + +До начала Build система завершала обработку сделки сразу после построения канонической модели `Trade`. + +Конвейер обработки имел следующий вид. + +```text +Transport Message + │ + ▼ +Schema Validation + │ + ▼ +Parser + │ + ▼ +Value Validation + │ + ▼ +Mapper + │ + ▼ +Trade +``` + +После завершения Build между построением канонической модели и её публикацией появляется дополнительный уровень. + +```text +Transport Message + │ + ▼ +Schema Validation + │ + ▼ +Parser + │ + ▼ +Value Validation + │ + ▼ +Mapper + │ + ▼ +Trade + │ + ▼ +Trade Stream Consistency + │ + ▼ +Canonical Trade Stream +``` + +Появление данного уровня является принципиальным изменением архитектуры Acquisition Layer. + +Если ранее система гарантировала корректность отдельных объектов `Trade`, то теперь она гарантирует корректность всей последовательности опубликованных сделок. + +Именно последовательность становится новой доменной сущностью. + +--- + +# Архитектурное решение + +Во время проектирования рассматривались несколько вариантов реализации проверки согласованности. + +Первый вариант предполагал распределение логики между различными компонентами Acquisition Layer. + +Например: + +- часть проверки выполнять внутри Feed; +- часть — внутри Handler; +- часть — внутри Runtime. + +После анализа архитектуры данный подход был отклонён. + +Подобное распределение приводило к нескольким серьёзным недостаткам. + +Во-первых, логика проверки потока оказывалась размазанной между различными уровнями системы. + +Во-вторых, различные источники данных могли реализовывать разные правила проверки. + +В-третьих, последующее развитие механизмов Recovery и REST Backfill существенно усложнялось. + +Поэтому было принято другое решение. + +Вся логика проверки согласованности концентрируется внутри специализированной подсистемы. + +Она становится единственной точкой формирования канонического потока сделок. + +--- + +# Архитектура новой подсистемы + +В рамках Build реализованы четыре новых компонента. + +```text +TradeStreamConsistencyProtocol + +TradeStreamConsistencyController + +TradeStreamState + +Trade Stream Exceptions +``` + +Каждый компонент обладает собственной областью ответственности. + +Ни один компонент не выполняет обязанности другого. + +Подобное разделение полностью соответствует принципу **Single Responsibility**, принятому в архитектуре Dzentra. + +--- + +# TradeStreamConsistencyProtocol + +Одной из целей Build являлось формирование полноценного контрактного уровня новой подсистемы. + +До начала Build соответствующий Protocol отсутствовал. + +В рамках реализации был добавлен новый контракт. + +```text +TradeStreamConsistencyProtocol +``` + +Protocol определяет единственную публичную операцию. + +```text +accept(trade) + │ + ▼ +Trade | None +``` + +Никаких других обязанностей Protocol не содержит. + +Он не определяет: + +- внутреннее состояние; +- способы хранения данных; +- размер окна дедупликации; +- алгоритмы проверки; +- механизмы публикации. + +Все перечисленные детали относятся исключительно к реализации Controller. + +Благодаря подобному разделению любые последующие компоненты системы смогут зависеть только от контракта, а не от конкретной реализации. + +Это полностью соответствует принципу **Dependency Inversion**, принятому в проекте Dzentra. + +--- + +# Почему выбран минимальный Protocol + +Во время проектирования рассматривались различные варианты публичного API. + +В частности анализировались варианты: + +- возврата логического значения (`bool`); +- использования специализированного объекта результата; +- публикации событий вместо возврата значения. + +После анализа было принято решение оставить контракт максимально простым. + +Контроллер принимает объект `Trade` и возвращает либо этот же объект, либо `None`. + +Нарушения архитектурных инвариантов выражаются специализированными исключениями. + +Такой контракт оказался наиболее устойчивым к дальнейшему развитию системы. + +Он одинаково хорошо подходит для: + +- WebSocket Feed; +- REST Backfill; +- Recovery Pipeline; +- Integration Tests; +- Unit Tests. + +При этом публичный интерфейс остаётся минимальным и легко читаемым. + +--- + +# Новые доменные исключения + +Следующим результатом Build становится появление специализированных исключений подсистемы Stream Consistency. + +До настоящего Build подобные ошибки отсутствовали. + +В рамках реализации добавлены два новых класса. + +```text +TradeOrderingError + +TradeConsistencyError +``` + +Каждый тип ошибки отражает отдельное нарушение архитектурных инвариантов канонического потока. + +Разделение исключений позволяет вызывающему компоненту принимать различные решения в зависимости от характера проблемы. + +Например, нарушение порядка и конфликтующий дубликат имеют различную природу и требуют различной стратегии обработки. + +Поэтому использование единственного общего исключения было признано нецелесообразным. + +Оба класса наследуются от общей иерархии исключений подсистемы Market Data Acquisition и полностью соответствуют существующей архитектуре проекта. + +--- + +# TradeStreamConsistencyController + +Центральным компонентом настоящего Build становится + +```text +TradeStreamConsistencyController +``` + +Именно он завершает формирование новой подсистемы **Trade Stream Consistency** внутри Acquisition Layer. + +До начала Build все компоненты системы были stateless. + +Parser не хранил состояние. + +Value Validation не хранила состояние. + +Mapper не хранил состояние. + +Trade Adapter не хранил состояние. + +Handler не хранил состояние. + +Feed также не содержал собственного состояния. + +Появление Trade Stream Consistency впервые вводит в подсистему компонент, принимающий решения на основании ранее обработанных сделок. + +Именно поэтому Controller становится первой stateful-службой внутри Acquisition Layer. + +--- + +## Архитектура Controller + +Конструкция Controller намеренно сделана максимально простой. + +Он содержит только одно внутреннее хранилище. + +```text +symbol + │ + ▼ +TradeStreamState +``` + +Для каждого торгового символа существует собственный экземпляр состояния. + +Например, + +```text +BTCUSDT + │ + ▼ +TradeStreamState +``` + +и + +```text +ETHUSDT + │ + ▼ +TradeStreamState +``` + +обслуживаются полностью независимо друг от друга. + +Controller не хранит информацию о самих сделках. + +Он лишь определяет, какому состоянию необходимо передать очередную сделку для проверки. + +--- + +## Ответственность Controller + +Во время проектирования особое внимание уделялось разделению ответственности между компонентами новой подсистемы. + +В результате Controller получил исключительно координационные обязанности. + +Он отвечает за: + +- выбор состояния по symbol; +- ленивое создание нового состояния; +- маршрутизацию сделки; +- возврат результата проверки вызывающему компоненту. + +При этом Controller сознательно не реализует: + +- алгоритм дедупликации; +- проверку порядка; +- хранение окна сделок; +- сравнение Trade; +- управление FIFO. + +Все перечисленные задачи полностью принадлежат TradeStreamState. + +Подобное разделение значительно упрощает дальнейшее развитие системы. + +--- + +# TradeStreamState + +Вторым ключевым компонентом новой подсистемы становится + +```text +TradeStreamState +``` + +Если Controller представляет собой уровень координации, то State представляет собой уровень хранения состояния и проверки инвариантов потока. + +Каждый экземпляр TradeStreamState обслуживает только один торговый символ. + +Это является одним из фундаментальных архитектурных принципов Build. + +Благодаря подобному решению состояние различных инструментов никогда не смешивается между собой. + +Поток BTCUSDT полностью независим от потока ETHUSDT. + +Даже совпадающие значения `trade_id` не создают никаких конфликтов. + +--- + +## Внутреннее состояние + +Каждый экземпляр TradeStreamState хранит минимальный объём информации, необходимый для проверки согласованности потока. + +Внутреннее состояние включает: + +```text +symbol + +last_trade_id + +FIFO Window + +Dictionary Trade Cache +``` + +Этого набора данных достаточно для реализации всех архитектурных инвариантов Build 060.18. + +--- + +## last_trade_id + +Поле + +```text +last_trade_id +``` + +хранит максимальный идентификатор сделки, успешно опубликованной системой для данного символа. + +Следует подчеркнуть, что речь идёт именно об опубликованной сделке. + +Если поступает сделка, нарушающая порядок, + +```text +100 + +101 + +95 +``` + +то состояние не изменяется. + +После возникновения `TradeOrderingError` + +значение + +```text +last_trade_id +``` + +остаётся равным + +```text +101 +``` + +Подобное поведение обеспечивает атомарность всех операций Controller. + +Ошибочная сделка никогда не влияет на состояние потока. + +--- + +## FIFO Window + +Для проверки повторов TradeStreamState использует ограниченное окно ранее опубликованных сделок. + +Во время проектирования рассматривались различные варианты хранения. + +В частности анализировались: + +- полная история сделок; +- HashSet; +- OrderedDict; +- LRU Cache; +- Ring Buffer. + +После анализа было принято решение использовать ограниченное FIFO-окно. + +Такое решение наиболее точно соответствует природе непрерывного потока сделок. + +По мере поступления новых сделок самые старые записи постепенно удаляются из памяти. + +Это позволяет ограничить использование памяти независимо от времени работы процесса. + +--- + +## Используемые структуры данных + +Внутренняя реализация TradeStreamState построена на совместном использовании двух стандартных структур данных Python. + +Для хранения порядка поступления используется + +```text +deque +``` + +Для быстрого поиска зарегистрированных сделок используется + +```text +dict +``` + +Совместное использование этих структур обеспечивает выполнение всех основных операций со средней сложностью + +```text +O(1) +``` + +В частности: + +- поиск зарегистрированной сделки; +- регистрация новой сделки; +- удаление самой старой записи; +- поддержание фиксированного размера окна. + +Подобная комбинация оказалась наиболее простой, эффективной и полностью удовлетворяющей требованиям настоящего Build. + +--- + +## Размер окна дедупликации + +По умолчанию TradeStreamState использует окно размером + +```text +10 000 +``` + +сделок. + +Данное значение было выбрано как разумный компромисс между объёмом используемой памяти и вероятностью повторной доставки одной и той же сделки. + +При этом размер окна не является архитектурным ограничением. + +Он задаётся отдельным параметром при создании состояния и может быть изменён без каких-либо изменений внутреннего алгоритма Controller. + +Таким образом Build фиксирует только принцип использования ограниченного окна, но не навязывает конкретный объём хранения для всех последующих реализаций. + +--- + +# Алгоритм обработки сделки + +После завершения реализации Build поведение новой подсистемы становится полностью детерминированным. + +Каждая входящая сделка проходит одну и ту же последовательность проверок. + +Результат обработки зависит исключительно от текущего состояния потока и содержимого входящей сделки. + +Никакие внешние факторы не влияют на принятие решения. + +Полный алгоритм обработки выглядит следующим образом. + +```text +Trade + │ + ▼ +Получение состояния symbol + │ + ▼ +Поиск trade_id + │ + ├─────────────── Найден ───────────────┐ + │ │ + ▼ ▼ +Сравнение объектов Trade Новый trade_id + │ │ + ├─────────────── Совпадают ─────────────┤ + │ │ + ▼ ▼ +Duplicate Проверка порядка + │ │ + ▼ ▼ +return None trade_id < last_trade_id ? + │ + ┌───────────────┴───────────────┐ + ▼ ▼ + TradeOrderingError Регистрация сделки + │ + ▼ + Обновление состояния + │ + ▼ + return Trade +``` + +Подобная последовательность является обязательной. + +Изменение порядка выполнения проверок приведёт к нарушению архитектурных инвариантов новой подсистемы. + +--- + +# Почему сначала проверяется дедупликация + +Во время проектирования отдельно анализировалась последовательность выполнения проверок. + +На первый взгляд могло показаться естественным сначала проверять порядок, а уже затем выполнять поиск повторов. + +Однако данный вариант оказался ошибочным. + +Рассмотрим следующий поток. + +```text +100 + +101 + +100 +``` + +Если первой выполняется проверка порядка, система обнаруживает уменьшение `trade_id` и немедленно генерирует `TradeOrderingError`. + +Однако в действительности последняя сделка не является нарушением порядка. + +Она представляет собой корректную повторную доставку уже опубликованной сделки. + +Следовательно, подобная ситуация должна обрабатываться как обычный Duplicate. + +Именно поэтому поиск зарегистрированной сделки всегда выполняется раньше проверки порядка. + +Данное решение было окончательно закреплено в процессе реализации и подтверждено соответствующими unit-тестами. + +--- + +# Проверка порядка + +Если входящая сделка отсутствует в окне дедупликации, она рассматривается как новая. + +После этого выполняется проверка монотонности последовательности. + +Единственным критерием является значение + +```text +trade_id +``` + +Для каждого символа должно выполняться условие. + +```text +trade_id(new) >= last_trade_id +``` + +При этом система сознательно допускает наличие разрывов последовательности. + +Например, + +```text +100 + +101 + +150 +``` + +является полностью корректным потоком. + +Настоящий Build не занимается анализом причин возникновения подобных разрывов. + +Он лишь фиксирует факт отсутствия нарушения порядка. + +Задача обнаружения пропусков относится к следующему этапу развития подсистемы. + +--- + +# Обработка повторов + +После определения идентичности сделки возможны два различных варианта повторной доставки. + +Первый вариант представляет собой полный повтор ранее опубликованной сделки. + +Во втором случае идентичность совпадает, однако содержимое сделки отличается. + +Эти ситуации обладают различной природой и требуют различной реакции системы. + +--- + +## Полный дубликат + +Если ранее зарегистрированная сделка полностью совпадает с новой по всем каноническим полям, повтор считается корректным. + +Контроллер не публикует такую сделку повторно. + +Ошибки при этом не возникает. + +Метод + +```python +accept() +``` + +возвращает + +```python +None +``` + +Подобное поведение рассматривается как нормальная рабочая ситуация при повторной доставке сообщений транспортным уровнем. + +--- + +## Конфликтующий дубликат + +Совершенно иной характер имеет ситуация, при которой идентичность сделки совпадает, но содержимое отличается. + +Например, + +```text +trade_id = 500 + +price = 100 +``` + +позже + +```text +trade_id = 500 + +price = 101 +``` + +Подобная ситуация означает внутреннее противоречие данных. + +С точки зрения канонической модели две различные сделки не могут обладать одинаковой идентичностью. + +В этом случае Controller немедленно прекращает обработку и генерирует + +```text +TradeConsistencyError +``` + +Состояние потока при этом остаётся неизменным. + +--- + +# Атомарность операций + +Одним из фундаментальных требований настоящего Build являлось обеспечение атомарности обработки каждой сделки. + +Любая ошибка должна приводить к полному откату текущей операции. + +Если в процессе проверки возникает: + +- `TradeOrderingError`; +- `TradeConsistencyError`; + +никакие изменения внутреннего состояния не выполняются. + +Не изменяются: + +- `last_trade_id`; +- окно дедупликации; +- словарь зарегистрированных сделок. + +Таким образом после возникновения ошибки Controller остаётся в том же состоянии, в котором находился до начала обработки. + +Подобное решение существенно упрощает последующую реализацию Recovery Pipeline и гарантирует внутреннюю согласованность состояния независимо от количества ошибок транспортного уровня. + +--- + +# Производительность + +Во время проектирования новой подсистемы одним из обязательных требований являлось сохранение постоянной сложности основных операций. + +Поскольку Controller будет использоваться при обработке непрерывного потока сделок, любые алгоритмы с линейной сложностью быстро стали бы узким местом всей Acquisition Layer. + +Именно поэтому внутренняя реализация была построена таким образом, чтобы обеспечить среднюю сложность + +```text +O(1) +``` + +для всех наиболее часто выполняемых операций. + +В частности: + +| Операция | Средняя сложность | +|----------|-------------------| +| Поиск зарегистрированной сделки | O(1) | +| Регистрация новой сделки | O(1) | +| Проверка дубликата | O(1) | +| Удаление самой старой записи | O(1) | +| Получение состояния symbol | O(1) | + +Использование памяти ограничивается исключительно размером окна дедупликации. + +В результате объём памяти остаётся постоянным независимо от продолжительности работы процесса. + +Это позволяет использовать Controller в составе долгоживущих потоковых сервисов без риска неограниченного роста потребления памяти. + +--- + +# Изменённые файлы + +В рамках Build были добавлены четыре новых компонента подсистемы **Trade Stream Consistency**. + +Все изменения были сознательно локализованы внутри нового каталога + +```text +src/market_data/acquisition/consistency/ +``` + +Подобное решение позволило полностью изолировать новую функциональность от уже существующих компонентов Acquisition Layer. + +Ни один ранее реализованный Parser, Mapper, Handler или Adapter не потребовал изменений. + +--- + +## Trade Stream Protocol + +```text +src/market_data/acquisition/consistency/trade_stream_protocol.py +``` + +Добавлен новый Protocol. + +```text +TradeStreamConsistencyProtocol +``` + +Protocol определяет единственный публичный контракт новой подсистемы. + +```python +accept(trade: Trade) -> Trade | None +``` + +Благодаря этому все последующие компоненты системы смогут зависеть исключительно от абстракции, а не от конкретной реализации Controller. + +--- + +## Trade Stream Exceptions + +```text +src/market_data/acquisition/consistency/trade_stream_exceptions.py +``` + +Добавлены два специализированных доменных исключения. + +```text +TradeOrderingError + +TradeConsistencyError +``` + +Оба класса наследуются от общей иерархии исключений Acquisition Layer и используются исключительно новой подсистемой Stream Consistency. + +Разделение ошибок позволяет вызывающему компоненту различать нарушение порядка и внутреннюю противоречивость потока. + +--- + +## Trade Stream State + +```text +src/market_data/acquisition/consistency/trade_stream_state.py +``` + +Реализована внутренняя модель состояния одного торгового символа. + +TradeStreamState отвечает за: + +- хранение последнего опубликованного `trade_id`; +- проверку принадлежности символа; +- дедупликацию сделок; +- обнаружение конфликтующих повторов; +- поддержку ограниченного FIFO-окна; +- обновление внутреннего состояния после успешной публикации сделки. + +Вся бизнес-логика проверки согласованности сосредоточена именно внутри данного компонента. + +--- + +## Trade Stream Consistency Controller + +```text +src/market_data/acquisition/consistency/trade_stream_consistency_controller.py +``` + +Реализован компонент + +```text +TradeStreamConsistencyController +``` + +Controller отвечает исключительно за: + +- выбор состояния по symbol; +- ленивое создание новых состояний; +- маршрутизацию сделки; +- возврат результата вызывающему компоненту. + +Внутренняя логика проверки полностью делегируется соответствующему экземпляру `TradeStreamState`. + +Благодаря подобному разделению Controller остаётся компактным координационным компонентом и не зависит от деталей хранения состояния. + +--- + +# Добавленные unit-тесты + +Настоящий Build сопровождается полноценным покрытием новой подсистемы unit-тестами. + +Все тесты написаны исключительно через публичный API компонентов. + +Внутренние структуры данных не используются напрямую. + +Подобный подход позволяет свободно изменять внутреннюю реализацию без изменения тестового набора. + +--- + +## TradeStreamState + +Добавлен новый файл. + +```text +tests/unit/market_data/acquisition/consistency/test_trade_stream_state.py +``` + +Проверяются следующие сценарии. + +--- + +### Создание состояния + +Подтверждается корректная инициализация нового состояния. + +Проверяются: + +- symbol; +- размер окна; +- отсутствие зарегистрированных сделок. + +--- + +### Первая сделка + +Подтверждается успешная публикация первой сделки нового символа. + +--- + +### Возрастающий trade_id + +Проверяется корректная обработка монотонно возрастающей последовательности сделок. + +--- + +### Разрыв последовательности + +Подтверждается, что пропуски идентификаторов не рассматриваются как ошибка. + +--- + +### Полный дубликат + +Подтверждается возврат + +```python +None +``` + +при повторной доставке идентичной сделки. + +--- + +### Конфликтующий дубликат + +Проверяется генерация + +```text +TradeConsistencyError +``` + +при несовпадении содержимого сделки. + +--- + +### Нарушение порядка + +Проверяется генерация + +```text +TradeOrderingError +``` + +при уменьшении `trade_id`. + +--- + +### Работа FIFO + +Проверяется корректное удаление наиболее старой записи после переполнения окна. + +--- + +### Повтор после выхода из окна + +Подтверждается, что сделка, удалённая из окна дедупликации, больше не рассматривается как известная системе. + +--- + +### Проверка symbol + +Подтверждается невозможность передачи сделки другого торгового символа. + +--- + +### Проверка конфигурации + +Проверяется корректная обработка недопустимого размера окна дедупликации. + +--- + +## TradeStreamConsistencyController + +Добавлен новый файл. + +```text +tests/unit/market_data/acquisition/consistency/test_trade_stream_consistency_controller.py +``` + +Проверяются следующие сценарии. + +--- + +### Ленивое создание состояния + +Подтверждается автоматическое создание нового `TradeStreamState` при первом появлении символа. + +--- + +### Повторное использование состояния + +Проверяется, что для всех последующих сделок одного символа используется уже существующее состояние. + +--- + +### Независимость символов + +Подтверждается полная независимость состояний различных торговых инструментов. + +--- + +### Обработка полного дубликата + +Подтверждается корректная передача результата + +```python +None +``` + +вызывающему компоненту. + +--- + +### Передача исключений + +Проверяется, что Controller не подавляет: + +- `TradeOrderingError`; +- `TradeConsistencyError`; + +а корректно передаёт их вызывающему компоненту. + +--- + +# Результаты тестирования + +После завершения реализации был выполнен запуск полного набора unit-тестов новой подсистемы. + +Использовались команды. + +```bash +python -m pytest -q \ +tests/unit/market_data/acquisition/consistency/test_trade_stream_state.py +``` + +Результат. + +```text +13 passed +``` + +После этого была выполнена проверка Controller. + +```bash +python -m pytest -q \ +tests/unit/market_data/acquisition/consistency/test_trade_stream_consistency_controller.py +``` + +Результат. + +```text +6 passed +``` + +Итоговый результат настоящего Build. + +```text +19 passed +``` + +Все предусмотренные сценарии успешно пройдены. + +Тестирование подтвердило: + +- корректность проверки порядка; +- корректность дедупликации; +- независимость символов; +- атомарность операций; +- корректную работу FIFO-окна; +- соответствие Controller утверждённой архитектуре. + +Ни одного отклонения от архитектурной спецификации обнаружено не было. + +--- + +# Итоги Build + +Build 060.18 полностью завершает построение слоя **Trade Stream Consistency** внутри подсистемы Market Data Acquisition. + +До начала настоящего Build система гарантировала корректность отдельных объектов `Trade`. + +После завершения Build система дополнительно гарантирует корректность последовательности публикуемых сделок. + +Таким образом ответственность Acquisition Layer расширяется. + +Теперь подсистема обеспечивает: + +- построение канонической модели сделки; +- проверку корректности потока; +- обнаружение повторной доставки; +- обнаружение конфликтующих дубликатов; +- контроль монотонности последовательности; +- формирование единственного канонического потока сделок. + +При этом все существующие компоненты системы продолжают работать без изменений. + +Build не нарушил обратную совместимость и не потребовал модификации ранее реализованных Parser, Mapper, Adapter или Handler. + +--- + +# Архитектурный результат + +Главным архитектурным результатом настоящего Build становится появление нового самостоятельного уровня Acquisition Layer. + +Теперь общая архитектура обработки сделок принимает следующий вид. + +```text +Transport Message + │ + ▼ +Schema Validation + │ + ▼ +Parser + │ + ▼ +Value Validation + │ + ▼ +Mapper + │ + ▼ +Trade + │ + ▼ +Trade Stream Consistency + │ + ▼ +Canonical Trade Stream +``` + +Данный уровень полностью изолирован от транспортной реализации и работает исключительно с канонической моделью предметной области. + +Это означает, что независимо от источника данных — WebSocket, REST Backfill или Recovery Pipeline — все сделки будут проходить через единый механизм проверки согласованности. + +Подобное решение исключает дублирование логики и гарантирует единообразное поведение системы. + +--- + +# Что сознательно НЕ реализовано + +В соответствии с утверждёнными границами Build 060.18 ряд задач был сознательно оставлен за пределами реализации. + +Настоящий Build **не включает**: + +- интеграцию с существующим REST `TradesFeed`; +- реализацию WebSocket Trades Feed; +- обнаружение пропусков последовательности (`Gap Detection`); +- автоматический REST Backfill; +- механизм восстановления потока (`Recovery Pipeline`); +- повторную синхронизацию после разрыва соединения; +- публикацию событий подписчикам; +- буферизацию или повторную доставку сообщений; +- управление жизненным циклом WebSocket-соединений. + +Все перечисленные возможности относятся к последующим этапам развития подсистемы и будут использовать уже реализованный слой Trade Stream Consistency в качестве готового архитектурного фундамента. + +--- + +# Влияние на последующие Build + +Реализация настоящего Build существенно упрощает дальнейшее развитие подсистемы Trades Feed. + +Следующие этапы смогут опираться на уже готовый механизм проверки согласованности и сосредоточиться исключительно на задачах получения и восстановления потока данных. + +В частности, слой Trade Stream Consistency станет обязательной частью конвейера обработки: + +```text +WebSocket + │ + ▼ +Trade Adapter + │ + ▼ +WebSocket Trades Feed + │ + ▼ +Trade Stream Consistency + │ + ▼ +Gap Detection + │ + ▼ +REST Backfill + │ + ▼ +Canonical Trade Stream +``` + +Такое разделение обязанностей позволяет развивать каждый уровень независимо, не изменяя уже реализованные компоненты. + +--- + +# Заключение + +Build 060.18 успешно достиг всех поставленных целей. + +В рамках реализации: + +- сформирован самостоятельный слой **Trade Stream Consistency**; +- реализован контракт `TradeStreamConsistencyProtocol`; +- реализован координирующий `TradeStreamConsistencyController`; +- реализована модель состояния `TradeStreamState`; +- введены специализированные исключения предметной области; +- выполнено полное покрытие новой подсистемы unit-тестами; +- подтверждено соответствие реализации утверждённой архитектурной спецификации; +- проведён архитектурный аудит существующей инфраструктуры Trades Feed, результаты которого зафиксированы в настоящем документе. + +Полученные результаты формируют прочную основу для следующего этапа развития подсистемы — построения полноценного **WebSocket Trades Feed**, который станет первым источником непрерывного канонического потока сделок и позволит интегрировать реализованный механизм проверки согласованности в общий конвейер обработки рыночных данных. + +--- + +# Следующий Build + +Build 060.19 — Trade Gap Detection & REST Backfill \ No newline at end of file diff --git a/docs/migrations/build_060_18_architecture.md b/docs/migrations/build_060_18_architecture.md new file mode 100644 index 0000000..8898c45 --- /dev/null +++ b/docs/migrations/build_060_18_architecture.md @@ -0,0 +1,2726 @@ +# 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 возвращают одинаковый объект: + +```python +Trade +``` + +Однако наличие канонической модели ещё не означает существование канонического потока. + +Например, поток может содержать: + +```text +100 +101 +100 +``` + +или + +```text +100 +101 +99 +``` + +или + +```text +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 гарантирует корректность каждой отдельной сделки. + +Однако система ещё не гарантирует корректность последовательности сделок. + +Например: + +```text +Trade +Trade +Trade +``` + +не обязательно образуют корректный поток. + +В частности отсутствуют гарантии: + +- монотонности; +- отсутствия повторов; +- отсутствия конфликтующих дублей; +- целостности опубликованного потока. + +Следовательно, потребитель Trade не может считать получаемый поток достоверным. + +Build 060.18 устраняет именно эту проблему. + +--- + +# Основная идея Build + +Главная идея Build заключается в разделении двух понятий. + +До настоящего момента существовала только каноническая модель Trade. + +После завершения Build появляются две независимые сущности. + +Первая: + +```text +Canonical Trade +``` + +Вторая: + +```text +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 не допускается существование нескольких файлов с одинаковым именем независимо от их расположения в каталогах. + +Единственным исключением является служебный файл: + +```text +__init__.py +``` + +Имя каждого файла должно однозначно отражать его назначение и позволять определить его содержимое без открытия файла. + +Например: + +```text +trade_stream_consistency_controller.py + +trade_stream_state.py + +trade_stream_protocol.py + +trade_stream_exceptions.py +``` + +Не допускается использование общих имён файлов, таких как: + +```text +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 отвечает на вопрос: + +> Может ли данная сделка стать частью уже существующего потока? + +Это принципиально разные уровни проверки. + +Корректная сделка вполне может быть отвергнута контроллером потока. + +Например: + +```text +Trade #100 +Trade #101 +Trade #99 +``` + +Все три объекта являются корректными `Trade`. + +Однако третья сделка нарушает согласованность потока. + +Следовательно, она никогда не становится частью Canonical Trade Stream. + +--- + +# Граница формирования потока + +До прохождения проверки согласованности существует только набор независимых объектов Trade. + +```text +REST + │ +WebSocket + │ + ▼ +Parser + ▼ +Validation + ▼ +Mapper + ▼ +Trade +``` + +После прохождения Stream Consistency появляется новая сущность. + +```text +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 + +```text +trade_id +``` + +Отклонён. + +Причина: + +Архитектура не должна предполагать глобальную уникальность `trade_id`. + +Разные торговые инструменты могут использовать одинаковые идентификаторы. + +--- + +## Вариант 2 + +```text +(symbol, trade_id) +``` + +Принят. + +Данная пара однозначно определяет одну сделку внутри канонического потока. + +--- + +## Вариант 3 + +```text +(symbol, trade_id, source) +``` + +Отклонён. + +После Mapper источник происхождения сделки перестаёт иметь доменное значение. + +REST и WebSocket обязаны порождать одну и ту же каноническую сущность. + +--- + +## Вариант 4 + +```text +(symbol, trade_id, timestamp) +``` + +Отклонён. + +Время исполнения является атрибутом сделки, а не частью её идентичности. + +--- + +# Официальная модель идентичности + +Во всей системе Build 060.18 официальной идентичностью сделки считается: + +```text +(symbol, trade_id) +``` + +Никакие другие поля не участвуют в определении идентичности. + +--- + +# Ordering Model + +## Назначение + +После определения идентичности необходимо определить правило упорядочивания потока. + +--- + +# Рассматриваемые варианты + +--- + +## Ordering по timestamp + +Отклонён. + +Причины: + +- различные источники могут получать данные с различной задержкой; +- время может совпадать; +- транспортная задержка не должна влиять на доменную последовательность. + +--- + +## Ordering по executed_at + +Отклонён. + +`executed_at` является характеристикой сделки, но не механизмом восстановления порядка. + +--- + +## Ordering по (timestamp, trade_id) + +Отклонён. + +Избыточно. + +Усложняет систему без появления дополнительных гарантий. + +--- + +## Ordering по trade_id + +Принят. + +Биржа уже определяет последовательность исполнения сделок. + +Следовательно, система должна использовать именно её. + +--- + +# Официальное правило Ordering + +Для каждого символа поток обязан удовлетворять условию: + +```text +trade_id(new) >= trade_id(last) +``` + +При этом допускаются разрывы последовательности. + +Например: + +```text +100 +101 +103 +120 +``` + +является корректным потоком. + +Build 060.18 не занимается анализом пропущенных идентификаторов. + +Эта задача относится к Build 060.19. + +--- + +# Нарушение порядка + +Если новая сделка имеет меньший `trade_id`, чем последняя опубликованная сделка данного символа, поток считается нарушенным. + +Пример: + +```text +100 +101 +98 +``` + +В этом случае сделка отвергается. + +Контроллер генерирует `TradeOrderingError`. + +Состояние потока при этом не изменяется. + +--- + +# Deduplication Model + +## Назначение + +После проверки порядка необходимо определить правила обработки повторных сделок. + +--- + +# Определение дубликата + +Дубликатом считается сделка, имеющая ту же идентичность: + +```text +(symbol, trade_id) +``` + +что и уже зарегистрированная сделка. + +--- + +# Полный дубликат + +Если все канонические поля совпадают, повтор считается корректным. + +Например: + +```text +BTC +trade_id = 150 +price = 100 +quantity = 5 +``` + +и повтор: + +```text +BTC +trade_id = 150 +price = 100 +quantity = 5 +``` + +представляют одну и ту же сделку. + +Контроллер не публикует её повторно. + +Ошибки не возникает. + +--- + +# Конфликтующий дубликат + +Если идентичность совпадает, но хотя бы одно бизнес-поле отличается: + +- price; +- quantity; +- executed_at; +- aggressor_side; + +поток считается противоречивым. + +Например: + +```text +trade_id = 500 + +price = 100 +``` + +позже: + +```text +trade_id = 500 + +price = 101 +``` + +Такая ситуация невозможна внутри корректного канонического потока. + +Контроллер обязан немедленно завершить обработку ошибкой. + +Публикация сделки запрещается. + +--- + +# Что не считается дубликатом + +Следующие сделки никогда не считаются повтором: + +```text +trade_id = 500 + +trade_id = 501 +``` + +Даже если совпадают: + +- цена; +- объём; +- время исполнения. + +Идентичность определяется исключительно парой: + +```text +(symbol, trade_id) +``` + +--- + +# Итоговая модель обработки сделки + +Каждая поступающая сделка проходит последовательность проверок. + +```text +Trade + │ + ▼ +Определение символа + │ + ▼ +Получение состояния потока + │ + ▼ +Проверка существования trade_id + │ + ├─────────────── Да ───────────────┐ + │ │ + ▼ ▼ +Проверка совпадения Полный дубликат +бизнес-полей │ + │ ▼ + │ Не публиковать + │ + ▼ +Конфликт? + │ + ├──── Да ───► TradeConsistencyError + │ + ▼ +Нет + │ + ▼ +Проверка порядка + │ + ├──── Нарушение ─► TradeOrderingError + │ + ▼ +Регистрация сделки + │ + ▼ +Публикация +``` + +--- + +# Архитектурный результат + +После завершения Build 060.18 в системе появляется новый уровень доменной модели. + +До Build: + +```text +Trade +``` + +После Build: + +```text +Trade + │ + ▼ +Canonical Trade Stream +``` + +Именно канонический поток становится единственным допустимым источником данных для всех последующих компонентов Market Intelligence Pipeline. + +--- + +# Архитектура компонентов + +После определения модели Canonical Trade Stream необходимо определить архитектурные компоненты, обеспечивающие его существование. + +Build 060.18 вводит в Acquisition Layer первую stateful-подсистему. + +До настоящего момента все компоненты Acquisition являлись stateless. + +Появление Stream Consistency является первой точкой, где система начинает хранить собственное состояние. + +Именно поэтому данный Build вводит строго определённые границы ответственности. + +--- + +# Архитектура подсистемы + +Подсистема Stream Consistency состоит из двух компонентов. + +```text +TradeStreamConsistencyController + │ + ▼ + TradeStreamState +``` + +Оба компонента являются внутренними компонентами Acquisition Layer. + +Никакие другие части системы не имеют права изменять их состояние. + +--- + +# TradeStreamConsistencyController + +## Назначение + +TradeStreamConsistencyController является единственной точкой формирования Canonical Trade Stream. + +Никакой другой компонент системы не имеет права принимать решение о публикации сделки. + +Контроллер определяет: + +- может ли сделка стать частью потока; +- нарушает ли она инварианты; +- является ли она повтором; +- должна ли она быть опубликована. + +--- + +# Ответственность Controller + +Контроллер отвечает исключительно за доменные решения. + +К ним относятся: + +- маршрутизация по символам; +- проверка порядка; +- определение повторов; +- обнаружение конфликтующих дублей; +- формирование доменных исключений; +- публикация только корректных сделок. + +--- + +# Controller НЕ отвечает + +Контроллер сознательно не отвечает за: + +- хранение данных; +- структуру дедупликационного окна; +- алгоритм удаления старых записей; +- реализацию FIFO; +- внутреннее устройство состояния. + +Все перечисленные задачи принадлежат исключительно TradeStreamState. + +--- + +# Главный принцип Controller + +Controller принимает решения. + +State хранит данные. + +Это фундаментальное архитектурное правило Build 060.18. + +--- + +# TradeStreamState + +## Назначение + +TradeStreamState представляет собой внутреннюю модель состояния одного потока сделок. + +Каждый экземпляр состояния соответствует ровно одному торговому символу. + +Например: + +```text +BTCUSDT + │ + ▼ +TradeStreamState +``` + +или + +```text +ETHUSDT + │ + ▼ +TradeStreamState +``` + +Состояние различных символов никогда не смешивается. + +--- + +# Почему состояние существует отдельно + +Во время проектирования рассматривалась возможность хранения всех данных непосредственно внутри Controller. + +Этот вариант был отклонён. + +Причина: + +Controller должен выражать бизнес-правила. + +State должен выражать состояние предметной области. + +Разделение этих двух ролей существенно упрощает дальнейшее развитие системы. + +--- + +# Инварианты TradeStreamState + +Каждый экземпляр состояния обязан удовлетворять следующим требованиям. + +--- + +## Инвариант №1 + +Состояние обслуживает ровно один symbol. + +--- + +## Инвариант №2 + +Состояние никогда не принимает Trade другого symbol. + +Даже если Controller ошибётся. + +Это является дополнительной защитой целостности системы. + +--- + +## Инвариант №3 + +Состояние хранит только уже опубликованные сделки. + +Непринятые сделки никогда не изменяют состояние. + +--- + +## Инвариант №4 + +Состояние никогда самостоятельно не принимает бизнес-решения. + +Оно лишь предоставляет операции хранения. + +--- + +## Инвариант №5 + +Состояние полностью скрывает собственную реализацию. + +Controller не знает: + +- используется ли OrderedDict; +- используется ли deque; +- используется ли Ring Buffer; +- используется ли специализированная структура данных. + +--- + +# Внутреннее состояние + +Минимальная модель состояния включает: + +```text +symbol + +last_trade_id + +recent_trades +``` + +Этого достаточно для реализации всех инвариантов Build 060.18. + +--- + +# symbol + +Каждое состояние хранит собственный символ. + +Например: + +```text +BTCUSDT +``` + +Это необходимо по двум причинам. + +Во-первых, + +для проверки принадлежности входящей сделки. + +Во-вторых, + +для полноценной диагностики ошибок. + +--- + +# last_trade_id + +last_trade_id означает: + +> максимальный trade_id, +> успешно опубликованный системой +> для данного символа. + +Это очень важно. + +last_trade_id НЕ означает: + +> последний увиденный trade_id. + +Например: + +```text +100 + +101 + +95 +``` + +После возникновения OrderingError состояние остаётся: + +```text +last_trade_id = 101 +``` + +Никаких изменений не происходит. + +--- + +# recent_trades + +recent_trades представляет собой окно уже опубликованных сделок. + +Окно необходимо исключительно для проверки повторов. + +После выхода сделки из окна + +она перестаёт участвовать в дедупликации. + +--- + +# Почему используется окно + +Полная история сделок потенциально бесконечна. + +Следовательно, + +невозможно хранить все Trade. + +Используется ограниченное окно фиксированного размера. + +Размер окна определяется конфигурацией. + +--- + +# Требования к окну + +Окно обязано обеспечивать: + +поиск + +```text +O(1) +``` + +вставку + +```text +O(1) +``` + +удаление самой старой записи + +```text +O(1) +``` + +Эти требования являются обязательными. + +--- + +# FIFO Window + +Во время проектирования рассматривались различные структуры хранения. + +--- + +## Полная история + +Отклонена. + +Причина: + +неограниченный рост памяти. + +--- + +## HashSet + +Отклонён. + +Причина: + +невозможно проверить конфликтующий повтор. + +--- + +## LRU Cache + +Отклонён. + +Причина: + +LRU ориентирован на обращения. + +Trade Stream ориентирован на порядок появления. + +--- + +## FIFO Window + +Принят. + +FIFO полностью соответствует природе потока сделок. + +Самые старые сделки постепенно забываются. + +Новые добавляются в конец окна. + +--- + +# Внутренний ключ + +Поскольку один экземпляр TradeStreamState обслуживает только один symbol, + +ключом окна становится исключительно: + +```text +trade_id +``` + +Полная идентичность + +```text +(symbol, trade_id) +``` + +используется только на уровне Controller. + +Это уменьшает объём памяти + +и упрощает внутреннюю структуру состояния. + +--- + +# Жизненный цикл состояния + +Build 060.18 определяет простой жизненный цикл. + +--- + +## Создание + +Состояние создаётся лениво. + +Первое появление сделки данного символа приводит к созданию нового TradeStreamState. + +--- + +## Использование + +После создания состояние используется всеми последующими сделками данного символа. + +--- + +## Уничтожение + +Build 060.18 не удаляет состояния. + +Они существуют до завершения процесса. + +Управление жизненным циклом относится к Runtime Layer и будет рассматриваться отдельно. + +--- + +# Публичный контракт Controller + +Controller предоставляет единственную публичную операцию. + +```python +accept(trade: Trade) -> Trade | None +``` + +Других публичных методов Build 060.18 не вводит. + +--- + +# Семантика accept() + +Если сделка успешно прошла проверку, + +Controller возвращает исходный объект Trade. + +Если поступил полный повтор, + +возвращается: + +```python +None +``` + +Если обнаружено нарушение инвариантов, + +генерируется соответствующее исключение. + +--- + +# Почему возвращается Trade + +Во время проектирования рассматривались альтернативы. + +--- + +## bool + +Отклонён. + +Причина: + +вызывающий код вынужден хранить исходный объект отдельно. + +--- + +## Result Object + +Отклонён. + +Причина: + +избыточен для Build 060.18. + +--- + +## Исключение для дубликатов + +Отклонено. + +Повтор является нормальной ситуацией. + +Он не считается ошибкой. + +--- + +## Trade | None + +Принят. + +Контракт минимален, + +естественен + +и легко расширяется в будущем. + +--- + +# Исключения Controller + +Build вводит только два новых доменных исключения. + +--- + +## TradeOrderingError + +Возникает, + +если сделка нарушает монотонность потока. + +Например: + +```text +100 + +101 + +95 +``` + +--- + +## TradeConsistencyError + +Возникает, + +если найден конфликтующий повтор. + +Например: + +```text +trade_id = 500 + +price = 100 +``` + +позже + +```text +trade_id = 500 + +price = 101 +``` + +Такой поток считается внутренне противоречивым. + +--- + +# Взаимодействие Controller и State + +Важнейшим архитектурным принципом Build является инкапсуляция состояния. + +Controller никогда не обращается к внутренним структурам данных напрямую. + +Вместо этого он взаимодействует со State исключительно через его операции. + +Концептуально взаимодействие выглядит следующим образом. + +```text +Controller + + │ + + ▼ + +TradeStreamState + + │ + + ├── получить последнюю сделку + + ├── получить зарегистрированную сделку + + ├── зарегистрировать новую сделку + + └── поддерживать размер окна +``` + +Таким образом: + +- Controller ничего не знает о реализации хранения; +- State ничего не знает о бизнес-правилах проверки согласованности. + +Именно это разделение делает подсистему устойчивой к дальнейшему развитию. + +--- + +# Алгоритм работы TradeStreamConsistencyController + +После определения архитектуры компонентов необходимо формально определить алгоритм обработки каждой сделки. + +Build 060.18 рассматривает Controller как **детерминированный автомат** (Deterministic State Machine). + +Это означает, что результат обработки полностью определяется двумя величинами: + +```text +Текущее состояние + ++ + +Входящая Trade + +↓ + +Новое состояние + ++ + +Результат обработки +``` + +Никакие внешние факторы не влияют на принятие решения. + +--- + +# Почему выбран детерминированный автомат + +Во время проектирования рассматривались несколько моделей. + +--- + +## Процедурный алгоритм + +Обычная последовательность условий. + +```text +if + +if + +if + +if +``` + +Работает. + +Но по мере развития Build начинает быстро усложняться. + +--- + +## Таблица правил + +Возможна. + +Однако становится плохо читаемой. + +--- + +## Deterministic State Machine + +Принята. + +Причины: + +- полностью предсказуемое поведение; +- простое тестирование; +- возможность восстановления состояния; +- одинаковое поведение REST и WebSocket; +- естественное расширение для Recovery Build. + +--- + +# Формальная модель + +Для каждой входящей сделки существует единственный возможный результат. + +```text +TradeStreamState + ++ + +Trade + +↓ + +TradeStreamState' + ++ + +Result +``` + +где + +Result представляет собой одно из следующих состояний: + +```text +Accepted + +Duplicate + +TradeOrderingError + +TradeConsistencyError +``` + +Других исходов Build 060.18 не предусматривает. + +--- + +# Последовательность обработки + +Каждая входящая сделка проходит одинаковую последовательность шагов. + +```text +Trade + +↓ + +Получить состояние символа + +↓ + +Поиск trade_id + +↓ + +Определение дубликата + +↓ + +Проверка порядка + +↓ + +Регистрация + +↓ + +Публикация +``` + +Порядок этапов является обязательным. + +Изменение последовательности проверок запрещается. + +--- + +# Почему сначала проверяется дубликат + +Во время проектирования рассматривался альтернативный вариант. + +Сначала Ordering. + +Потом Deduplication. + +Например: + +```text +100 + +101 + +100 +``` + +Ordering немедленно сообщил бы, + +что сделка старая. + +Но на самом деле это корректный повтор. + +Следовательно, + +Ordering не должен выполняться первым. + +--- + +# Официальная последовательность + +Контроллер всегда выполняет проверки в следующем порядке. + +--- + +## Шаг 1 + +Получение состояния символа. + +Если состояние отсутствует, + +оно создаётся. + +--- + +## Шаг 2 + +Поиск существующей сделки. + +Если trade_id найден, + +Controller обязан сравнить всю каноническую модель. + +--- + +## Шаг 3 + +Определение типа повтора. + +Если совпадают все бизнес-поля, + +сделка считается полным дубликатом. + +Возвращается: + +```python +None +``` + +--- + +Если найдено хотя бы одно различие, + +генерируется + +```text +TradeConsistencyError +``` + +После этого обработка завершается. + +--- + +## Шаг 4 + +Если сделка новая, + +выполняется проверка порядка. + +Если + +```text +trade_id < last_trade_id +``` + +генерируется + +```text +TradeOrderingError +``` + +Состояние не изменяется. + +--- + +## Шаг 5 + +Если все проверки успешно завершены, + +сделка регистрируется. + +--- + +## Шаг 6 + +После регистрации + +Trade публикуется вызывающему компоненту. + +--- + +# Почему регистрация выполняется перед публикацией + +Это принципиальное решение Build. + +Рассмотрим последовательность. + +Плохой вариант. + +```text +Trade + +↓ + +Publication + +↓ + +Registration +``` + +Если между публикацией и регистрацией произойдёт исключение, + +система потеряет согласованность. + +Потребитель уже увидел сделку, + +а внутреннее состояние ещё нет. + +--- + +Правильный вариант. + +```text +Trade + +↓ + +Registration + +↓ + +Publication +``` + +После регистрации + +Controller уже находится в согласованном состоянии. + +--- + +# Неизменяемость состояния при ошибках + +Любая ошибка обязана обладать свойством атомарности. + +Если возникает: + +- TradeOrderingError; +- TradeConsistencyError; + +состояние Controller не изменяется. + +Это означает: + +- last_trade_id остаётся прежним; +- окно дедупликации не изменяется; +- история публикаций остаётся неизменной. + +--- + +# Архитектурная диаграмма + +Полный путь обработки сделки выглядит следующим образом. + +```text +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: + +```text +Handler + +↓ + +Consumer +``` + +После Build: + +```text +Handler + +↓ + +TradeStreamConsistencyController + +↓ + +Consumer +``` + +Никакие существующие компоненты не меняют собственную ответственность. + +--- + +# Integration Boundary + +TradeStreamConsistencyController располагается между Handler и Consumer. + +Он получает уже канонический Trade. + +И публикует только Canonical Trade Stream. + +Это официальная граница между: + +```text +Canonical Trade + +↓ + +Canonical Trade Stream +``` + +--- + +# Dependency Injection + +Controller создаётся в Composition Root. + +Никакой компонент не создаёт Controller самостоятельно. + +Итоговая схема выглядит следующим образом. + +```text +Composition Root + + │ + + ├──────────────┐ + + ▼ ▼ + +TradesHandler TradeStreamConsistencyController + + │ │ + + └──────┬───────┘ + + ▼ + + TradesFeed +``` + +Таким образом: + +- Feed ничего не знает о реализации Controller; +- Controller ничего не знает о Feed; +- Handler ничего не знает о Controller. + +Каждый компонент получает только необходимые зависимости. + +--- + +# План изменений файлов + +Build 060.18 вводит новую внутреннюю подсистему. + +Предлагаемая структура. + +```text +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 + +Содержит: + +```text +TradeStreamConsistencyController +``` + +Отвечает исключительно за доменные решения. + +--- + +## trade_stream_state.py + +Содержит: + +```text +TradeStreamState +``` + +Отвечает исключительно за хранение состояния. + +--- + +## trade_stream_exceptions.py + +Содержит: + +```text +TradeOrderingError + +TradeConsistencyError +``` + +--- + +## trade_stream_protocol.py + +Содержит публичный контракт подсистемы Stream Consistency. + +Все внешние компоненты должны зависеть от данного Protocol, а не от конкретной реализации Controller. + +--- + +# Соглашение об именовании файлов + +Build 060.18 закрепляет правило уникальности имён файлов в проекте Dzentra. + +Все новые файлы должны иметь уникальные имена в пределах всего репозитория. + +Имя файла должно отражать его содержимое, а не только роль внутри текущего каталога. + +Например: + +```text +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 + +**Решение** + +Идентичность сделки определяется парой: + +```text +(symbol, trade_id) +``` + +Другие поля в идентичность не входят. + +--- + +## ADR-060.18-02 — Stream State Ownership + +**Решение** + +Единственным владельцем состояния потока является `TradeStreamConsistencyController`, который управляет набором `TradeStreamState` по символам. + +--- + +## ADR-060.18-03 — Deterministic Stream State Machine + +**Решение** + +Контроллер рассматривается как детерминированный автомат: + +```text +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. + +```text +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 + +```text +100 + +101 + +102 +``` + +### Ожидаемый результат + +Все сделки публикуются. + +--- + +## 6. Разрыв последовательности + +```text +100 + +105 + +150 +``` + +### Ожидаемый результат + +Ошибок нет. + +--- + +## 7. Уменьшение trade_id + +```text +100 + +101 + +95 +``` + +### Проверяется + +Ordering. + +### Ожидаемый результат + +TradeOrderingError. + +--- + +## 8. Полный дубликат + +```text +100 + +100 +``` + +### Проверяется + +Deduplication. + +### Ожидаемый результат + +Возвращается + +```python +None +``` + +--- + +## 9. Конфликтующий повтор + +```text +trade_id = 500 + +price = 100 +``` + +позже + +```text +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. + +Однако контроллер обязан удовлетворять следующим требованиям. + +--- + +## Поиск сделки + +Средняя сложность + +```text +O(1) +``` + +--- + +## Добавление сделки + +Средняя сложность + +```text +O(1) +``` + +--- + +## Удаление старой записи + +Средняя сложность + +```text +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. + +--- + +## Производительность + +Все операции соответствуют требованиям: + +```text +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. \ No newline at end of file