Conformance¶
Обвязка conformance доказывает, что брокер соблюдает контракт ядра. У неё две точки входа, и обе паникуют с внятным сообщением на первом же нарушении контракта:
harness::run_suiteпроверяет поверхность маршрутизации на вашем внутрипроцессном транспорте (том самомTestableBroker, который вы поставляете).harness::lifecycleпроверяет последовательность переходов жизненного цикла целиком на настоящем брокере.
Запускайте обе: run_suite - ради гарантий диспетчеризации, lifecycle - чтобы убедиться, что
new -> connect(self) -> подписка -> публикация -> ack -> shutdown(self) работает на реальном
транспорте.
Фича conformance включает testing, поэтому один-единственный TestableBroker, который
поставляет ваш крейт, годится и для run_suite здесь, и для обвязки
TestApp, которую пишут пользователи.
Набор проверок маршрутизации¶
harness::run_suite принимает синхронную фабрику (Fn() -> B), которая строит свежий
внутрипроцессный транспорт под каждый сценарий, поэтому ни один сценарий не видит состояние
другого. Каждый сценарий подключает брокер и работает с его подключённой формой - с вашим
TestableBroker, который реализует и Subscribe. Ниже дословно приведён прогон набора у
эталонного in-memory брокера; подставьте в фабрику конструктор своего транспорта:
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.
Контракт трейтов описан в разделе Как написать брокер.