Перейти к содержанию

Conformance

Нормативная версия документации - английская

Эта страница переведена с английского языковой моделью. При любом расхождении верен английский оригинал.

Обвязка conformance доказывает, что брокер соблюдает контракт ядра. У неё две точки входа, и обе паникуют с внятным сообщением на первом же нарушении контракта:

  • harness::run_suite проверяет поверхность маршрутизации на вашем внутрипроцессном транспорте (том самом TestableBroker, который вы поставляете).
  • harness::lifecycle проверяет последовательность переходов жизненного цикла целиком на настоящем брокере.

Запускайте обе: run_suite - ради гарантий диспетчеризации, lifecycle - чтобы убедиться, что new -> connect(self) -> подписка -> публикация -> ack -> shutdown(self) работает на реальном транспорте.

[dev-dependencies]
ruststream = { version = "0.7", features = ["conformance"] }

Фича conformance включает testing, поэтому один-единственный TestableBroker, который поставляет ваш крейт, годится и для run_suite здесь, и для обвязки TestApp, которую пишут пользователи.

Набор проверок маршрутизации

harness::run_suite принимает синхронную фабрику (Fn() -> B), которая строит свежий внутрипроцессный транспорт под каждый сценарий, поэтому ни один сценарий не видит состояние другого. Каждый сценарий подключает брокер и работает с его подключённой формой - с вашим TestableBroker, который реализует и Subscribe. Ниже дословно приведён прогон набора у эталонного in-memory брокера; подставьте в фабрику конструктор своего транспорта:

tests/conformance_self.rs
use ruststream::conformance::harness;

/// The suites that read a broker's publish log back - the routing contract's log check and the
/// seeking capability - need a broker that keeps one.
fn replaying() -> MemoryBroker<Retaining> {
    MemoryBroker::retaining(Retention::Messages(nonzero!(64)))
}

#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn memory_broker_passes_conformance_suite() {
    harness::run_suite(replaying).await;
}

Что именно он проверяет

Сценарий Что утверждается
порядок доставки сообщения доставляются в порядке публикации
публикация после подписки подписчик получает только сообщения, опубликованные после подписки; более ранние публикации не буферизуются
ack поглощает доставку подтверждённое сообщение не доставляется повторно
nack с requeue доставляет заново nack(requeue = true) доставляет сообщение ещё раз
nack без requeue отбрасывает после nack(requeue = false) повторной доставки нет
заголовки передаются заголовки сообщения приходят к подписчику без изменений
лог публикаций фиксирует публикации published(name) записывает каждое опубликованное сообщение

Транспорт, который не умеет подтверждать (ZeroMQ, MQTT QoS 0, Redis pub/sub, Core NATS), возвращает AckError::Unsupported из ack и nack, и набор принимает такой ответ везде, где сам завершает доставку. Отвечайте во внутрипроцессном транспорте ровно так же, как отвечает настоящий, а не заявляйте о подтверждении, которого в бою не происходит. Исключение одно - сценарий повторной доставки: nack(requeue = true) с ответом Unsupported этот сценарий завершает, потому что у транспорта, который ничего не забирает обратно, наблюдать нечего. Остальное проверяется по-прежнему, включая сценарий отбрасывания: доставка, которую никто не может завершить, всё равно не должна прийти снова.

Ответ набор читает у доставки, а не у брокера, поэтому он различается по подпискам и сообщениям так же, как различается у транспорта: повторная постановка в очередь бывает совещательной при одном режиме коммитов и перематывающей при другом, подтверждение - доступным на одном уровне качества обслуживания и недоступным на другом. Неизменным остаётся смысл успеха: Ok(()) из nack(requeue = true) обещает, что сообщение вернётся, и путь повторов в рантайме читает этот ответ так же.

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

У каждой совместимости есть собственный набор, описанный ниже. Наборы для реализованных совместимостей вызывает ваш крейт. Среди них capabilities::batches - единственная проверка того, что пакет не приходит больше размера, с которым его открыли. Крейт, который запустил только run_suite, свои пакеты не проверил.

Проверка жизненного цикла

harness::lifecycle проходит последовательность переходов жизненного цикла на настоящем Broker: синхронное конструирование без I/O, затем connect, который поглощает self и даёт типизированную подключённую форму, подписка через собственный SubscriptionSource брокера, публикация, которую подписка получает и подтверждает, и shutdown, который поглощает self и даёт терминального свидетеля.

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

Проверка принимает три фабрики, и они делают её независимой от брокера:

use ruststream::conformance::harness;
use ruststream_nats::{NatsBroker, SubscribeOptions};

#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[ignore = "needs a running nats-server; set NATS_TEST_URL"]
async fn passes_lifecycle() {
    let url = std::env::var("NATS_TEST_URL").unwrap();
    harness::lifecycle(
        || NatsBroker::new(url.clone()), // sync construction (no I/O)
        |subject| SubscribeOptions::new(subject), // the broker's SubscriptionSource
        |connected| connected.publisher(), // a publisher from the connected form
    )
    .await;
}
  • make_broker - синхронная (Fn() -> B). Брокер, который можно построить только асинхронно, ей не удовлетворяет: конструируйте дёшево, а подключайтесь в Broker::connect.
  • make_source строит дескриптор подписки на субъект (путь макро-подписчика).
  • make_publisher создаёт издателя из подключённой формы.

