Tarantool CE/EE Documentation portal logo
Помощь
Обновлена 15 сентября 2026 г. в 08:55

Модуль swim

Общие сведения

Модуль swim содержит реализацию SWIM в Tarantool — SWIM (Scalable Weakly-consistent Infection-style Process Group Membership Protocol) — масштабируемого слабо-согласованного протокола членства в группе процессов с инфекционным типом распространения. Он рекомендуется для любого типа кластера Tarantool, где количество узлов может быть большим. Назначение протокола — обнаружение и мониторинг других участников кластера, а также хранение информации о них в «таблице участников». Протокол работает путем периодической отправки и получения сообщений по UDP в фоновом цикле событий.

Каждое сообщение состоит из нескольких частей, включая:

  • ping, например «Я проверяю, жив ли ты»,
  • событие, например «Я присоединяюсь»,
  • антиэнтропию, например «Я знаю, что существует другой участник»,
  • полезную нагрузку, например «У меня или у другого участника могут быть пользовательские данные».

Максимальный размер сообщения — около 1500 байт.

SWIM периодически отправляет сообщения случайному подмножеству из таблицы участников. Ответы от этих участников обрабатываются асинхронно.

Каждая запись в таблице участников содержит:

  • UUID,
  • статус ("alive", "suspected", "dead" или "left"). Если участник не подтверждает определенное количество ping-запросов, его статус меняется с "alive" на "suspected", то есть он подозревается в недоступности. Но SWIM старается избегать ложных срабатываний (ошибочного определения участников как недоступных), которые могут возникать, когда участник перегружен и слишком медленно отвечает на ping-запросы, или когда есть проблемы с сетью и пакеты не могут пройти через некоторые каналы. Когда участник оказывается под подозрением, SWIM случайным образом выбирает других участников и отправляет им запросы: «пожалуйста, отправьте ping этому подозреваемому участнику». Это называется косвенным ping (indirect ping). Таким образом, через разные маршруты и дополнительные переходы подозреваемый участник получает дополнительные шансы ответить и тем самым «опровергнуть» подозрение.

Поскольку выбор случаен, обеспечивается равномерная сетевая нагрузка — около одного сообщения на участника за один шаг протокола, независимо от размера кластера. Это ключевая особенность SWIM. Поскольку протокол опирается на передачу информации между участниками, также известную как «обмен сплетнями» (gossiping), участникам не нужно рассылать сообщения всем остальным, что привело бы к сетевой нагрузке в N сообщений на участника за один шаг протокола, где N — количество участников в кластере. Однако выбор не является полностью случайным: предпочтение отдается участникам, которые дольше всего не получали ping-запросов, по принципу round-robin.

Что касается антиэнтропийной части сообщения: она необходима для поддержания статуса в записях таблицы участников. Рассмотрим пример: два участника, №1 и №2, оба активны. Событий не происходит, поэтому периодически отправляются только ping-запросы. Затем появляется третий участник, №3. Он знает об одном из существующих участников, №2. Как ему обнаружить другого участника? Конечно, №1 мог бы уведомить №2, а №2 — №3, но сообщения передаются по UDP, поэтому любое уведомление может быть потеряно. Однако обычные сообщения, содержащие "ping" и/или "событие", также могут содержать секцию "антиэнтропия", которая берется из случайно выбранной части таблицы участников. Так, в данном примере №2 в конечном итоге случайным образом добавит к обычному сообщению антиэнтропийную заметку о том, что №1 активен, и таким образом №3 обнаружит №1, даже не получив прямого сообщения-события «Я активен» от №1.

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

Что касается полезной нагрузки (payload) в сообщении: она нужна не всегда — это функция, позволяющая передавать пользовательскую информацию через SWIM вместо прямой связи между узлами. Модуль swim предоставляет методы для указания «полезной нагрузки» — произвольных пользовательских данных максимальным размером около 1,2 КБ. Полезная нагрузка может быть любой и в конечном итоге будет распространена по кластеру и доступна другим участникам. Каждый участник может иметь собственную полезную нагрузку.

Сообщения могут быть зашифрованы. Шифрование может не требоваться в закрытой сети, но необходимо для безопасности, если кластер находится в публичном интернете. Пользователи могут указать алгоритм шифрования, режим шифрования и приватный ключ. Все части всех сообщений (включая ping, подтверждение, событие, полезную нагрузку, URI и UUID) будут зашифрованы этим приватным ключом, а также случайным публичным ключом, генерируемым для каждого сообщения, чтобы предотвратить атаки на основе шаблонов.

Теоретически скорость распространения событий (количество переходов для передачи информации по всему кластеру) составляет O(log(cluster_size)). Подробнее об этом и другой теоретической информации см. в статье Университета Корнелла, где SWIM был описан изначально.

swim.new([cfg])

Создание нового экземпляра SWIM. Экземпляр SWIM поддерживает таблицу участников и взаимодействует с другими участниками. В одном процессе Tarantool можно создать несколько экземпляров SWIM.

