Коллатор

Чтение сообщений

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

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

Группа сообщений и слоты

Группа устроена как набор слотов, но привязка выстроена от аккаунта: не слот целиком закреплён за одним аккаунтом-получателем, а аккаунт живёт ровно в одном слоте и не может оказаться сразу в двух — попытка добавить его сообщения во второй слот отклоняется. Один и тот же слот может донабираться сообщениями другого аккаунта: если первый аккаунт отдал слоту все свои сообщения раньше, чем слот заполнился до предела, освободившееся место займёт другой аккаунт из буфера. Аккаунт не может занимать два слота одновременно, поэтому слоты не пересекаются по аккаунтам, и группу можно исполнять параллельно. Порядок при этом сохраняется в очереди каждого аккаунта: его сообщения идут строго друг за другом, а не как единый сквозной порядок «внутри слота», если в слоте успели побывать сообщения нескольких аккаунтов.

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

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

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

Диапазоны и позиция обработки

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

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

Параметр open_ranges_limit ограничивает число диапазонов, которые чтение сообщений одновременно держит открытыми в позиции обработки.

Просроченные внешние сообщения

У внешнего сообщения есть срок жизни: параметр externals_expire_timeout отсчитывает секунды от времени цепочки якоря, из которого сообщение пришло. Коллатор проверяет срок дважды, потому что сообщение может пролежать в буфере несколько групп подряд, а время цепочки за это время уйдёт вперёд — одной проверки на чтении недостаточно.

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

Дозаполнение буферов после перезапуска

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

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

Партиция аккаунта в очереди

Партицию очереди сообщение получает один раз, в момент записи дифа очереди, решением роутера партиций. Партиция 0 — основная: в неё попадает сообщение, если роутер не смог отнести его ни к какой другой партиции. Партиция 1 низкоприоритетная, и какие аккаунты в неё переводятся, решает не сам роутер, а чтение сообщений: при завершении сборки каждого блока оно пересчитывает это решение заново.

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

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

Транзакционность чтения

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

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

Значения по умолчанию

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

ПараметрЗначениеЕдиницаСмысл
buffer_limit10 000сообщенийпредельное число сообщений в одном буфере
group_limit100слотовсколько слотов может быть в группе
group_vert_size10сообщенийсколько сообщений умещается в одном слоте
externals_expire_timeout58секундсрок жизни внешнего сообщения от времени цепочки его якоря
open_ranges_limit20диапазоновсколько диапазонов чтения хранится в позиции обработки
par_0_int_msgs_count_limit100 000сообщенийпорог накопленных внутренних сообщений на аккаунт, после которого он переходит в низкоприоритетную партицию
par_0_ext_msgs_count_limit10 000 000сообщенийтот же порог для внешних сообщений, но в версии 0.3.11 в решении о переводе не участвует — правило смотрит только на внутренние сообщения
group_slots_fractions{0: 80, 1: 10}% от group_limitдоли слотов группы по партициям
range_messages_limit10 000сообщенийсколько сообщений своего шарда может покрывать один диапазон внутренних сообщений (ноль читается как 10 000)