Брокер без семантики ack (Core NATS) проходит проверку, возвращая из ack значение AckError::Unsupported: проверка принимает и его, и успешный ack. lifecycle выполняет настоящий connect, поэтому запускайте её против живого сервера и включайте только при заданной переменной окружения вроде NATS_TEST_URL; in-memory брокер проходит её прямо в процессе.

Наборы проверок для совместимостей

Если ваш брокер реализует трейт-совместимость, запустите соответствующий набор из conformance::capabilities: он доказывает, что реализация соблюдает контракт трейта. Брокеры без этой совместимости его не вызывают. Каждый набор принимает фабрики той же формы, что и lifecycle, и выполняет настоящий connect, поэтому включайте его по той же переменной окружения:

Набор Требует Что утверждается
capabilities::request_reply RequestReply запрос доходит до отвечающей стороны с пригодным заголовком reply-to, соотнесённый ответ завершает запрос, а запрос без ответа возвращает ошибку по истечении своего таймаута
capabilities::batches BatchSubscriber каждое опубликованное сообщение приходит в порядке публикации, распределённое по непустым пакетам
capabilities::transactions TransactionalPublisher ничто внутри транзакции не видно до commit, commit публикует буфер по порядку, abort его отбрасывает; неверный вызов возвращает ошибку - commit / abort без открытой транзакции, второй begin_transaction при уже открытой (он обязан оставить её нетронутой)
capabilities::owned_transactions OwnedTransactions и его Transaction ничто опубликованное в открытую транзакцию не видно до commit, commit доставляет весь буфер в порядке публикации, abort его отбрасывает, две транзакции, открытые одновременно на одном издателе, завершаются независимо, а сам издатель продолжает публиковать напрямую, пока одна из них открыта
capabilities::seeking Seekable, сообщения Positioned перемотка назад на позицию, снятую с доставленного сообщения, повторно доставляет именно это сообщение и всё, что идёт за ним, по порядку; перемотка вперёд пропускает доставки, стоявшие в очереди до цели; после перемотки подписка продолжает доставлять новые публикации
use ruststream::conformance::capabilities;
use ruststream_nats::{NatsBroker, SubscribeOptions};

#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[ignore = "needs a running nats-server; set NATS_TEST_URL"]
async fn passes_request_reply() {
    let url = std::env::var("NATS_TEST_URL").unwrap();
    capabilities::request_reply(
        || NatsBroker::new(url.clone()),
        |subject| SubscribeOptions::new(subject),
        |connected| connected.publisher(), // the RequestReply publisher under test
        |connected| connected.publisher(), // the plain publisher the responder replies through
    )
    .await;
}

Каждый набор называет субъект заново на каждом прогоне, поэтому прогон читает только то, что опубликовал сам. Фиксированный субъект проходит на свежем сервере и не проходит на втором прогоне у любого брокера, который хранит оставленное первым: удерживаемый лог воспроизведёт в одну подписку оба прогона, durable-очередь всё ещё держит ранние сообщения, пространство ключей - ранний тип. Для этого наборы вызывают conformance::helpers::unique_subject; у вашего сквозного набора тестов та же проблема и то же решение.

In-memory брокер реализует все совместимости нативно и проходит все пять наборов прямо в процессе (см. Memory); он и есть исполняемый эталон того, чего ждёт каждый набор.

Чеклист автора

Прежде чем публиковать крейт брокера:

  • [ ] Реализованы Broker, ConnectedBroker, Subscribe (или SubscriptionSource), Subscriber, IncomingMessage, Publisher и PublishPolicy, которая его конструирует.
  • [ ] shutdown выполняет всё освобождение ресурсов, которое может вернуть ошибку, и никогда не блокирует и не паникует.
  • [ ] Ack поглощает self; nack учитывает флаг requeue.
  • [ ] Крейт владеет своим Config; поля без разумного умолчания не получают Default.
  • [ ] Трейты-совместимости реализованы только там, где брокер их действительно поддерживает, и каждая реализованная совместимость проходит свой набор из conformance::capabilities.
  • [ ] Внутрипроцессный транспорт, реализующий TestableBroker на подключённой форме, поставляется под фичей testing (только маршрутизация ядра) и зарегистрирован через register_testable_broker!.
  • [ ] harness::run_suite проходит (поверхность маршрутизации).
  • [ ] harness::lifecycle проходит против настоящего сервера и включается по переменной окружения (та самая последовательность: синхронный new, поглощающий connect, подписка, ack, поглощающий shutdown и ошибка разделяемого дескриптора после него).
  • [ ] Сквозной набор тестов покрывает специфичную для брокера семантику и включается по той же переменной.
  • [ ] Метаданные Cargo.toml заполнены полностью (description, license, repository, keywords, categories), а CI проверяет --no-default-features и --all-features.

Контракт трейтов описан в разделе Как написать брокер.