Параметры:

  • cfg (table) — необязательный параметр конфигурации.

Если cfg не указан или равен nil, новый экземпляр SWIM не привязывается к сокету и имеет атрибуты со значением nil, поэтому он не может взаимодействовать с другими участниками, и только несколько методов доступны до вызова swim_object:cfg().

Если cfg указан, эффект такой же, как при вызове s = swim.new() s:cfg(), за исключением generation. Описание конфигурации см. в swim_object:cfg().

Параметр generation в cfg можно указать только при вызове new(), его нельзя указать позже при вызове cfg(). Generation является частью incarnation. Обычно generation не указывается, поскольку значения по умолчанию (временная метка) достаточно, но если есть причины не доверять временным меткам (из-за изменения времени или запуска экземпляра на другой машине), пользователи могут указать swim.new({generation = <number>}). В этом случае последнее значение следует каким-либо образом сохранить (например, в файле, спейсе или глобальном сервисе), а новое значение должно быть больше любого предыдущего значения generation.

Возвращает

объект swim

Пример:

swim_object = swim.new({uri = 3333, uuid = '00000000-0000-1000-8000-000000000001', heartbeat_rate = 0.1})

Объект swim

Объект swim — это объект, возвращаемый swim.new(). Он имеет следующие методы: cfg(), delete(), is_configured(), size(), quit(), add_member(), remove_member(), probe_member(), broadcast(), set_payload(), set_payload_raw(), set_codec(), self(), member_by_uuid(), pairs(), on_member_event(), а также свойство cfg.

swim_object:cfg(cfg)

Настройка или перенастройка экземпляра SWIM.

Параметры:

  • cfg (table) — параметры, описывающие поведение экземпляра

Таблица cfg может содержать следующие компоненты:

  • heartbeat_rate (double) — частота отправки сообщений раунда в секундах. Установка значения X для параметра heartbeat_rate не означает, что каждый участник будет проверяться каждые X секунд — X определяет скорость протокола. Период протокола зависит от количества участников и значения heartbeat_rate. По умолчанию = 1.

  • ack_timeout (double) — время в секундах, после которого ping считается неподтвержденным. По умолчанию = 30.

  • gc_mode (enum) — режим сборки недоступных участников.

    Если gc_mode == 'off', SWIM никогда не удаляет недоступных участников из таблицы участников (хотя пользователи могут удалить их с помощью swim_object:remove_member()), и SWIM продолжит отправлять им ping-запросы, как если бы они были активны.

    Если gc_mode == 'on', SWIM удаляет недоступных участников из таблицы участников после одного раунда.

    По умолчанию = 'on'.

  • uri (string или number) — либо адрес 'ip:port', либо только номер порта (если IP опущен, предполагается 127.0.0.1). Если port == 0, ядро выберет любой свободный порт для IP-адреса.

  • uuid (string или cdata struct tt_uuid) — значение, которое должно быть уникальным среди экземпляров SWIM. Пользователи могут выбрать любое значение, но рекомендуется использовать box.cfg.instance_uuid — UUID экземпляра Tarantool.

Все компоненты cfg являются динамическими — swim_object:cfg() можно вызывать несколько раз. Если метод вызывается не в первый раз и какой-либо компонент не указан, этот компонент сохраняет свое предыдущее значение. Если метод вызывается впервые, uri и uuid обязательны, поскольку экземпляр SWIM не может работать без URI и UUID.

swim_object:cfg() выполняется атомарно — при возникновении ошибки ничего не изменяется.

Возвращает

true в случае успешной настройки

Возвращает

nil, err в случае ошибки. err — объект ошибки

Пример:

swim_object:cfg({heartbeat_rate = 0.5})

После вызова swim_object:cfg() все остальные методы swim_object становятся доступными.

swim_object.cfg

Предоставляет доступ ко всем компонентам со значением, отличным от nil, из таблицы только для чтения, которая была настроена или изменена с помощью swim_object:cfg().

Пример:

tarantool> swim_object.cfg---- gc_mode: off  uri: 3333  uuid: 00000000-0000-1000-8000-000000000001...

swim_object:delete()

Немедленное удаление экземпляра SWIM. Его память освобождается, запись из таблицы участников удаляется, и экземпляр больше не может использоваться. Остальные участники будут считать этого участника «мертвым».

После swim_object:delete() любая попытка использовать удаленный экземпляр приведет к возникновению исключения.

Возвращает

нет; этот метод не завершается ошибкой

Пример: swim_object:delete()

swim_object:is_configured()

Возврат false, если экземпляр SWIM был создан с помощью swim.new() без необязательного аргумента cfg и не был настроен с помощью swim_object:cfg(). В противном случае возвращается true.

Возвращает

boolean — true, если экземпляр настроен, иначе false

Пример: swim_object:is_configured()

swim_object:size()

Возврат размера таблицы участников. Размер будет не меньше 1, так как в таблицу включен собственный участник («self»).

Возвращает

integer — размер

Пример: swim_object:size()

swim_object:quit()

