Gonka GitHub Discussions · Discussion #954

Оптимистичное параллельное выполнение сообщений для inference-chain

Отдельная страница дискуссии с русским переводом и параллельным режимом RU / Original с точным сопоставлением предложений.

Как пользоваться: включите RU / Original и кликните по предложению — соответствующее предложение в другой колонке доскроллится и подсветится.
Discussion #954

Оптимистичное параллельное выполнение сообщений для inference-chain

Оригинал: Optimistic parallel execution of messages for inference-chain

akup avatar
akupMaintainerАвтор
2026-03-26

Двухуровневый кэш блоков/передач с обнаружением конфликтов для OCC

Мотивация

Cosmos SDK по умолчанию обрабатывает транзакции последовательно. Каждое чтение из хранилища KV требует демаршаллинга protobuf, а каждая запись требует маршалинга — операций, которые одновременно нагружают ЦП и повторяются между транзакциями в одном и том же блоке.

Цель этого предложения двоякая:

  • Этап 1 (текущий): внедрение двухуровневого кэша блоков/передач, реализованного OptimisticStore (название сохранено для совместимости кода), чтобы исключить избыточную маршализацию/демаршализацию внутри блока, сохраняя при этом полную детерминированность и совместимость с существующей моделью последовательного выполнения.

Этап 1 (текущий): внедрение двухуровневого кэша блоков/передач, реализованного OptimisticStore (название сохранено для совместимости кода), чтобы исключить избыточную маршализацию/демаршализацию внутри блока, сохраняя при этом полную детерминированность и совместимость с существующей моделью последовательного выполнения.

  • Этап 2 (в будущем). Повторное использование той же инфраструктуры отслеживания конфликтов для включения Optimistic Concurrency Control (OCC) — параллельного выполнения транзакций с автоматическим обнаружением конфликтов и откатом.

Этап 2 (в будущем). Повторное использование той же инфраструктуры отслеживания конфликтов для включения Optimistic Concurrency Control (OCC) — параллельного выполнения транзакций с автоматическим обнаружением конфликтов и откатом.

Постановка задачи

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

Дополнительный оптимистический режим выполнения Cosmos SDK не является универсальной и универсальной функцией. Его реализация сильно зависит от шаблонов доступа к модулю. Поэтому Gonka реализует собственную схему OCC, адаптированную к рабочим нагрузкам цепочки вывода.

Дизайн: OCC для Gonka

Планирование

Планировщик принимает N сообщений из мемпула (где N масштабируется в зависимости от ядер ЦП). Сообщения сортируются и группируются по сходству — похожие сообщения обычно имеют доступ к одинаковым ключам хранилища и имеют сопоставимое время выполнения.

Выбранные N сообщений выполняются как параллельный пакет. Планировщик ожидает завершения всех сообщений в пакете, а затем принимает следующий пакет.

Обнаружение конфликтов

Во время выполнения наборы чтения и записи каждой транзакции записываются конфликтом Tracker внутри каждого OptimisticStore . После завершения пакета:

  • Конфликт чтения-записи: транзакция A прочитала ключ K, транзакция B записала ключ K → A необходимо откатить и перепланировать.
  • Конфликт записи-записи: транзакции A и B записали ключ K → все, кроме одной (той, что была раньше в порядке блоков), необходимо откатить и перепланировать.

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

Порядок разрешения конфликтов

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

Группировка по сходству

Критерием группировки является прогнозируемый шаблон доступа типа сообщения. Два MsgValidation для одной и той же модели и эпохи «похожи» — они обращаются к одним и тем же ключам и выполняются примерно одинаково. Статический анализ каждого типа сообщения обеспечивает прогнозирование доступа. Большинство горячих сообщений исчезнут, когда шардчейны начнут работать, но в любом случае этот подход продолжает работать и будет полезен.

Газ для переоформленных сделок

Когда транзакция откатывается из-за конфликта, счетчик газа сбрасывается. Перенесенное выполнение — это новая попытка с новым счетчиком. Таким образом, учет газа Cosmos остается правильным — лимит газа пользователя применяется к успешному выполнению, а не к неудачным спекулятивным попыткам.

Высокая конфликтность

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

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

Улучшение пропускной способности

Реализация этой схемы OCC может повысить пропускную способность основной цепи в N × k раз, где:

N = количество ядер ЦП

k = коэффициент эффективности (< 1,0, зависит от конфликтности)

Для рабочих нагрузок с незначительным перекрытием (например, выводы, касающиеся разных моделей), k приближается к 1,0, а пропускная способность масштабируется почти линейно в зависимости от ядер.

Этап 1: двухуровневый кэш блоков/передач (текущая реализация)

Фаза 1 полностью реализована и объединена. Параллельного планировщика пока нет. Все транзакции по-прежнему выполняются последовательно. Это не оптимистическое исполнение само по себе; это детерминированный двухуровневый кеш (tx-черновик + блочный кеш) с обнаружением конфликтов, который можно использовать позже в оптимистичном параллельном режиме.

Этот кэш предоставляет следующие преимущества:

  • Устраните повторную маршализацию/демаршализацию для значений хранилища, к которым обращаются несколько раз в блоке (например, Параметры, EpochGroupData).
  • Заложите основу для этапа 2, отслеживая наборы операций чтения/записи для каждой транзакции, чтобы обнаружение конфликтов можно было включить с помощью одной переменной среды ( COSMOS_OCC_ENABLED=1 ).

Общие системные данные и область кэша

Мы используем этот двухуровневый кэш блоков/передач для часто читаемых общих системных данных:

Параметры

EpochGroupData

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

Протобуф Маршал/Унмаршал и Газ

Cosmos SDK маршалирует и демаршалирует данные protobuf при операциях чтения и записи в хранилище, в том числе в рамках учета газа в хранилище. Когда операции чтения/записи выполняются из этого кэша, повторяющиеся операции кодека пропускаются, и на эти кэшированные операции в этих системных путях не тратится дополнительный газ.

В нашем проекте это приемлемо, поскольку мы уже используем политики отсутствия газа или фиксированного газа для системных сообщений (чтобы снизить риск DDoS-спама, сохраняя при этом детерминированное выполнение). Например, чтение EpochGroupData может использоваться логикой защиты конечных точек для идентификации участников, сокращенных из-за отсутствия выводов или отсутствия cPoC; это системная проверка, при которой быстрое выполнение предпочтительнее повторной маршаллинга/демаршаллинга.

