Коллатор

Внутренняя очередь сообщений

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

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

Диф очереди

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

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

Зоны очереди

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

Коммит дифов

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

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

Проверка перед записью

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

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

Сброс незакоммиченной зоны

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

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

Партиции и роутер

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

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

Сборка мусора

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

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

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

База данных очереди

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

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

Импорт из устойчивого состояния

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

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

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

ПараметрЗначение по умолчаниюСмысл
gc_interval"5s"период тика фонового сборщика мусора очереди