Выход из кластера.

Это корректный аналог swim_object:delete() — экземпляр удаляется, но перед удалением он отправляет каждому участнику из своей таблицы участников сообщение о том, что данный экземпляр покинул кластер и его не следует считать «мертвым».

Другие экземпляры пометят такого участника в своих таблицах как 'left' и удалят его после одного раунда распространения.

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

Возвращает

нет; метод не завершается ошибкой

Пример: swim_object:quit()

swim_object:add_member(cfg)

Явное добавление участника в таблицу участников.

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

Параметры:

  • cfg (table) — описание участника

Таблица cfg имеет два обязательных компонента, uuid и uri, формат которых совпадает с форматом uuid и uri в таблице для swim_object:cfg().

Возвращает

true, если участник добавлен

Возвращает

nil, err, если произошла ошибка. err — объект ошибки

Пример:

swim_member_object = swim_object:add_member({uuid = ..., uri = ...})

swim_object:remove_member(uuid)

Явное и немедленное удаление участника из таблицы участников.

Параметры:

  • uuid (string или cdata struct tt_uuid) — UUID

Возвращает

true, если участник удален

Возвращает

nil, err, если произошла ошибка. err — объект ошибки.

Пример: swim_object:remove_member('00000000-0000-1000-8000-000000000001')

swim_object:probe_member(uri)

Отправка ping-запроса по указанному адресу uri. Если по этому адресу прослушивает другой участник, он получит ping и ответит сообщением ACK (подтверждением), содержащим такую информацию, как UUID. Эта информация будет добавлена в таблицу участников.

swim_object:probe_member() похож на swim_object:add_member(), но не требует UUID и не является надежным, так как использует UDP.

Параметры:

  • uri (string или number) — URI. Формат такой же, как у uri в swim_object:cfg().

Возвращает

true, если участник пропингован

Возвращает

nil, err, если произошла ошибка. err — объект ошибки.

Пример: swim_object:probe_member(3333)

swim_object:broadcast([port])

Рассылка ping-запроса на все сетевые интерфейсы системы.

swim_object:broadcast() — это то же самое, что swim_object:probe_member(), но для многих участников одновременно.

Параметры:

  • port (number) — все отправленные ping-запросы имеют этот порт в качестве порта назначения в своих UDP-заголовках. По умолчанию используется текущий привязанный порт.

Возвращает

true, если широковещательное сообщение отправлено

Возвращает

nil, err, если произошла ошибка. err — объект ошибки.

Пример:

tarantool> fiber = require('fiber')---...tarantool> swim = require('swim')---...tarantool> s1 = swim.new({uri = 3333, uuid = '00000000-0000-1000-8000-000000000001', heartbeat_rate = 0.1})---...tarantool> s2 = swim.new({uri = 3334, uuid = '00000000-0000-1000-8000-000000000002', heartbeat_rate = 0.1})---...tarantool> s1:size()---- 1...tarantool> s1:add_member({uri = s2:self():uri(), uuid = s2:self():uuid()})---- true...tarantool> s1:size()---- 1...tarantool> s2:size()---- 1...tarantool> fiber.sleep(0.2)---...tarantool> s1:size()---- 2...tarantool> s2:size()---- 2...tarantool> s1:remove_member(s2:self():uuid()) s2:remove_member(s1:self():uuid())---...tarantool> s1:size()---- 1...tarantool> s2:size()---- 1...tarantool> s1:probe_member(s2:self():uri())---- true...tarantool> fiber.sleep(0.1)---...tarantool> s1:size()---- 2...tarantool> s2:size()---- 2...tarantool> s1:remove_member(s2:self():uuid()) s2:remove_member(s1:self():uuid())---...tarantool> s1:size()---- 1...tarantool> s2:size()---- 1...tarantool> s1:broadcast(3334)---- true...tarantool> fiber.sleep(0.1)---...tarantool> s1:size()---- 2...tarantool> s2:size()---- 2...

swim_object:set_payload(payload)

Установка полезной нагрузки (payload) в виде форматированных данных.

Полезная нагрузка — произвольные пользовательские данные размером до 1200 байт, которые распространяются по кластеру. В результате каждый участник кластера со временем узнает полезную нагрузку других участников, так как она хранится в таблице участников и может быть получена с помощью swim_member_object:payload().

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

Параметры:

  • payload (object) — произвольный Lua-объект для распространения. Установите nil, чтобы удалить полезную нагрузку, — в этом случае она со временем будет удалена и на других экземплярах. Объект сериализуется в MessagePack.

Возвращает

true, если полезная нагрузка установлена; иначе nil, errerr является объектом ошибки

Пример:

swim_object:set_payload({field1 = 100, field2 = 200})

swim_object:set_payload_raw(payload[, size])

Установка полезной нагрузки в виде «сырых» данных.

Иногда полезной нагрузке не обязательно быть Lua-объектом. Например, у пользователя может быть уже готовый объект MessagePack, который нужно просто установить как полезную нагрузку. Или требуется передать cdata.