Архитектура

┌──────────────────────────────────────────────┐ │ OptimisticStore │ │ │ │ Уровень 1: черновой вариант передачи (с учетом контекста, для каждой передачи) │ │ ↓ промах │ │ Уровень 2: Кэш блоков (в памяти, поблочно) │ │ ↓ промах │ │ Уровень 3: Серверная часть хранилища (постоянный KV/protobuf) │ │ │ │ + конфликт-трекер (наборы чтения/записи для каждой передачи) │ └───────────────────────────────────────────────────┘
  • Tx Draft: буфер отложенной записи, ограниченный одной транзакцией через context.Value . Создан в AnteHandler, зафиксирован для блокировки кеша при успешной передаче в PostHandler, отброшен в случае сбоя.
  • Кэш блоков: общая карта в памяти для текущего блока. Заполняется при первом чтении (шаблон кэширования). Сбрасывается в серверную часть постоянного хранилища в EndBlock .
  • Серверная часть хранилища: подключаемый интерфейс (загрузка/сохранение/удаление/клонирование), который оборачивает любое постоянное хранилище (карту коллекций, необработанный KV и т. д.).

Жизненный цикл: CheckTx против DeliverTx

CheckTx (проверка мемпула)

Во время CheckTx SDK проверяет транзакцию перед ее попаданием в мемпул. Оптимистичный магазин не фиксирует черновики во время CheckTx:

PostHandler: если ctx. IsCheckTx () || симулировать { return // пропустить фиксацию — только проверка мемпула }

Это означает, что CheckTx видит состояние постоянного хранилища (или блокирует кеш, если он уже теплый), но его записи отбрасываются. Это правильно, поскольку CheckTx не должен иметь побочных эффектов.

DeliverTx (выполнение блока)

Во время DeliverTx выполняется полный жизненный цикл:

AnteHandler (OptimisticStoreDraftDecorator): → StoreGroup.WithDraftAll(ctx) // прикрепляем черновики для каждой передачи → StoreGroup.RegisterTxAll(ctx) // регистрируемся для отслеживания OCC Выполнение сообщения: → store.Get(ctx, key) // читает: черновик → кэш блоков → серверная часть → store.Set(ctx, key) // пишет: перейти к черновику PostHandler (OptimisticStoreCommitPostDecorator): в случае успеха: → StoreGroup.CommitDraftAll(ctx) // объединяем черновик → блокируем кеш else: → StoreGroup.ReleaseDraftAll(ctx) // отбрасываем черновик EndBlock: → StoreGroup.FlushAll(ctx) // блокируем кеш → постоянное хранилище BeginBlock (следующий блок): → StoreGroup.InvalidateAll() // очищаем кеш блоков + средство отслеживания конфликтов

Отраслевые шашки

Для операций, требующих спекулятивного выполнения в пределах одной передачи (например. Шаблоны CacheContext), OptimisticStore поддерживает черновики ветвей:

кэшCtx, writeCache:= хранитель. CacheContext ( ctx ) // спекулятивная работа над кэшем Ctx... writeCache () // объединяем черновик ветки с черновиком родительского tx + кеш SDK

Метод StoreGroup.CacheContext создает черновики ветвей для каждого зарегистрированного магазина и возвращает функцию объединенной фиксации.

Как OptimisticStore работает с существующим хранилищем

OptimisticStore — это декоратор — он оборачивает существующее хранилище, не заменяя его. Интерфейс StoreBackend делает это явным образом:

type StoreBackend [ K сравнимый , V любой ] struct { Load func ( ctx context. Контекст, клавиша K) (V, bool) Сохранить функцию (ctx context. Контекст, ключ K, val V) Удалить func (ctx context. Контекст, ключ K) Функция клонирования (val V) V }

Любое существующее хранилище Collections.Map, Collections.Item или необработанное хранилище KV можно обернуть, предоставив эти четыре функции.

Пример: упаковка карты коллекций

Для EpochGroupData, который хранится в коллекции.Map[Pair[uint64,string], EpochGroupData]:

// В keeper.go NewKeeper: epochGroupStore : NewOptimisticCollMap [ epochGroupCacheKey , Collections. Пара [ uint64 , string ], типы. EpochGroupData , ]( sb , типы . EpochGroupDataPrefix, «epoch_group_data», коллекции. PairKeyCodec (коллекции . Uint64Key, коллекции. StringKey), кодек. CollValue [типы. EpochGroupData ]( cdc ),cacheConfig , func (ключ epochGroupCacheKey) коллекции. Pair [uint64, string] {возврат коллекций. Присоединяйтесь (ключ. Эпоха, ключ. ModelId ) }, ),

NewOptimisticCollMap создает коллекции.Map И обертывает их OptimisticStore за один вызов. Функции загрузки/сохранения/удаления делегируются коллекции; Clone использует proto.Clone.

Пример: упаковка синглтона (Params)

Для Params хранится как один объект protobuf с фиксированным ключом KV:

paramsStore: NewOptimisticProtoItem [types. Параметры](storeService, cdc, типы. ParamsKey, "params",cacheConfig, ),

NewOptimisticProtoItem обрабатывает внутреннюю маршализацию/демаршализацию. GetParams и SetParams просто вызывают paramsStore.Get(ctx)/paramsStore.Set(ctx, val) .

Регистрация в StoreGroup

Каждый оптимистичный магазин должен быть зарегистрирован в StoreGroup хранителя, чтобы методы жизненного цикла ( InvalidateAll , FlushAll , WithDraftAll и т. д.) применялись ко всем магазинам единообразно:

к. группа магазинов. Зарегистрируйтесь (k. epochGroupStore. OptimisticStore) k. группа магазинов. Зарегистрируйтесь (k. paramsStore. Магазин ())

Добавление нового оптимистичного магазина в Keeper

Чтобы обернуть новую коллекцию оптимистическим кэшированием:

Определите тип ключа кэша (должен быть сопоставимым): введите myNewCacheKey struct { Field1 uint64 Field2 string }

Определите тип ключа кэша (должен быть сопоставимым):

введите myNewCacheKey struct { Field1 uint64 Field2 string }

  • Добавьте в Keeper поле оптимистического магазина: myNewStore * OptimisticCollMap [myNewCacheKey, Collections. Пара [ uint64 , string ], типы. МойНовыйТип]

Добавьте в Keeper поле оптимистичного магазина:

myNewStore * OptimisticCollMap [myNewCacheKey, Collections. Пара [ uint64 , string ], типы. МойНовыйТип]
  • Инициализируйте в NewKeeper, используя NewOptimisticCollMap (или NewOptimisticProtoItem для одиночных элементов).

Инициализируйте в NewKeeper, используя NewOptimisticCollMap (или NewOptimisticProtoItem для одиночных элементов).

  • Зарегистрируйтесь в группе магазинов: k . группа магазинов. Зарегистрируйтесь (k.myNewStore. ОптимистическийМагазин)

Зарегистрируйтесь в группе магазина:

к. группа магазинов. Зарегистрируйтесь (k.myNewStore. ОптимистическийМагазин)
  • Напишите в Keeper методы получения/установки, которые делегируют хранилище: func ( k Keeper ) GetMyNewData ( ctx context. Контекст, ключ myNewCacheKey) (types. MyNewType, bool) { return k. мойНовыйМагазин. Get (ctx, key) } func (k Keeper) SetMyNewData (ctx context. Контекст, типы значений. MyNewType) {ключ:= myNewCacheKey {Field1: val. Поле1, Поле2: значение. Поле2 } k . мойНовыйМагазин. Установить (ctx, ключ, val)}

Напишите методы получения/установки в Keeper, которые делегируют хранилище:

func (k Keeper) GetMyNewData (ctx context. Контекст, ключ myNewCacheKey) (types. MyNewType, bool) { return k. мойНовыйМагазин. Get (ctx, key) } func (k Keeper) SetMyNewData (ctx context. Контекст, типы значений. MyNewType) {ключ:= myNewCacheKey {Field1: val. Поле1, Поле2: значение. Поле2 } k . мойНовыйМагазин. Установить (ctx, ключ, val)}

Никаких изменений в AnteHandler, PostHandler, BeginBlock или EndBlock не требуется — StoreGroup автоматически обрабатывает все зарегистрированные хранилища.

Подключение: AnteHandler, PostHandler, BeginBlock, EndBlock

АнтеХандлер ( ante.go )

тип OptimisticStoreDraftDecorator struct { InferenceKeeper * inferencemodulekeeper. Keeper } func (d OptimisticStoreDraftDecorator) AnteHandle (ctx sdk. Контекст, передача SDK. Tx, симулируем bool, следующий SDK. AnteHandler) (sdk. Контекст, ошибка) { g := d . Хранитель выводов. StoreGroup () newCtx:= ctx. Сконтекстом (g. WithDraftAll (ctx. Контекст ())) г . RegisterTxAll (newCtx) return next (newCtx, tx, имитация) }

Постхандлер ( ante.go )

func (d OptimisticStoreCommitPostDecorator) PostHandle (ctx sdk. Контекст, передача SDK. Tx, симуляция, успех bool, следующий SDK. PostHandler) (sdk. Контекст, ошибка) {newCtx, err:= next (ctx, tx, имитация, успех), если ctx. IsCheckTx () || имитация { return newCtx, nil // никаких побочных эффектов во время CheckTx } g := d . Хранитель выводов. StoreGroup () если успех {g. CommitDraftAll (newCtx) // объединяем черновик → кэш блока } else { g . ReleaseDraftAll (newCtx) // отбрасываем черновик } return newCtx, nil }

БегинБлок (модуль.go)

func (am AppModule) BeginBlock (ctx context. Контекст) ошибка {am. хранитель. Группа магазинов (). InvalidateAll () // очистка кеша блоков + трекер конфликтов // ... оставшаяся часть начального блока }

EndBlock (модуль.go)

func (am AppModule) EndBlock (ctx context. Контекст) ошибка {отложить утра. хранитель. Группа магазинов (). FlushAll ( ctx ) // сохранение кеша блоков для хранения // ... остальная часть конечного блока }

Дорожная карта этапа 2: Планировщик OCC

На втором этапе будет добавлен детерминированный параллельный планировщик. Текущий конфликттрекер внутри OptimisticStore уже отслеживает наборы чтения/записи для каждой передачи и может обнаруживать конфликты с помощью DetectConflicts() . Что осталось:

  • Пакетный планировщик — выберите N сообщений, сгруппируйте их по прогнозируемому шаблону доступа, выполняйте параллельно.
  • Разрешение конфликтов — после завершения пакета вызовите DetectConflicts() в каждом магазине, откатите проигравших (по порядку блоков), перепланируйте следующий пакет.
  • Retry Budget — принудительное последовательное выполнение после N неудачных попыток.
  • Интеграция — замена последовательного цикла DeliverTx пакетным планировщиком; сохраните тот же жизненный цикл проекта AnteHandler/PostHandler.

API OptimisticStore уже предназначен для этого:

// Включить обнаружение конфликтов при запуске: // COSMOS_OCC_ENABLED=1 // После завершения параллельного пакета:conflictedReads,conflictedWrites:= store. DetectConflicts () // Откат и повторное планирование конфликтующих txID... store . СброситьКонфликтТрекер ()

На этапе 2 не ожидается никаких изменений в уровне хранилища, хранителя или кэша — только добавление планировщика и замена цикла DeliverTx.

tcharchian avatar
tcharchianMaintainerMaintainer
2026-03-26
Русский перевод
akup avatar
akupMaintainerАвтор
2026-03-26

Двухуровневый кэш блоков/передач с обнаружением конфликтов для OCC

Мотивация

Cosmos SDK по умолчанию обрабатывает транзакции последовательно. Каждое чтение из хранилища KV требует демаршаллинга protobuf, а каждая запись требует маршалинга — операций, которые одновременно нагружают ЦП и повторяются между транзакциями в одном и том же блоке.

Цель этого предложения двоякая:

  • Этап 1 (текущий): внедрение двухуровневого кэша блоков/передач, реализованного OptimisticStore (название сохранено для совместимости кода), чтобы исключить избыточную маршализацию/демаршализацию внутри блока, сохраняя при этом полную детерминированность и совместимость с существующей моделью последовательного выполнения.

Этап 1 (текущий): внедрение двухуровневого кэша блоков/передач, реализованного OptimisticStore (название сохранено для совместимости кода), чтобы исключить избыточную маршализацию/демаршализацию внутри блока, сохраняя при этом полную детерминированность и совместимость с существующей моделью последовательного выполнения.

  • Этап 2 (в будущем). Повторное использование той же инфраструктуры отслеживания конфликтов для включения Optimistic Concurrency Control (OCC) — параллельного выполнения транзакций с автоматическим обнаружением конфликтов и откатом.

Этап 2 (в будущем). Повторное использование той же инфраструктуры отслеживания конфликтов для включения Optimistic Concurrency Control (OCC) — параллельного выполнения транзакций с автоматическим обнаружением конфликтов и откатом.

Постановка задачи

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

Дополнительный оптимистический режим выполнения Cosmos SDK не является универсальной и универсальной функцией. Его реализация сильно зависит от шаблонов доступа к модулю. Поэтому Gonka реализует собственную схему OCC, адаптированную к рабочим нагрузкам цепочки вывода.

Дизайн: OCC для Gonka

Планирование

Планировщик принимает N сообщений из мемпула (где N масштабируется в зависимости от ядер ЦП). Сообщения сортируются и группируются по сходству — похожие сообщения обычно имеют доступ к одинаковым ключам хранилища и имеют сопоставимое время выполнения.

Выбранные N сообщений выполняются как параллельный пакет. Планировщик ожидает завершения всех сообщений в пакете, а затем принимает следующий пакет.

Обнаружение конфликтов

Во время выполнения наборы чтения и записи каждой транзакции записываются конфликтом Tracker внутри каждого OptimisticStore . После завершения пакета:

  • Конфликт чтения-записи: транзакция A прочитала ключ K, транзакция B записала ключ K → A необходимо откатить и перепланировать.
  • Конфликт записи-записи: транзакции A и B записали ключ K → все, кроме одной (той, что была раньше в порядке блоков), необходимо откатить и перепланировать.

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

Порядок разрешения конфликтов

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

Группировка по сходству

Критерием группировки является прогнозируемый шаблон доступа типа сообщения. Два MsgValidation для одной и той же модели и эпохи «похожи» — они обращаются к одним и тем же ключам и выполняются примерно одинаково. Статический анализ каждого типа сообщения обеспечивает прогнозирование доступа. Большинство горячих сообщений исчезнут, когда шардчейны начнут работать, но в любом случае этот подход продолжает работать и будет полезен.

Газ для переоформленных сделок

Когда транзакция откатывается из-за конфликта, счетчик газа сбрасывается. Перенесенное выполнение — это новая попытка с новым счетчиком. Таким образом, учет газа Cosmos остается правильным — лимит газа пользователя применяется к успешному выполнению, а не к неудачным спекулятивным попыткам.

Высокая конфликтность

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

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

Улучшение пропускной способности

Реализация этой схемы OCC может повысить пропускную способность основной цепи в N × k раз, где:

N = количество ядер ЦП

k = коэффициент эффективности (< 1,0, зависит от конфликтности)

Для рабочих нагрузок с незначительным перекрытием (например, выводы, касающиеся разных моделей), k приближается к 1,0, а пропускная способность масштабируется почти линейно в зависимости от ядер.

Этап 1: двухуровневый кэш блоков/передач (текущая реализация)

Фаза 1 полностью реализована и объединена. Параллельного планировщика пока нет. Все транзакции по-прежнему выполняются последовательно. Это не оптимистическое исполнение само по себе; это детерминированный двухуровневый кеш (tx-черновик + блочный кеш) с обнаружением конфликтов, который можно использовать позже в оптимистичном параллельном режиме.

Этот кэш предоставляет следующие преимущества:

  • Устраните повторную маршализацию/демаршализацию для значений хранилища, к которым обращаются несколько раз в блоке (например, Параметры, EpochGroupData).
  • Заложите основу для этапа 2, отслеживая наборы операций чтения/записи для каждой транзакции, чтобы обнаружение конфликтов можно было включить с помощью одной переменной среды ( COSMOS_OCC_ENABLED=1 ).

Общие системные данные и область кэша

Мы используем этот двухуровневый кэш блоков/передач для часто читаемых общих системных данных:

Параметры

EpochGroupData

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

Протобуф Маршал/Унмаршал и Газ

Cosmos SDK маршалирует и демаршалирует данные protobuf при операциях чтения и записи в хранилище, в том числе в рамках учета газа в хранилище. Когда операции чтения/записи выполняются из этого кэша, повторяющиеся операции кодека пропускаются, и на эти кэшированные операции в этих системных путях не тратится дополнительный газ.

В нашем проекте это приемлемо, поскольку мы уже используем политики отсутствия газа или фиксированного газа для системных сообщений (чтобы снизить риск DDoS-спама, сохраняя при этом детерминированное выполнение). Например, чтение EpochGroupData может использоваться логикой защиты конечных точек для идентификации участников, сокращенных из-за отсутствия выводов или отсутствия cPoC; это системная проверка, при которой быстрое выполнение предпочтительнее повторной маршаллинга/демаршаллинга.

Архитектура

┌──────────────────────────────────────────────┐ │ OptimisticStore │ │ │ │ Уровень 1: черновой вариант передачи (с учетом контекста, для каждой передачи) │ │ ↓ промах │ │ Уровень 2: Кэш блоков (в памяти, поблочно) │ │ ↓ промах │ │ Уровень 3: Серверная часть хранилища (постоянный KV/protobuf) │ │ │ │ + конфликт-трекер (наборы чтения/записи для каждой передачи) │ └───────────────────────────────────────────────────┘
  • Tx Draft: буфер отложенной записи, ограниченный одной транзакцией через context.Value . Создан в AnteHandler, зафиксирован для блокировки кеша при успешной передаче в PostHandler, отброшен в случае сбоя.
  • Кэш блоков: общая карта в памяти для текущего блока. Заполняется при первом чтении (шаблон кэширования). Сбрасывается в серверную часть постоянного хранилища в EndBlock .
  • Серверная часть хранилища: подключаемый интерфейс (загрузка/сохранение/удаление/клонирование), который оборачивает любое постоянное хранилище (карту коллекций, необработанный KV и т. д.).

Жизненный цикл: CheckTx против DeliverTx

CheckTx (проверка мемпула)

Во время CheckTx SDK проверяет транзакцию перед ее попаданием в мемпул. Оптимистичный магазин не фиксирует черновики во время CheckTx:

PostHandler: если ctx. IsCheckTx () || симулировать { return // пропустить фиксацию — только проверка мемпула }

Это означает, что CheckTx видит состояние постоянного хранилища (или блокирует кеш, если он уже теплый), но его записи отбрасываются. Это правильно, поскольку CheckTx не должен иметь побочных эффектов.

DeliverTx (выполнение блока)

Во время DeliverTx выполняется полный жизненный цикл:

AnteHandler (OptimisticStoreDraftDecorator): → StoreGroup.WithDraftAll(ctx) // прикрепляем черновики для каждой передачи → StoreGroup.RegisterTxAll(ctx) // регистрируемся для отслеживания OCC Выполнение сообщения: → store.Get(ctx, key) // читает: черновик → кэш блоков → серверная часть → store.Set(ctx, key) // пишет: перейти к черновику PostHandler (OptimisticStoreCommitPostDecorator): в случае успеха: → StoreGroup.CommitDraftAll(ctx) // объединяем черновик → блокируем кеш else: → StoreGroup.ReleaseDraftAll(ctx) // отбрасываем черновик EndBlock: → StoreGroup.FlushAll(ctx) // блокируем кеш → постоянное хранилище BeginBlock (следующий блок): → StoreGroup.InvalidateAll() // очищаем кеш блоков + средство отслеживания конфликтов

Отраслевые шашки

Для операций, требующих спекулятивного выполнения в пределах одной передачи (например. Шаблоны CacheContext), OptimisticStore поддерживает черновики ветвей:

кэшCtx, writeCache:= хранитель. CacheContext ( ctx ) // спекулятивная работа над кэшем Ctx... writeCache () // объединяем черновик ветки с черновиком родительского tx + кеш SDK

Метод StoreGroup.CacheContext создает черновики ветвей для каждого зарегистрированного магазина и возвращает функцию объединенной фиксации.

Как OptimisticStore работает с существующим хранилищем

OptimisticStore — это декоратор — он оборачивает существующее хранилище, не заменяя его. Интерфейс StoreBackend делает это явным образом:

type StoreBackend [ K сравнимый , V любой ] struct { Load func ( ctx context. Контекст, клавиша K) (V, bool) Сохранить функцию (ctx context. Контекст, ключ K, val V) Удалить func (ctx context. Контекст, ключ K) Функция клонирования (val V) V }

Любое существующее хранилище Collections.Map, Collections.Item или необработанное хранилище KV можно обернуть, предоставив эти четыре функции.

Пример: упаковка карты коллекций

Для EpochGroupData, который хранится в коллекции.Map[Pair[uint64,string], EpochGroupData]:

// В keeper.go NewKeeper: epochGroupStore : NewOptimisticCollMap [ epochGroupCacheKey , Collections. Пара [ uint64 , string ], типы. EpochGroupData , ]( sb , типы . EpochGroupDataPrefix, «epoch_group_data», коллекции. PairKeyCodec (коллекции . Uint64Key, коллекции. StringKey), кодек. CollValue [типы. EpochGroupData ]( cdc ),cacheConfig , func (ключ epochGroupCacheKey) коллекции. Pair [uint64, string] {возврат коллекций. Присоединяйтесь (ключ. Эпоха, ключ. ModelId ) }, ),

NewOptimisticCollMap создает коллекции.Map И обертывает их OptimisticStore за один вызов. Функции загрузки/сохранения/удаления делегируются коллекции; Clone использует proto.Clone.

Пример: упаковка синглтона (Params)

Для Params хранится как один объект protobuf с фиксированным ключом KV:

paramsStore: NewOptimisticProtoItem [types. Параметры](storeService, cdc, типы. ParamsKey, "params",cacheConfig, ),

NewOptimisticProtoItem обрабатывает внутреннюю маршализацию/демаршализацию. GetParams и SetParams просто вызывают paramsStore.Get(ctx)/paramsStore.Set(ctx, val) .

Регистрация в StoreGroup

Каждый оптимистичный магазин должен быть зарегистрирован в StoreGroup хранителя, чтобы методы жизненного цикла ( InvalidateAll , FlushAll , WithDraftAll и т. д.) применялись ко всем магазинам единообразно:

к. группа магазинов. Зарегистрируйтесь (k. epochGroupStore. OptimisticStore) k. группа магазинов. Зарегистрируйтесь (k. paramsStore. Магазин ())

Добавление нового оптимистичного магазина в Keeper

Чтобы обернуть новую коллекцию оптимистическим кэшированием:

Определите тип ключа кэша (должен быть сопоставимым): введите myNewCacheKey struct { Field1 uint64 Field2 string }

Определите тип ключа кэша (должен быть сопоставимым):

введите myNewCacheKey struct { Field1 uint64 Field2 string }

  • Добавьте в Keeper поле оптимистического магазина: myNewStore * OptimisticCollMap [myNewCacheKey, Collections. Пара [ uint64 , string ], типы. МойНовыйТип]

Добавьте в Keeper поле оптимистичного магазина:

myNewStore * OptimisticCollMap [myNewCacheKey, Collections. Пара [ uint64 , string ], типы. МойНовыйТип]
  • Инициализируйте в NewKeeper, используя NewOptimisticCollMap (или NewOptimisticProtoItem для одиночных элементов).

Инициализируйте в NewKeeper, используя NewOptimisticCollMap (или NewOptimisticProtoItem для одиночных элементов).

  • Зарегистрируйтесь в группе магазинов: k . группа магазинов. Зарегистрируйтесь (k.myNewStore. ОптимистическийМагазин)

Зарегистрируйтесь в группе магазина:

к. группа магазинов. Зарегистрируйтесь (k.myNewStore. ОптимистическийМагазин)
  • Напишите в Keeper методы получения/установки, которые делегируют хранилище: func ( k Keeper ) GetMyNewData ( ctx context. Контекст, ключ myNewCacheKey) (types. MyNewType, bool) { return k. мойНовыйМагазин. Get (ctx, key) } func (k Keeper) SetMyNewData (ctx context. Контекст, типы значений. MyNewType) {ключ:= myNewCacheKey {Field1: val. Поле1, Поле2: значение. Поле2 } k . мойНовыйМагазин. Установить (ctx, ключ, val)}

Напишите методы получения/установки в Keeper, которые делегируют хранилище:

func (k Keeper) GetMyNewData (ctx context. Контекст, ключ myNewCacheKey) (types. MyNewType, bool) { return k. мойНовыйМагазин. Get (ctx, key) } func (k Keeper) SetMyNewData (ctx context. Контекст, типы значений. MyNewType) {ключ:= myNewCacheKey {Field1: val. Поле1, Поле2: значение. Поле2 } k . мойНовыйМагазин. Установить (ctx, ключ, val)}

Никаких изменений в AnteHandler, PostHandler, BeginBlock или EndBlock не требуется — StoreGroup автоматически обрабатывает все зарегистрированные хранилища.

Подключение: AnteHandler, PostHandler, BeginBlock, EndBlock

АнтеХандлер ( ante.go )

тип OptimisticStoreDraftDecorator struct { InferenceKeeper * inferencemodulekeeper. Keeper } func (d OptimisticStoreDraftDecorator) AnteHandle (ctx sdk. Контекст, передача SDK. Tx, симулируем bool, следующий SDK. AnteHandler) (sdk. Контекст, ошибка) { g := d . Хранитель выводов. StoreGroup () newCtx:= ctx. Сконтекстом (g. WithDraftAll (ctx. Контекст ())) г . RegisterTxAll (newCtx) return next (newCtx, tx, имитация) }

Постхандлер ( ante.go )

func (d OptimisticStoreCommitPostDecorator) PostHandle (ctx sdk. Контекст, передача SDK. Tx, симуляция, успех bool, следующий SDK. PostHandler) (sdk. Контекст, ошибка) {newCtx, err:= next (ctx, tx, имитация, успех), если ctx. IsCheckTx () || имитация { return newCtx, nil // никаких побочных эффектов во время CheckTx } g := d . Хранитель выводов. StoreGroup () если успех {g. CommitDraftAll (newCtx) // объединяем черновик → кэш блока } else { g . ReleaseDraftAll (newCtx) // отбрасываем черновик } return newCtx, nil }

БегинБлок (модуль.go)

func (am AppModule) BeginBlock (ctx context. Контекст) ошибка {am. хранитель. Группа магазинов (). InvalidateAll () // очистка кеша блоков + трекер конфликтов // ... оставшаяся часть начального блока }

EndBlock (модуль.go)

func (am AppModule) EndBlock (ctx context. Контекст) ошибка {отложить утра. хранитель. Группа магазинов (). FlushAll ( ctx ) // сохранение кеша блоков для хранения // ... остальная часть конечного блока }

Дорожная карта этапа 2: Планировщик OCC

На втором этапе будет добавлен детерминированный параллельный планировщик. Текущий конфликттрекер внутри OptimisticStore уже отслеживает наборы чтения/записи для каждой передачи и может обнаруживать конфликты с помощью DetectConflicts() . Что осталось:

  • Пакетный планировщик — выберите N сообщений, сгруппируйте их по прогнозируемому шаблону доступа, выполняйте параллельно.
  • Разрешение конфликтов — после завершения пакета вызовите DetectConflicts() в каждом магазине, откатите проигравших (по порядку блоков), перепланируйте следующий пакет.
  • Retry Budget — принудительное последовательное выполнение после N неудачных попыток.
  • Интеграция — замена последовательного цикла DeliverTx пакетным планировщиком; сохраните тот же жизненный цикл проекта AnteHandler/PostHandler.

API OptimisticStore уже предназначен для этого:

// Включить обнаружение конфликтов при запуске: // COSMOS_OCC_ENABLED=1 // После завершения параллельного пакета:conflictedReads,conflictedWrites:= store. DetectConflicts () // Откат и повторное планирование конфликтующих txID... store . СброситьКонфликтТрекер ()

На этапе 2 не ожидается никаких изменений в уровне хранилища, хранителя или кэша — только добавление планировщика и замена цикла DeliverTx.

tcharchian avatar
tcharchianMaintainerMaintainer
2026-03-26
Оригинал
akup avatar
akupMaintainerАвтор
2026-03-26

Two-Level Block/Tx Cache with Conflict Detection for OCC

Motivation

Cosmos SDK processes transactions sequentially by default. Every read from the KV store requires protobuf unmarshalling, and every write requires marshalling — operations that are both CPU-intensive and repeated across transactions within the same block.

The goal of this proposal is twofold:

  • Phase 1 (current): Introduce a two-level block/tx cache implemented by OptimisticStore (name kept for code compatibility) to eliminate redundant marshal/unmarshal within a block while remaining fully deterministic and compatible with the existing sequential execution model.

Phase 1 (current): Introduce a two-level block/tx cache implemented by OptimisticStore (name kept for code compatibility) to eliminate redundant marshal/unmarshal within a block while remaining fully deterministic and compatible with the existing sequential execution model.

  • Phase 2 (future): Reuse the same conflict-tracking infrastructure to enable Optimistic Concurrency Control (OCC) — parallel execution of transactions with automatic conflict detection and rollback.

Phase 2 (future): Reuse the same conflict-tracking infrastructure to enable Optimistic Concurrency Control (OCC) — parallel execution of transactions with automatic conflict detection and rollback.

Problem Statement

The issue is not a classical race condition — it is a root of non-determinism in optimistic (parallel) behaviour . When two transactions execute in parallel and touch overlapping store keys, the final state depends on execution timing rather than block ordering. This breaks consensus.

The Cosmos SDK's optional optimistic execution mode is not a universal, one-size-fits-all feature. Its implementation is highly case-specific to the module's access patterns. Gonka therefore implements its own OCC scheme tailored to inference-chain workloads.

Design: OCC for Gonka

Scheduling

A scheduler takes N messages from the mempool (where N scales with CPU cores). Messages are sorted and grouped by similarity — similar messages tend to access similar store keys and have comparable execution times.

The selected N messages run as a parallel batch . The scheduler waits until all messages in the batch complete, then takes the next batch.

Conflict Detection

During execution, each transaction's read-set and write-set are recorded by the conflictTracker inside every OptimisticStore . After the batch completes:

  • Read-Write conflict: Transaction A read key K, transaction B wrote key K → A must be rolled back and rescheduled.
  • Write-Write conflict: Transactions A and B both wrote key K → all but one (the one earlier in block order) must be rolled back and rescheduled.

Only the losing transactions are rescheduled to the next batch. When ordering and batch-fill rules are deterministic, the entire parallel execution remains deterministic across all validators.

Conflict Resolution Ordering

When two transactions in the same batch both write the same key, the winner is determined by block order — the transaction that appears earlier in the ordered block survives; the later one is rolled back. This ensures every validator makes the same choice deterministically.

Grouping by Similarity

Grouping criterion is the predicted access pattern of the message type. Two MsgValidation for the same model and epoch are "similar" — they access the same keys and take roughly the same time to execute. Static analysis per message type provides the access prediction. Most of hot messages will be gone when shardchains start to work, but anyway this approach continues to work and will be useful.

Gas for Rescheduled Transactions

When a transaction is rolled back due to a conflict, its gas meter resets . The rescheduled execution is a fresh attempt with a fresh meter. Cosmos gas accounting therefore remains correct — the user's gas limit applies to the successful execution, not to failed speculative attempts.

High-Contention Liveness

By design, there should be no persistent hot keys. Shared data is mostly read and then written in deterministic flow, so sustained conflicts should not appear in normal operation.

As a safety fallback only, if repeated conflicts still happen due to unexpected workload shape, a transaction can be forced into sequential execution after N retries (configurable), guaranteeing forward progress.

Throughput Improvement

Implementing this OCC scheme could yield N × k times better throughput for the mainchain, where:

N = number of CPU cores

k = efficiency coefficient (< 1.0, depends on conflict rate)

For workloads with low key overlap (e.g. inferences touching different models), k approaches 1.0 and throughput scales nearly linearly with cores.

Phase 1: Two-Level Block/Tx Cache (Current Implementation)

Phase 1 is fully implemented and merged. There is no parallel scheduler yet. All transactions still execute sequentially. This is not optimistic execution by itself; it is a deterministic 2-level cache (tx draft + block cache) with conflict detection that can be used later in optimistic concurrent mode.

This cache provides these benefits:

  • Eliminate repeated marshal/unmarshal for store values accessed multiple times within a block (e.g. Params , EpochGroupData ).
  • Lay the foundation for Phase 2 by tracking read/write sets per transaction, so conflict detection can be enabled with a single environment variable ( COSMOS_OCC_ENABLED=1 ).

Shared System Data and Cache Scope

We use this 2-level block/tx cache for frequently read shared system data:

Params

EpochGroupData

These are hot system reads, so caching removes repeated codec overhead on the same values inside a block and transaction flow.

Protobuf Marshal/Unmarshal and Gas

Cosmos SDK marshals/unmarshals protobuf data on store reads/writes, including as part of store gas accounting. When reads/writes are served from this cache, repeated codec operations are skipped and no extra gas is spent for those cached operations in these system paths.

This is acceptable in our design because we already use no-gas or fixed-gas policies for system messages (to reduce DDoS spam risk while keeping deterministic execution). For example, EpochGroupData reads may be used by endpoint protection logic to identify participants slashed for missing inferences or missing cPoCs; this is a system check where fast execution is preferred over repeated marshalling/unmarshalling.

Architecture

┌───────────────────────────────────────────────────┐ │ OptimisticStore │ │ │ │ Layer 1: Tx Draft (context-scoped, per-tx) │ │ ↓ miss │ │ Layer 2: Block Cache (in-memory, per-block) │ │ ↓ miss │ │ Layer 3: Store Backend (persistent KV / protobuf) │ │ │ │ + conflictTracker (read/write sets per tx) │ └───────────────────────────────────────────────────┘
  • Tx Draft: Write-behind buffer scoped to a single transaction via context.Value . Created in AnteHandler , committed to block cache on tx success in PostHandler , discarded on failure.
  • Block Cache: Shared in-memory map for the current block. Populated on first read (cache-aside pattern). Flushed to the persistent store backend in EndBlock .
  • Store Backend: Pluggable interface ( Load / Save / Delete / Clone ) that wraps any persistent storage (collections map, raw KV, etc.).

Lifecycle: CheckTx vs DeliverTx

CheckTx (Mempool Validation)

During CheckTx , the SDK validates a transaction before it enters the mempool. The optimistic store does not commit drafts during CheckTx :

PostHandler: if ctx . IsCheckTx () || simulate { return // skip commit — mempool validation only }

This means CheckTx sees the persistent store state (or block cache if already warm), but its writes are discarded. This is correct because CheckTx must be side-effect-free.

DeliverTx (Block Execution)

During DeliverTx , the full lifecycle executes:

AnteHandler (OptimisticStoreDraftDecorator): → StoreGroup.WithDraftAll(ctx) // attach per-tx drafts → StoreGroup.RegisterTxAll(ctx) // register for OCC tracking Message Execution: → store.Get(ctx, key) // reads: draft → block cache → backend → store.Set(ctx, key) // writes: go to draft PostHandler (OptimisticStoreCommitPostDecorator): if success: → StoreGroup.CommitDraftAll(ctx) // merge draft → block cache else: → StoreGroup.ReleaseDraftAll(ctx) // discard draft EndBlock: → StoreGroup.FlushAll(ctx) // block cache → persistent store BeginBlock (next block): → StoreGroup.InvalidateAll() // clear block cache + conflict tracker

Branch Drafts

For operations that need speculative execution within a single tx (e.g. CacheContext patterns), OptimisticStore supports branch drafts :

cacheCtx , writeCache := keeper . CacheContext ( ctx ) // speculative work on cacheCtx... writeCache () // merge branch draft into parent tx draft + SDK cache

The StoreGroup.CacheContext method creates branch drafts for every registered store and returns a merged commit function.

How OptimisticStore Works Around Existing Storage

OptimisticStore is a decorator — it wraps existing storage without replacing it. The StoreBackend interface makes this explicit:

type StoreBackend [ K comparable , V any ] struct { Load func ( ctx context. Context , key K ) ( V , bool ) Save func ( ctx context. Context , key K , val V ) Delete func ( ctx context. Context , key K ) Clone func ( val V ) V }

Any existing collections.Map , collections.Item , or raw KV store can be wrapped by providing these four functions.

Example: Wrapping a Collections Map

For EpochGroupData , which is stored in a collections.Map[Pair[uint64,string], EpochGroupData] :

// In keeper.go NewKeeper: epochGroupStore : NewOptimisticCollMap [ epochGroupCacheKey , collections. Pair [ uint64 , string ], types. EpochGroupData , ]( sb , types . EpochGroupDataPrefix , "epoch_group_data" , collections . PairKeyCodec ( collections . Uint64Key , collections . StringKey ), codec. CollValue [types. EpochGroupData ]( cdc ), cacheConfig , func ( key epochGroupCacheKey ) collections. Pair [ uint64 , string ] { return collections . Join ( key . Epoch , key . ModelId ) }, ),

NewOptimisticCollMap creates the collections.Map AND wraps it with an OptimisticStore in one call. The Load / Save / Delete functions delegate to the collection; Clone uses proto.Clone .

Example: Wrapping a Singleton (Params)

For Params , stored as a single protobuf blob at a fixed KV key:

paramsStore : NewOptimisticProtoItem [types. Params ]( storeService , cdc , types . ParamsKey , "params" , cacheConfig , ),

NewOptimisticProtoItem handles marshal/unmarshal internally. GetParams and SetParams simply call paramsStore.Get(ctx) / paramsStore.Set(ctx, val) .

Registering with the StoreGroup

Every optimistic store must be registered with the keeper's StoreGroup so lifecycle methods ( InvalidateAll , FlushAll , WithDraftAll , etc.) apply to all stores uniformly:

k . storeGroup . Register ( k . epochGroupStore . OptimisticStore ) k . storeGroup . Register ( k . paramsStore . Store ())

Adding a New Optimistic Store to the Keeper

To wrap a new collection with optimistic caching:

Define a cache key type (must be comparable ): type myNewCacheKey struct { Field1 uint64 Field2 string }

Define a cache key type (must be comparable ):

type myNewCacheKey struct { Field1 uint64 Field2 string }

  • Add the optimistic store field to Keeper : myNewStore * OptimisticCollMap [ myNewCacheKey , collections. Pair [ uint64 , string ], types. MyNewType ]

Add the optimistic store field to Keeper :

myNewStore * OptimisticCollMap [ myNewCacheKey , collections. Pair [ uint64 , string ], types. MyNewType ]
  • Initialize in NewKeeper using NewOptimisticCollMap (or NewOptimisticProtoItem for singletons).

Initialize in NewKeeper using NewOptimisticCollMap (or NewOptimisticProtoItem for singletons).

  • Register with the store group: k . storeGroup . Register ( k . myNewStore . OptimisticStore )

Register with the store group:

k . storeGroup . Register ( k . myNewStore . OptimisticStore )
  • Write getter/setter methods on Keeper that delegate to the store: func ( k Keeper ) GetMyNewData ( ctx context. Context , key myNewCacheKey ) (types. MyNewType , bool ) { return k . myNewStore . Get ( ctx , key ) } func ( k Keeper ) SetMyNewData ( ctx context. Context , val types. MyNewType ) { key := myNewCacheKey { Field1 : val . Field1 , Field2 : val . Field2 } k . myNewStore . Set ( ctx , key , val ) }

Write getter/setter methods on Keeper that delegate to the store:

func ( k Keeper ) GetMyNewData ( ctx context. Context , key myNewCacheKey ) (types. MyNewType , bool ) { return k . myNewStore . Get ( ctx , key ) } func ( k Keeper ) SetMyNewData ( ctx context. Context , val types. MyNewType ) { key := myNewCacheKey { Field1 : val . Field1 , Field2 : val . Field2 } k . myNewStore . Set ( ctx , key , val ) }

No changes are needed to AnteHandler , PostHandler , BeginBlock , or EndBlock — the StoreGroup handles all registered stores automatically.

Wiring: AnteHandler, PostHandler, BeginBlock, EndBlock

AnteHandler ( ante.go )

type OptimisticStoreDraftDecorator struct { InferenceKeeper * inferencemodulekeeper. Keeper } func ( d OptimisticStoreDraftDecorator ) AnteHandle ( ctx sdk. Context , tx sdk. Tx , simulate bool , next sdk. AnteHandler ) (sdk. Context , error ) { g := d . InferenceKeeper . StoreGroup () newCtx := ctx . WithContext ( g . WithDraftAll ( ctx . Context ())) g . RegisterTxAll ( newCtx ) return next ( newCtx , tx , simulate ) }

PostHandler ( ante.go )

func ( d OptimisticStoreCommitPostDecorator ) PostHandle ( ctx sdk. Context , tx sdk. Tx , simulate , success bool , next sdk. PostHandler ) (sdk. Context , error ) { newCtx , err := next ( ctx , tx , simulate , success ) if ctx . IsCheckTx () || simulate { return newCtx , nil // no side effects during CheckTx } g := d . InferenceKeeper . StoreGroup () if success { g . CommitDraftAll ( newCtx ) // merge draft → block cache } else { g . ReleaseDraftAll ( newCtx ) // discard draft } return newCtx , nil }

BeginBlock ( module.go )

func ( am AppModule ) BeginBlock ( ctx context. Context ) error { am . keeper . StoreGroup (). InvalidateAll () // clear block caches + conflict tracker // ... rest of begin block }

EndBlock ( module.go )

func ( am AppModule ) EndBlock ( ctx context. Context ) error { defer am . keeper . StoreGroup (). FlushAll ( ctx ) // persist block caches to store // ... rest of end block }

Phase 2 Roadmap: OCC Scheduler

Phase 2 will add a deterministic parallel scheduler. The current conflictTracker inside OptimisticStore already tracks per-tx read/write sets and can detect conflicts via DetectConflicts() . What remains:

  • Batch Scheduler — select N messages, group by predicted access pattern, execute in parallel.
  • Conflict Resolution — after batch completion, call DetectConflicts() on each store, roll back losers (by block order), reschedule to next batch.
  • Retry Budget — force sequential execution after N failed attempts.
  • Integration — replace the sequential DeliverTx loop with the batched scheduler; keep the same AnteHandler / PostHandler draft lifecycle.

The OptimisticStore API is already designed for this:

// Enable conflict detection at startup: // COSMOS_OCC_ENABLED=1 // After parallel batch completes: conflictedReads , conflictedWrites := store . DetectConflicts () // Roll back and reschedule conflicted txIDs... store . ResetConflictTracker ()

No changes to the store, keeper, or cache layer are expected for Phase 2 — only the addition of the scheduler and the DeliverTx loop replacement.

tcharchian avatar
tcharchianMaintainerMaintainer
2026-03-26