set_payload_raw позволяет установить полезную нагрузку как есть, без сериализации в MessagePack.

Параметры:

  • payload (string или cdata) — любое значение

  • size (number) — размер полезной нагрузки в байтах. Если payload — строка, то size необязателен, а если указан, то не должен превышать фактический размер payload. Если size меньше фактического размера payload, используются только первые size байтов payload. Если payload — cdata, то size обязателен.

Возвращает

true, если полезная нагрузка установлена; иначе nil, errerr является объектом ошибки

Пример:

tarantool> ffi = require('ffi')---...tarantool> fiber = require('fiber')---...tarantool> swim = require('swim')---...tarantool> s1 = swim.new({uri = 0, uuid = '00000000-0000-1000-8000-000000000001', heartbeat_rate = 0.1})---...tarantool> s2 = swim.new({uri = 0, uuid = '00000000-0000-1000-8000-000000000002', heartbeat_rate = 0.1})---...tarantool> s1:add_member({uri = s2:self():uri(), uuid = s2:self():uuid()})---- true...tarantool> s1:set_payload({a = 100, b = 200})---- true...tarantool> s2:set_payload('any payload')---- true...tarantool> fiber.sleep(0.2)---...tarantool> s1_view = s2:member_by_uuid(s1:self():uuid())---...tarantool> s2_view = s1:member_by_uuid(s2:self():uuid())---...tarantool> s1_view:payload()---- {'a': 100, 'b': 200}...tarantool> s2_view:payload()---- any payload...tarantool> cdata = ffi.new('char[?]', 2)---...tarantool> cdata[0] = 1---...tarantool> cdata[1] = 2---...tarantool> s1:set_payload_raw(cdata, 2)---- true...tarantool> fiber.sleep(0.2)---...tarantool> cdata, size = s1_view:payload_cdata()---...tarantool> cdata[0]---- 1...tarantool> cdata[1]---- 2...tarantool> size---- 2...

swim_object:set_codec(codec_cfg)

Включение шифрования всех последующих сообщений.

Краткое описание алгоритмов шифрования см. в enum_crypto_algo и enum crypto_mode в исходном файле Tarantool crypto.h.

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

Параметры:

  • codec_cfg (table) — описание шифрования

Компонентами таблицы codec_cfg могут быть:

  • algo (string) — название алгоритма шифрования. Поддерживаются все названия из модуля crypto: 'aes128', 'aes192', 'aes256', 'des'. Укажите 'none', чтобы отключить шифрование.

  • mode (string) — режим алгоритма шифрования. Поддерживаются все режимы из модуля crypto: 'ecb', 'cbc', 'cfb', 'ofb'. По умолчанию — 'cbc'.

  • key (cdata или string) — закрытый секретный ключ, который должен храниться в секрете и никогда не должен быть жестко прописан в исходном коде.

  • key_size (integer) — размер ключа в байтах.

    key_size обязателен, если ключ — cdata.

    key_size необязателен, если ключ — строка; если key_size меньше фактического размера ключа, ключ усекается.

Значения algo, mode, key и key_size должны быть одинаковыми для всех экземпляров SWIM, чтобы участники могли понимать сообщения друг друга.

Пример:

tarantool> swim = require('swim')---...tarantool> s1 = swim.new({uri = 0, uuid = '00000000-0000-1000-8000-000000000001'})---...tarantool> s1:set_codec({algo = 'aes128', mode = 'cbc', key = '1234567812345678'})---- true...

swim_object:self()

Возврат объекта участника swim (собственного) из таблицы участников или из кэша, содержащего результаты более ранних вызовов swim_object:self(), swim_object:member_by_uuid() или swim_object:pairs().

Возвращает

объект участника swim, не nil, так как self() не завершается ошибкой

Пример: swim_member_object = swim_object:self()

swim_object:member_by_uuid(uuid)

Возврат объекта участника swim (по заданному UUID) из таблицы участников или из кэша, содержащего результаты более ранних вызовов swim_object:self(), swim_object:member_by_uuid() или swim_object:pairs().

Параметры:

  • uuid (string или cdata struct tt_uuid) — UUID

Возвращает

объект участника swim или nil, если участник не найден

Пример:

swim_member_object = swim_object:member_by_uuid('00000000-0000-1000-8000-000000000001')

swim_object:pairs()

Создание итератора для возврата объектов участников swim из таблицы участников или из кэша, содержащего результаты более ранних вызовов swim_object:self(), swim_object:member_by_uuid() или swim_object:pairs().

Метод swim_object:pairs() следует использовать в цикле 'for', и одновременно может работать только один итератор. (Итератор реализован в сверхлегком стиле, поэтому на один экземпляр SWIM доступен только один объект итератора.)

Параметры:

  • generator+object+key (varies) — как у любых Lua-итераторов в стиле pairs(): функция-генератор, объект итератора (объект участника swim) и начальный ключ (UUID).

Пример:

tarantool> fiber = require('fiber')---...tarantool> swim = require('swim')---...tarantool> s1 = swim.new({uri = 0, uuid = '00000000-0000-1000-8000-000000000001', heartbeat_rate = 0.1})---...tarantool> s2 = swim.new({uri = 0, uuid = '00000000-0000-1000-8000-000000000002', heartbeat_rate = 0.1})---...tarantool> s1:add_member({uri = s2:self():uri(), uuid = s2:self():uuid()})---- true...tarantool> fiber.sleep(0.2)---...tarantool> s1:self()---- uri: 127.0.0.1:55845  status: alive  incarnation: cdata {generation = 1569353431853325ULL, version = 1ULL}  uuid: 00000000-0000-1000-8000-000000000001  payload_size: 0...tarantool> s1:member_by_uuid(s1:self():uuid())---- uri: 127.0.0.1:55845  status: alive  incarnation: cdata {generation = 1569353431853325ULL, version = 1ULL}  uuid: 00000000-0000-1000-8000-000000000001  payload_size: 0...tarantool> s1:member_by_uuid(s2:self():uuid())---- uri: 127.0.0.1:53666  status: alive  incarnation: cdata {generation = 1569353431865138ULL, version = 1ULL}  uuid: 00000000-0000-1000-8000-000000000002  payload_size: 0...tarantool> t = {}---...tarantool> for k, v in s1:pairs() do table.insert(t, {k, v}) end---...tarantool> t---- - '00000000-0000-1000-8000-000000000002'  - uri: 127.0.0.1:53666    status: alive    incarnation: cdata {generation = 1569353431865138ULL, version = 1ULL}    uuid: 00000000-0000-1000-8000-000000000002    payload_size: 0- - '00000000-0000-1000-8000-000000000001'  - uri: 127.0.0.1:55845    status: alive    incarnation: cdata {generation = 1569353431853325ULL, version = 1ULL}    uuid: 00000000-0000-1000-8000-000000000001    payload_size: 0...

swim_object:on_member_event(trigger-function[, ctx])

Создание триггера on_member. Функция trigger-function будет выполнена при изменении участника в таблице участников.

Параметры:

  • trigger-function (function) — функция, которая станет функцией триггера

  • ctx (cdata) — (необязательный параметр) значение, которое будет передано в trigger-function

Возвращает

nil или указатель на функцию

У функции trigger-function должно быть три объявления параметров (Tarantool передаст в них значения при вызове функции):

  • участник, с которым произошло событие участника;
  • объект события;
  • ctx — то же значение, что передано в swim_object:on_member_event.

Событием участника может быть любое из:

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

Объект события — объект, который функция триггера может использовать, чтобы определить, какой тип события участника произошел. Методы объекта — например, is_new_status(), is_new_uri(), is_new_incarnation(), is_new_payload(), is_drop() — возвращают булевы значения.

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

Пример функции триггера on_member_event:

swim = require('swim')local function on_event(member, event, ctx)    if event:is_new() then        ...    elseif event:is_drop() then        ...    end    if event:is_update() then        -- All next conditions can be        -- true simultaneously.        if event:is_new_status() then            ...        end        if event:is_new_uri() then            ...        end        if event:is_new_incarnation() then            ...        end        if event:is_new_payload() then            ...        end    endend

Обратите внимание, что в примере выше функция готова к возможности одновременного наступления нескольких событий при одной активации триггера. is_new() и is_drop() не могут быть истинными одновременно, но is_new() и is_update() могут, как и is_drop() и is_update(). Одновременные события особенно вероятны, если событий много, а функции триггеров медленные — в этом случае, например, участник может быть добавлен, затем через некоторое время изменен, и затем через некоторое время произойдет одна активация триггера.

Кроме того: is_new() и is_new_payload() могут быть истинными одновременно. Этот случай не связан с медленными функциями триггеров. Он возникает из-за того, что «опущенная полезная нагрузка» и «полезная нагрузка нулевого размера» — не одно и то же. Например: при получении ping-сообщения новый участник может быть добавлен, но ping-сообщения не содержат полезной нагрузки. Полезная нагрузка появится позже в другом сообщении. Если это важно для приложения, функция не должна предполагать при истинном is_new(), что у участника уже есть полезная нагрузка, и не должна судить о наличии или отсутствии полезной нагрузки по ее размеру.

Также: функции не должны предполагать, что is_new() и is_drop() обязательно будут увидены. Если новый участник появился, но был удален до того, как его появление вызвало активацию триггера, активации триггера не будет.

is_new_generation() вернет true, если изменилась часть generation incarnation. is_new_version() вернет true, если изменилась часть version. is_new_incarnation() вернет true, если изменилась любая из частей incarnation. Например, комбинацию этих методов можно использовать в пользовательском триггере, чтобы проверить, был ли перезапущен процесс или изменен участник:

swim = require('swim')s = swim.new()s:on_member_event(function(m, e)    if e:is_new_incarnation() then        if e:is_new_generation() then            -- Process restart.        end        if e:is_new_version() then            -- Process version update. It means            -- the member is somehow changed.        end    endend

swim_object:on_member_event(nil, old-trigger)

Удаление триггера on_member_event.

Параметры:

  • old-trigger (function) — старый триггер

Значение old-trigger должно быть значением, которое вернул swim_object:on_member_event(trigger-function[, ctx]).

swim_object:on_member_event(new-trigger, old-trigger[, ctx])

Это вариант метода swim_object:on_member_event(trigger-function[, ctx]).

Дополнительный параметр — old-trigger. Вместо добавления new-trigger в конец списка триггеров функция заменит в списке запись, совпадающую с old-trigger. Позиция в списке может быть важна, так как триггеры активируются последовательно, начиная с первого триггера в списке.

Значение old-trigger должно быть значением, которое вернул swim_object:on_member_event(trigger-function[, ctx]).

swim_object:on_member_event()

Возврат списка триггеров on_member_event.

Объект участника swim

Методы swim_object:member_by_uuid(), swim_object:self() и swim_object:pairs() возвращают объекты участников swim.

Объект участника swim имеет методы для чтения его атрибутов: status(), uuid(), uri(), incarnation(), payload_cdata(), payload_str(), payload(), is_dropped().

swim_member_object:status()

Возврат статуса, который может быть 'alive', 'suspected', 'left' или 'dead'.

Возвращает

string — 'alive' | 'suspected' | 'left' | 'dead'

swim_member_object:uuid()

Возврат UUID в виде cdata-структуры tt_uuid.

Возвращает

cdata struct tt_uuid — UUID

swim_member_object:uri()

Возврат URI в виде строки 'ip:port'. С помощью этого метода можно узнать реально назначенный порт, если в swim_object:cfg() был указан port = 0.

Возвращает

string — ip:port

swim_member_object:incarnation()

Возврат cdata-объекта с инкарнацией. Cdata-объект имеет два атрибута: incarnation().generation и incarnation().version.

Инкарнации можно сравнивать между собой с помощью любого оператора сравнения (==, <, >, <=, >=, ~=).

При выводе инкарнации отображаются как строки, содержащие и generation, и version.

Возвращает

cdata — инкарнация

swim_member_object:payload_cdata()

Возврат полезной нагрузки участника.

Возвращает

указатель на cdata — полезная нагрузка и размер в байтах

swim_member_object:payload_str()

Возврат полезной нагрузки в виде строкового объекта. Полезная нагрузка не декодируется. Она просто возвращается как строка вместо cdata. Если полезная нагрузка не была задана с помощью swim_object:set_payload() или swim_object:set_payload_raw(), то её размер равен 0 и возвращается nil.

Возвращает

строковый объект — полезная нагрузка или nil, если полезной нагрузки нет

swim_member_object:payload()

Поскольку модуль swim является Lua-модулем, для полезной нагрузки, скорее всего, будут использоваться Lua-объекты — таблицы, числа, строки и т.д. Естественно ожидать, что swim_member_object:payload() вернёт тот же объект, который был передан в swim_object:set_payload() другим экземпляром. swim_member_object:payload() пытается интерпретировать полезную нагрузку как MessagePack, а если это не удаётся, возвращает её как строку.

swim_member_object:payload() кэширует свой результат. Поэтому только первый вызов фактически декодирует полезную нагрузку cdata. Все последующие вызовы возвращают указатель на тот же результат, если только полезная нагрузка не изменена с новой инкарнацией. Если полезная нагрузка не задана (её размер равен 0), возвращается nil.

swim_member_object:is_dropped()

Возврат true, если этот объект участника является потерянной ссылкой на участника, который уже удален из таблицы участников.

Возвращает

boolean — true, если участник удален, иначе false

Пример:

tarantool> swim = require('swim')---...tarantool> s = swim.new({uri = 0, uuid = '00000000-0000-1000-8000-000000000001'})---...tarantool> self = s:self()---...tarantool> self:status()---- alive...tarantool> self:uuid()---- 00000000-0000-1000-8000-000000000001...tarantool> self:uri()---- 127.0.0.1:56367...tarantool> self:incarnation()---- - cdata {generation = 1569354463981551ULL, version = 1ULL}...tarantool> self:is_dropped()---- false...tarantool> s:set_payload_raw('123')---- true...tarantool> self:payload_cdata()---- 'cdata<const char *>: 0x0103500050'- 3...tarantool> self:payload_str()---- '123'...tarantool> s:set_payload({a = 100})---- true...tarantool> self:payload_cdata()---- 'cdata<const char *>: 0x0103500050'- 4...tarantool> self:payload_str()---- !!binary gaFhZA==...tarantool> self:payload()---- {'a': 100}...

Внутреннее устройство SWIM

Раздел о внутреннем устройстве SWIM не обязателен для изучения программистам, которые хотят использовать модуль SWIM, — он предназначен для программистов, которые хотят изменить или заменить модуль SWIM.

Сетевой протокол SWIM открыт, остаётся обратно совместимым при любых изменениях и может быть реализован пользователями, которые хотят эмулировать собственные узлы кластера SWIM, поскольку используют другой язык вместо Lua или другую среду, не связанную с Tarantool. Протокол кодируется в формате MsgPack.

SWIM packet structure:+-----------------Public data, not encrypted------------------+|                                                             ||      Initial vector, size depends on chosen algorithm.      ||                   Next data is encrypted.                   ||                                                             |+----------Meta section, handled by transport level-----------+| map {                                                       ||     0 = SWIM_META_TARANTOOL_VERSION: uint, Tarantool        ||                                      version ID,            ||     1 = SWIM_META_SRC_ADDRESS: uint, ip,                    ||     2 = SWIM_META_SRC_PORT: uint, port,                     ||     3 = SWIM_META_ROUTING: map {                            ||         0 = SWIM_ROUTE_SRC_ADDRESS: uint, ip,               ||         1 = SWIM_ROUTE_SRC_PORT: uint, port,                ||         2 = SWIM_ROUTE_DST_ADDRESS: uint, ip,               ||         3 = SWIM_ROUTE_DST_PORT: uint, port                 ||     }                                                       || }                                                           |+-------------------Protocol logic section--------------------+| map {                                                       ||     0 = SWIM_SRC_UUID: 16 byte UUID,                        ||                                                             ||                 AND                                         ||                                                             ||     2 = SWIM_FAILURE_DETECTION: map {                       ||         0 = SWIM_FD_MSG_TYPE: uint, enum swim_fd_msg_type,  ||         1 = SWIM_FD_GENERATION: uint,                       ||         2 = SWIM_FD_VERSION: uint                           ||     },                                                      ||                                                             ||               OR/AND                                        ||                                                             ||     3 = SWIM_DISSEMINATION: array [                         ||         map {                                               ||             0 = SWIM_MEMBER_STATUS: uint,                   ||                                     enum member_status,     ||             1 = SWIM_MEMBER_ADDRESS: uint, ip,              ||             2 = SWIM_MEMBER_PORT: uint, port,               ||             3 = SWIM_MEMBER_UUID: 16 byte UUID,             ||             4 = SWIM_MEMBER_GENERATION: uint,               ||             5 = SWIM_MEMBER_VERSION: uint,                  ||             6 = SWIM_MEMBER_PAYLOAD: bin                    ||         },                                                  ||         ...                                                 ||     ],                                                      ||                                                             ||               OR/AND                                        ||                                                             ||     1 = SWIM_ANTI_ENTROPY: array [                          ||         map {                                               ||             0 = SWIM_MEMBER_STATUS: uint,                   ||                                     enum member_status,     ||             1 = SWIM_MEMBER_ADDRESS: uint, ip,              ||             2 = SWIM_MEMBER_PORT: uint, port,               ||             3 = SWIM_MEMBER_UUID: 16 byte UUID,             ||             4 = SWIM_MEMBER_GENERATION: uint,               ||             5 = SWIM_MEMBER_VERSION: uint,                  ||             6 = SWIM_MEMBER_PAYLOAD: bin                    ||         },                                                  ||         ...                                                 ||     ],                                                      ||                                                             ||               OR/AND                                        ||                                                             ||     4 = SWIM_QUIT: map {                                    ||         0 = SWIM_QUIT_GENERATION: uint,                     ||         1 = SWIM_QUIT_VERSION: uint                         ||     }                                                       || }                                                           |+-------------------------------------------------------------+

Раздел Initial vector (вектор инициализации) присутствует только при включённом шифровании. Этот раздел содержит открытый ключ. Например, для алгоритмов AES это 16-байтовый вектор инициализации, сохранённый как есть. Если шифрование не используется, размер раздела равен 0.

Последующие разделы (Meta и Protocol Logic) шифруются как один большой блок данных, если включено шифрование.

Раздел Meta (метаданные) обрабатывает маршрутизацию и совместимость версий протокола. Он работает на «транспортном» уровне протокола SWIM и присутствует всегда. Ключи в разделе метаданных:

  • SWIM_META_TARANTOOL_VERSION — обязательное поле. Tarantool записывает сюда свою версию в виде 3-байтового целого числа:

    • 1 байт для major,
    • 1 байт для minor,
    • 1 байт для patch.

    Например, версия Tarantool 2.1.3 кодируется так: (((2 << 8) | 1) << 8) | 3;. Это поле будет использоваться для поддержки нескольких версий протокола.

  • SWIM_META_SRC_ADDRESS и SWIM_META_SRC_PORT — обязательные поля. IP-адрес и порт источника. IP кодируется как 4 байта: "xxx.xxx.xxx.xxx", где каждый 'xxx' — кодировка одного байта. Порт кодируется как целое число. Пример кодирования "127.0.0.1:3313":

    struct in_addr addr;inet_aton("127.0.0.1", &addr);pos = mp_encode_uint(pos, SWIM_META_SRC_ADDRESS);pos = mp_encode_uint(pos, addr->s_addr);pos = mp_encode_uint(pos, SWIM_META_SRC_PORT);pos = mp_encode_uint(pos, 3313);
  • Подраздел SWIM_META_ROUTING — необязательный. Отвечает за пересылку пакетов. Используется механизмом подозрений SWIM. Подробнее о подозрениях см. в документе SWIM.

    Если этот подраздел присутствует, следующие поля обязательны:

    • SWIM_ROUTE_SRC_ADDRESS и SWIM_ROUTE_SRC_PORT (IP-адрес и порт источника) (должны быть адресом инициатора сообщения (могут отличаться от SWIM_META_SRC_ADDRESS и от SWIM_META_SRC_ADDRESS_PORT);
    • SWIM_ROUTE_DST_ADDRESS и SWIM_ROUTE_DST_PORT (IP-адрес и порт назначения, конечный адресат сообщения).

    Если сообщение было отправлено опосредованно с помощью SWIM_META_ROUTING, ответ следует отправить обратно тем же маршрутом.

Пример того, как SWIM использует маршрутизацию для косвенных ping-запросов: предположим, есть 3 узла: S1, S2, S3. S1 отправляет сообщение S3 через S2. Для доставки сообщения выполняются следующие шаги:

S1 -> S2{ src: S1, routing: {src: S1, dst: S3}, body: ... }

S2 получает сообщение и видит, что routing.dst не равен S2, следовательно, это чужой пакет. S2 пересылает пакет на S3, сохраняя все данные, включая body и разделы routing.

S2 -> S3

S3 получает сообщение и видит, что routing.dst равен S3, следовательно, сообщение доставлено. Если S3 хочет ответить, он отправляет ответ через тот же прокси. S3 знает, что сообщение было доставлено от S2, поэтому отправляет ответ через S2.

Раздел Protocol logic (логика протокола) обрабатывает логические шаги и действия протокола SWIM.

  • SWIM_SRC_UUID — обязательное поле. SWIM использует UUID как уникальный идентификатор участника, а не IP/порт. В этом поле хранится UUID отправителя. Тип — MP_BIN. Размер всегда 16 байт. UUID кодируется в порядке байтов хоста, замена байтов не требуется.

После SWIM_SRC_UUID могут следовать четыре возможных подраздела: SWIM_FAILURE_DETECTION, SWIM_DISSEMINATION, SWIM_ANTI_ENTROPY, SWIM_QUIT. Может присутствовать любой или все эти подразделы. Реализация коннектора должна быть готова к обработке любой комбинации.

  • Подраздел SWIM_FAILURE_DETECTION — описывает ping или ACK. В подразделе SWIM_FAILURE_DETECTION содержатся:

    • SWIM_FD_MSG_TYPE (0 — ping, 1 — ack);
    • SWIM_FD_GENERATION + SWIM_FD_VERSION (инкарнация).
  • Подраздел SWIM_DISSEMINATION — список изменённых участников кластера. Может включать только часть изменённых участников кластера, если изменений слишком много для одного UDP-пакета.

    В подразделе SWIM_DISSEMINATION содержатся:

    • SWIM_MEMBER_STATUS (обязательное) (0 = alive, 1 = suspected, 2 = dead, 3 = left);
    • SWIM_MEMBER_ADDRESS и SWIM_MEMBER_PORT (обязательные) IP и порт участника;
    • SWIM_MEMBER_UUID (обязательное) (UUID участника);
    • SWIM_MEMBER_GENERATION + SWIM_MEMBER_VERSION (обязательные) (инкарнация участника);
    • SWIM_MEMBER_PAYLOAD (необязательное) (полезная нагрузка участника) (тип MessagePack — MP_BIN).

    Обратите внимание, что отсутствие SWIM_MEMBER_PAYLOAD ничего не означает — это не то же самое, что полезная нагрузка нулевого размера.

  • Подраздел SWIM_ANTI_ENTROPY — вспомогательный механизм для распространения (dissemination). Содержит те же поля, что и подраздел dissemination, но все они обязательны, включая полезную нагрузку, даже если её размер равен 0. Anti-entropy в конечном итоге распространяет изменения, которые по какой-либо причине не были распространены через dissemination.

  • Подраздел SWIM_QUIT — уведомление о том, что отправитель корректно покинул кластер, например с помощью swim_object:quit(), и его не следует считать «мёртвым». Статус отправителя должен быть изменён на 'left'.

    В подразделе SWIM_QUIT содержатся:

Инкарнация

Инкарнация (incarnation) — это 128-битное значение cdata, которое является частью конфигурации каждого участника и присутствует в большинстве сообщений. Оно состоит из двух частей: поколения (generation) и версии (version).

Поколение (generation) постоянно. По умолчанию оно равно числу микросекунд с начала эпохи (сравните со значением, возвращаемым clock_realtime64()). При необходимости пользователь может задать поколение при вызове new().

Версия (version) — изменчивое значение. Изначально равна 0. Она автоматически увеличивается при каждом изменении.

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

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

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

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

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