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

Как написать брокер

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

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

Брокер - это самостоятельный крейт, который реализует трейты ядра. Он зависит от ruststream с выключенными фичами по умолчанию, поэтому получает поверхность трейтов и рантайм без встроенного JSON-кодека и без других брокеров:

[dependencies]
ruststream = { version = "0.7", default-features = false }

Эта страница и есть контракт. Реализуйте обязательные трейты, заведите собственный Config, добавьте трейты-совместимости под то, что ваш брокер умеет, и подтвердите результат обвязкой conformance. Полная реализация поверх настоящего клиента разобрана в примере с NATS.

Обязательные трейты

Broker и ConnectedBroker

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

pub trait Broker: Send + Sync + Sized {
    type Error: std::error::Error + Send + Sync + 'static;
    type Connected: ConnectedBroker;
    async fn connect(self) -> Result<Self::Connected, Self::Error>;
}

pub trait ConnectedBroker: Send + Sync + Sized + 'static {
    type Error: std::error::Error + Send + Sync + 'static;
    type Closed: Send;
    async fn shutdown(self) -> Result<Self::Closed, Self::Error>;
}

В shutdown нельзя блокировать и паниковать: всё освобождение ресурсов, которое может вернуть ошибку, делайте здесь и возвращайте Result. Closed - свидетель остановки: сложите в него диагностику (результаты сброса буферов, счётчики потерь) как обычные данные или возьмите ().

Конструирование синхронно и без I/O: new(addrs) только записывает конфигурацию. Вся сетевая работа происходит в connect, который рантайм вызывает один раз на старте. Подключённая форма держит живого клиента напрямую, поэтому её операции не проверяют состояние «вроде бы подключены».

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

Обвязка conformance проверяет всю последовательность переходов, а пример с NATS проходит её на настоящем клиенте.

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

Всю последовательность переходов in-memory брокер проходит за несколько строк, и каждый его пример на этой странице вырезан из того же файла, поэтому код страницы меняется вместе с контрактом:

src/memory/mod.rs
impl<Log: LogMode> Broker for MemoryBroker<Log> {
    type Error = MemoryError;
    type Connected = ConnectedMemoryBroker<Log>;

    /// Connecting is free for an in-process bus. A shut-down bus (a clone lineage may have shut
    /// the shared state down) is revived with a fresh, empty registration map, so the connected
    /// form always starts live; a live bus keeps its registrations.
    fn connect(self) -> impl Future<Output = Result<Self::Connected, Self::Error>> {
        {
            let mut bus = self
                .state
                .subscribers
                .lock()
                .expect("memory broker mutex poisoned");
            if matches!(*bus, Bus::ShutDown) {
                *bus = Bus::Live(HashMap::new());
            }
        }
        ready(Ok(ConnectedMemoryBroker {
            state: self.state,
            mode: PhantomData,
        }))
    }
}

/// The connected form of [`MemoryBroker`]: the typed witness that [`Broker::connect`] ran.
///
/// Cheap to clone: the in-memory bus is shared state by nature, so the connected form is a
/// shareable handle on it, exactly like the unconnected broker. Subscriptions (the
/// [`Subscribe`] capability, [`MemorySource`]) resolve against this form, and carry over the
/// broker's log mode: only a [`Retaining`] one opens repositionable subscriptions.
pub struct ConnectedMemoryBroker<Log = Discarding> {
    state: Arc<MemoryState>,
    mode: PhantomData<Log>,
}

impl<Log> Clone for ConnectedMemoryBroker<Log> {
    fn clone(&self) -> Self {
        Self {
            state: Arc::clone(&self.state),
            mode: PhantomData,
        }
    }
}

impl<Log> fmt::Debug for ConnectedMemoryBroker<Log> {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        f.debug_struct("ConnectedMemoryBroker")
            .finish_non_exhaustive()
    }
}

impl<Log: LogMode> ConnectedMemoryBroker<Log> {
    /// Returns a publisher bound to this broker.
    #[must_use]
    pub fn publisher(&self) -> MemoryPublisher {
        MemoryPublisher {
            state: Arc::clone(&self.state),
            txn: Mutex::new(None),
        }
    }

    /// Returns a request / reply-capable publisher bound to this broker.
    ///
    /// See [`MemoryBroker::requester`] for why its operations report [`RequestError`] rather
    /// than [`MemoryError`].
    #[must_use]
    pub fn requester(&self) -> MemoryRequester {
        MemoryRequester::new(Arc::clone(&self.state))
    }
}

impl<Log: LogMode> ConnectedBroker for ConnectedMemoryBroker<Log> {
    type Error = MemoryError;
    type Closed = ClosedMemoryBroker;

    /// Enters the terminal shut-down state: the bus itself flips to its `ShutDown` variant, so
    /// every aliased handle that would touch it (a publisher's publish or transaction commit, a
    /// request) errors with [`MemoryError::ShutDown`]. Consuming `self` makes any further use
    /// of this handle a compile error; the returned witness reports how many subscriber
    /// registrations the teardown dropped.
    fn shutdown(self) -> impl Future<Output = Result<Self::Closed, Self::Error>> {
        let dropped = {
            let mut bus = self
                .state
                .subscribers
                .lock()
                .expect("memory broker mutex poisoned");
            match std::mem::replace(&mut *bus, Bus::ShutDown) {
                Bus::Live(subscribers) => subscribers.values().map(Vec::len).sum(),
                Bus::ShutDown => 0,
            }
        };
        ready(Ok(ClosedMemoryBroker {
            subscribers_dropped: dropped,
        }))
    }
}

ClosedMemoryBroker - это тот самый свидетель с диагностикой остановки, о котором сказано выше: он сообщает, сколько регистраций подписчиков отбросила остановка.

Subscribe

Реализуйте Subscribe на подключённой форме, чтобы подписаться можно было по имени темы, субъекта или очереди. Через него подписывается #[subscriber("name")].

pub trait Subscribe: ConnectedBroker {
    type Subscriber: Subscriber;

    // Кто публикует копии для повторов подписки по имени и кто называет их адрес.
    // AddressedCopies там, где публикация по имени подписки доходит до открытой под
    // ним подписки (субъект, топик, стрим, имя очереди), - тогда имя и есть адрес.
    // NamedCopies там, где это не так: фильтр MQTT читает много топиков и не
    // называет ни одного.
    type Copies: CopyPath;

    async fn subscribe(&self, name: &str) -> Result<Self::Subscriber, Self::Error>;

    // По умолчанию: что регистрация, смонтированная по этому имени, объявила о своих
    // повторах. Отобразите предел и адрес на подписку, которую открывает это имя,
    // если у брокера есть для них механизм.
    fn declare_retry(&self, name: &str, declaration: &RetryDeclaration)
        -> Result<(), DeclareRetryError>;
}

Всё, что требуется, - открыть подписку и сказать, по какому адресу до неё доходит публикация:

src/memory/mod.rs
impl<Log: LogMode> Subscribe for ConnectedMemoryBroker<Log> {
    type Subscriber = MemorySubscriber<Log>;
    // One subject is both ends of the bus here, so a publish under the name a subscription reads
    // reaches that subscription, and the name is the address.
    type Copies = AddressedCopies;

    fn subscribe(&self, name: &str) -> impl Future<Output = Result<Self::Subscriber, Self::Error>> {
        let (tx, rx) = mpsc::unbounded_channel();
        let name = name.to_owned();
        if let Err(err) = self.state.register(&name, tx.clone()) {
            return ready(Err(err));
        }
        ready(Ok(MemorySubscriber {
            name,
            rx,
            requeue: tx,
            state: Arc::clone(&self.state),
            seek: Arc::new(SeekControl::default()),
            mode: PhantomData,
        }))
    }
}

type Copies - это то, что источник по имени сообщает на вашем брокере, и от него зависит, как повторяет #[subscriber("orders")]: при AddressedCopies адресом служит само имя и писать больше нечего, при NamedCopies адрес копий называет место монтирования.

Отвечайте NamedCopies там, где имя подписки не служит адресом публикации. Так устроен фильтр топиков MQTT: devices/+/telemetry читает топик каждого устройства и не называет ни одного, поэтому адрес копии называет место монтирования - через .out_retry(политика).to(имя) или через преобразование публикации, которое называет его для каждой доставки. Подписке, которой одного имени мало, нужен ваш собственный дескриптор.

Subscriber

Подписчик - это Stream входящих сообщений. Обратное давление даёт сам Stream.

pub trait Subscriber: Send {
    type Message: IncomingMessage;
    type Error: std::error::Error + Send + Sync + 'static;
    fn stream(&mut self) -> impl Stream<Item = Result<Self::Message, Self::Error>> + Send + '_;
}

stream берёт &mut self, поэтому любое состояние, буферизованное между опросами, хранится за изменяемым заимствованием - это и даёт безопасность при отмене.

IncomingMessage

Доставленное сообщение отдаёт полезную нагрузку и заголовки; его подтверждают через ack или отклоняют через nack. ack поглощает self, поэтому двойной ack - ошибка компиляции.

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

pub trait IncomingMessage: Send + Sync {
    fn payload(&self) -> &[u8];
    fn headers(&self) -> &HeaderMap;
    async fn ack(self) -> Result<(), AckError>;
    async fn nack(self, requeue: bool) -> Result<(), AckError>;

    // Defaulted: false. The runtime reads this first and never calls
    // nack_after without it, so override the pair together.
    fn supports_nack_after(&self) -> bool;

    // Defaulted: AckError::Unsupported. Override when the transport has native
    // delayed redelivery (JetStream NAK with delay); handlers reach it through
    // HandlerOutcome::retry_after.
    async fn nack_after(self, delay: Duration) -> Result<(), AckError>;

    // Defaulted: None. Override (with the Partitioned capability) to feed the
    // runtime's keyed worker lanes, workers(n, by_key).
    fn partition_key(&self) -> Option<&[u8]>;

    // Defaulted: None. Override where the transport counts its own deliveries
    // (JetStream num_delivered, SQS ApproximateReceiveCount, Pub/Sub
    // delivery_attempt): a registration's max_attempts(..) cap then counts the
    // broker's redeliveries and not only the copies the runtime published.
    // The first delivery of a message answers 1.
    fn redelivery_count(&self) -> Option<u64>;
}

Отложенная повторная доставка - это два метода, и рантайм спрашивает supports_nack_after. Если переопределить только nack_after, флаг останется в false и переопределение никто не вызовет. По умолчанию nack_after возвращает AckError::Unsupported, а не завершает доставку обычным nack(true): транспорт, который не умеет придержать сообщение, обязан сказать об этом, иначе пауза перед повтором превращается в шторм повторных доставок.

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

Пока redelivery_count возвращает None, этот заголовок - единственный счётчик, и предел max_attempts(..) читается по нему. Переопределите метод - и счётчик брокера станет единственным, который читает предел: фреймворковый заголовок к нему не прибавляется.

Если задержку вы соблюдаете, публикуя копию сами - очередь ожидания за обменником мёртвых писем, отдельный топик повторов, - увеличивайте на этой копии RETRY_COUNT_HEADER (он экспортируется из ruststream::runtime), но только там, где транспорт не считает ничего. Там, где он считает, ваша копия - новое сообщение, и брокер начинает считать его заново, так что доставка, прошедшая круг по очереди ожидания, приходит к пределу первой попыткой. Это поведение вашей схемы задержки: опишите его в документации крейта, а заголовок не трогайте.

Показать «не переопределено ничего» не на ком: в этом рабочем пространстве эти методы переопределяет каждый брокер. Поэтому поведение закреплено тестом в ядре:

src/message.rs
struct Stub {
    payload: Vec<u8>,
    headers: HeaderMap,
}

impl IncomingMessage for Stub {
    fn payload(&self) -> &[u8] {
        &self.payload
    }

    fn headers(&self) -> &HeaderMap {
        &self.headers
    }

    fn ack(self) -> impl Future<Output = Result<(), AckError>> {
        ready(Ok(()))
    }

    fn nack(self, _requeue: bool) -> impl Future<Output = Result<(), AckError>> {
        ready(Ok(()))
    }
}

let stub = Stub {
    payload: b"body".to_vec(),
    headers: HeaderMap::new(),
};
assert_eq!(stub.payload(), b"body");
// The default partition_key is None (no key).
assert!(stub.partition_key().is_none());
// The default reports no native delayed redelivery, so the runtime uses its fallback.
assert!(!stub.supports_nack_after());
// The default nack_after signals "not honored" rather than silently degrading.
assert!(matches!(
    stub.nack_after(Duration::from_secs(1)).await,
    Err(AckError::Unsupported)
));

Именно ответ Unsupported позволяет рантайму отличить транспорт без отложенной доставки от соблюдённой задержки и включить собственный запасной путь.

Publisher

pub trait Publisher: Send + Sync {
    /// Как ваш транспорт распоряжается нагрузкой: `Lend`, если он её читает, `Take`, если ваш
    /// клиент оставляет её себе.
    type Payload: PayloadForm;

    type Error: std::error::Error + Send + Sync + 'static;

    /// Настройки вашего брокера на одно сообщение. Все поля необязательные; `()`, если их нет.
    type Options: Clone + Send + Sync + 'static;

    async fn publish(
        &self,
        msg: OutgoingFor<'_, Self::Payload>,
        options: Option<&Self::Options>,
    ) -> Result<(), Self::Error>;

    /// С реализацией по умолчанию: заголовки, которые этот издатель кладёт под каждую публикацию.
    fn base_headers(&self) -> Option<&HeaderMap> { None }
}

Объявите Take, если клиент держит нагрузку после вызова: он принимает Vec<u8>, Bytes или другое собственное значение. Объявите Lend, если транспорт только читает байты, потому что складывает их в свой кадр, пакет или буфер сокета. Почти всякий протокол на проводе относится ко второму виду.

Сообщение, которое вы получаете, следует за объявлением, а вместе с ним и работа фреймворка над вами:

// Take: буфер, который записал фреймворк, и он ваш.
async fn publish(&self, msg: OutgoingMessage<'_, BytesMut>, options: Option<&()>) -> Result<(), Error>

// Lend: байты там, где они уже лежат, годные на время вызова.
async fn publish(&self, msg: OutgoingMessage<'_, &[u8]>, options: Option<&()>) -> Result<(), Error>

Издателю с Take достаётся буфер кодека таким, каким его записали. Vec::from(payload) бесплатен, потому что этот буфер и есть вектор; payload.freeze() стоит один блок, который делает владение разделяемым; копируется по дороге только та публикация, которая одалживала чужие байты.

Издателю с Lend достаётся &[u8] и ничего, что нужно освобождать. Внутри цикла диспетчеризации фреймворк одалживает ему один буфер на цикл и переписывает его для каждого сообщения, поэтому ответ через ваш брокер не выделяет памяти ни разу за доставку. Ради этого объявление и существует: ядро не может одолжить то, что вы можете оставить себе.

OutgoingMessage заимствует имя, а нагрузку и ассоциативный массив заголовков ваш транспорт может забрать себе, а не копировать. Прочитать их можно через msg.payload() и msg.headers(), которые отдают &[u8] и &HeaderMap при любой объявленной форме. Транспорт, который поглощает сообщение, забирает его части через msg.into_parts(): адрес, нагрузку и заголовки одним перемещением, без копирования.

Сервис пишет не этот метод, а билдер: publisher.message(&value).publish() выбирает адрес, кодек и заголовки и делает ровно один вызов publish. Реализуйте publish, и весь билдер заработает поверх него.

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

options равно None на всех путях, где настройки некому изменить: ответ обработчика, отложенная повторная доставка. Тогда действуют настройки политики.

Clone и 'static нужны тестовой обвязке: она снимает копию настроек публикации через слот Out и отдаёт их тесту этим же типом. Поэтому сервис, который тестирует ваш брокер, проверяет значение, полученное вашим publish, а не поле протокола, в которое оно превратилось. Выведите ещё Debug и PartialEq, и проверка запишется как with_options(&YourOptions { .. }) (проверки по слотам Out).

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

Transaction называет собственные Options и Payload и объявляет тот же base_headers с реализацией по умолчанию. Транзакция - отдельная поверхность публикации, поэтому она может учитывать настройки, которых не учитывает издатель, от которого она открыта, а клиентский буфер оставляет нагрузку себе там, где прямая публикация её только читает. У большинства брокеров в обоих местах стоят типы самого издателя. Дескриптор, у которого своей постоянной величины нет, в обоих местах оставляет base_headers с реализацией по умолчанию.

PublishPolicy

Издатель брокера - это политика (обменник, таймаут очереди, транзакционный идентификатор) и живое соединение. Поставьте отдельный тип политики: он конструируется где угодно и держит опции билдера.

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

pub trait PublishPolicy<C: ConnectedBroker> {
    type Live; // the live publisher (or live wiring form, for combinator stacks)
    async fn pair(self, connected: &C) -> Result<Self::Live, PairError>;
}

Ошибка - PairError со стёртым типом: оберните в него ошибку своего брокера конструктором PairError::new. Политика инстанцирует издателя один раз, на старте, поэтому на горячий путь pair не попадает.

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

Минимальный эталон - MemoryPublish и MemoryRequest in-memory брокера: опций у них нет, поэтому это пустые структуры.

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

Если обычная политика годится со своими умолчаниями (а так почти всегда), реализуйте на подключённой форме ещё и DefaultPublish и назовите её там. Тогда рантайм сам инстанцирует издателя ответа, когда обработчик с publish("dest") монтируется без явного .out_reply(..), и b.include(def) компилируется сам по себе. Брокеры, издателям которых всегда нужны явные опции, DefaultPublish не реализуют, и их пользователи указывают политику при каждой регистрации обработчика.

pub trait DefaultPublish: ConnectedBroker {
    type Policy: PublishPolicy<Self> + Default + Send + 'static;
}

Обе половины - на брокере, у которого политика не держит вообще никаких опций:

src/memory/mod.rs
impl<Log: LogMode> PublishPolicy<ConnectedMemoryBroker<Log>> for MemoryPublish {
    type Live = MemoryPublisher;

    fn pair(
        self,
        connected: &ConnectedMemoryBroker<Log>,
    ) -> impl Future<Output = Result<Self::Live, PairError>> {
        ready(Ok(connected.publisher()))
    }
}

impl<Log: LogMode> DefaultPublish for ConnectedMemoryBroker<Log> {
    type Policy = MemoryPublish;
}

Источники подписки

Subscribe закрывает случай, когда подписке хватает имени. Если ей нужны опции вашего брокера (группа консьюмеров, durable-имя, политика доставки), заведите тип-дескриптор, который реализует SubscriptionSource:

pub trait SubscriptionSource<C: ConnectedBroker> {
    type Subscriber: Subscriber;

    // Кто публикует копии, из которых состоят повторы этой подписки, и кто называет
    // их адрес: AddressedCopies, NamedCopies или BrokerMoves.
    type Copies: CopyPath;

    fn name(&self) -> &str;
    fn subscribe(self, connected: &C) -> impl Future<Output = Result<Self::Subscriber, C::Error>> + Send;

    // По умолчанию дескриптор без изменений. Прочитайте здесь max_attempts(..) и
    // dead_letter(..) регистрации и примените их к подписке, которую собираетесь
    // открыть, если у брокера есть для них механизм.
    fn declare_retry(self, declaration: &RetryDeclaration) -> Self;

}

// Вторая половина AddressedCopies: адрес - свойство типа, а не ответ, который
// проверяют на старте. Спросите брокер там, где это знает только живое соединение.
pub trait RedeliveryAddressed<C>: SubscriptionSource<C, Copies = AddressedCopies> {
    async fn redelivery_address(&self, connected: &C) -> Result<RedeliveryAddress, C::Error>;
}

Дайте дескриптору ассоциированный конструктор (OrdersStream::new(..)), а не свободную функцию: тогда пользователь называет его прямо в атрибуте, #[subscriber(OrdersStream::new("orders", "workers"))].

Макрос читает тип из вызова конструктора и принимает на нём ещё и цепочку билдера (#[subscriber(OrdersStream::new("orders").durable("workers"))]), пока каждый метод возвращает Self.

type Subscriber объявлен на источнике, поэтому один брокер может предложить несколько видов подписки (pub/sub и стримы) с разными типами подписчиков или, как в примере с NATS, обслуживать их все одним дескриптором, который ветвится внутри.

Выведите на дескрипторе Clone: монтирование пересобирает конфигурацию на каждую регистрацию, поэтому одно определение можно смонтировать сразу на два брокера.

Кто публикует копию для повтора

Каждый дескриптор объявляет через type Copies одно из трёх.

AddressedCopies означает, что копии для повтора публикует этот процесс, и дескриптор знает их адрес. Рядом с SubscriptionSource он реализует RedeliveryAddressed, поэтому адрес - свойство типа, а не ответ, который проверяют на старте. Это ответ для субъекта, топика, стрима и очереди: одна подписка, один адрес, куда сервис умеет опубликовать обратно.

NamedCopies означает, что копии публикует этот процесс, но дескриптор не может их адресовать: субъект с шаблоном, фильтр MQTT, паттерн Pulsar, список топиков. Такая подписка читает много адресов, поэтому один называет место монтирования - статически через .out_retry(политика).to(имя) или для каждой доставки через преобразование публикации.

BrokerMoves означает, что доставку перекладывает сам сервер или клиентская библиотека: кворумная очередь с x-delivery-limit и x-dead-letter-exchange, подписка Pub/Sub с политикой мёртвых писем, политика redrive в SQS, консьюмер Pulsar с DeadLetterPolicy. Сервис тогда не публикует ничего, поэтому .out_retry(..) в любом месте монтирования этого дескриптора - ошибка компиляции, и ошибка называет сам дескриптор.

Первые два на каждую регистрацию связывают издателя повторов из политики DefaultPublish брокера, поэтому дескриптор, объявивший любой из них над брокером без DefaultPublish, не компилируется.

Subscribe объявляет то же самое для формы по имени: type Copies там - это то, что сообщает #[subscriber("orders")] на вашем брокере. Отвечайте AddressedCopies там, где публикация по имени подписки доходит до открытой под ним подписки - так обычно устроены субъект, топик, стрим и имя очереди, - и тогда имя и есть адрес, а писать больше нечего.

Там, где нативность зависит от значения поля, а не от типа - очередь RabbitMQ без .delay(..) не имеет своей отложенной доставки, - оставляйте путь открытым.

Что объявляет регистрация

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

Регистрация, смонтированная по голому имени, объявляет то же самое, но дескриптор ей достаётся ядерный - Name, у которого нет своей топологии. Объявление вместо этого принимает ваш брокер, в Subscribe::declare_retry: в той же точке и для того имени, которое он собирается открыть. Отобразите его там так же, как это делает ваш дескриптор, и только когда объявлены обе половины.

Реализация по умолчанию принимает регистрацию, которая не объявила ничего. Она принимает и любое объявление там, где ваш type Copies говорит, что копии публикует этот процесс: предел и адрес применит рантайм. На брокере с BrokerMoves она отклоняет непустое объявление на старте, называя подписку и место, которому объявление принадлежит: применить его больше некому, а сообщение с молча пропавшим пределом переживёт собственную политику мёртвых писем. Реализуйте метод там, где у брокера есть механизм, до которого голое имя дотягивается: политика мёртвых писем в Pub/Sub, политика redrive в SQS, DeadLetterPolicy у консьюмера Pulsar. В остальных случаях не трогайте его.

Куда публикуется копия для повтора

Без нативной отложенной доставки рантайм выполняет retry_after сам: по истечении задержки он публикует копию сообщения. Адрес для этой копии называет дескриптор с AddressedCopies.

src/memory/mod.rs
impl<Log: LogMode> SubscriptionSource<ConnectedMemoryBroker<Log>> for MemorySource {
    type Subscriber = MemorySubscriber<Log>;
    // The bus moves nothing on its own, and one subject is both ends of it, so a copy goes back
    // to the subject the subscription reads.
    type Copies = AddressedCopies;

    fn name(&self) -> &str {
        &self.name
    }

    async fn subscribe(
        self,
        connected: &ConnectedMemoryBroker<Log>,
    ) -> Result<Self::Subscriber, MemoryError> {
        Subscribe::subscribe(connected, &self.name).await
    }
}

impl<Log: LogMode> RedeliveryAddressed<ConnectedMemoryBroker<Log>> for MemorySource {
    fn redelivery_address(
        &self,
        _connected: &ConnectedMemoryBroker<Log>,
    ) -> impl Future<Output = Result<RedeliveryAddress, MemoryError>> + Send {
        // One subject is both ends of the bus, and no lookup is needed to say so.
        ready(Ok(RedeliveryAddress::new(self.name.clone())))
    }
}

Отвечайте тем именем, по которому публикация вашего брокера снова доходит до этой подписки: субъект в NATS, топик в Kafka, ключ стрима в Redis.

В NATS JetStream консьюмер привязан к стриму, а не к субъекту, поэтому ответ - субъект, на который этот стрим публикуют, и дескриптор, построенный по имени стрима, спрашивает его у сервера. Рантайм спрашивает один раз, на старте, а .to(имя) в месте монтирования этот ответ переопределяет.

.to(имя) называет канал, в который публикует эта регистрация, поэтому в сгенерированном документе он появляется вместе с операцией send. Адрес, которым отвечаете вы, в документе не появляется: это канал самой подписки, и он в документе уже есть.

harness::redelivery_address проверяет ваш ответ: публикация по названному адресу обязана прийти в ту подписку, которая его назвала. У дескриптора с NamedCopies проверять нечего, а остальную лестницу для обоих проходит harness::lifecycle.

Как назвать вид одной строкой

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

src/memory/mod.rs
impl FromName for MemorySource {
    fn from_name(name: impl Into<Cow<'static, str>>) -> Self {
        Self::new(name.into().into_owned())
    }
}

Тогда #[subscriber(OrdersStream)] законен: атрибут называет вид, а значение подставляет точка монтирования. Вид, которому нужно больше одного имени (тема и имя подписки), FromName не реализует, и такая форма для него не компилируется.

Настройки в терминах вашего брокера

Ядро не знает, что у подписки есть стрим, durable-имя или группа консьюмеров, поэтому даёт один хук: map_source - преобразование над источником, который собирает точка монтирования. Свой трейт вы кладёте сверху и привязываете к своему типу источника:

use ruststream::runtime::{Declared, SubscriberBuilder, SubscriberSettings};

pub trait NatsSubscriber {
    fn jetstream(self, stream: impl Into<String>) -> Self;
    fn durable(self, name: impl Into<String>) -> Self;
}

// Четыре слота состояния - это (воркеры, политики отказа, стартовая позиция, размер пакета);
// `Codec` - собственное переопределение кодека регистрации, `()` пока его никто не назвал. Оба
// передаются дальше без изменений.
impl<Def, Workers, Failures, StartPosition, Batch, Codec> NatsSubscriber
    for SubscriberBuilder<Def, SubscribeOptions, (Workers, Failures, StartPosition, Batch), Codec>
where
    Def: Declared,
{
    fn jetstream(self, stream: impl Into<String>) -> Self {
        self.map_source(|source| source.jetstream(stream))
    }

    fn durable(self, name: impl Into<String>) -> Self {
        self.map_source(|source| source.durable(name))
    }
}

Ограничение по типу источника означает, что на билдере другого брокера этих методов не существует. Тем же расширением ниже пользуется словарь слотов Out.

Одна настройка ядра меняет не слот состояния, а сам тип источника: start_at(..) оборачивает дескриптор в StartAt<SubscribeOptions, Position>. На подписках со стартовой позицией ваши методы выходят из области видимости, поэтому закройте этот случай вторым impl над обёрнутым источником. StartAt::map_inner берёт дескриптор изнутри обёртки и возвращает позицию нетронутой, так что каждый метод остаётся однострочным:

use ruststream::StartAt;
use ruststream::runtime::Fixed;

// Слот стартовой позиции здесь по построению равен `Fixed` - обёртку породил именно
// `start_at(..)`, - да и тип источника другой, поэтому этот impl и предыдущий не пересекаются.
impl<Def, Workers, Failures, Batch, Codec, Position> NatsSubscriber
    for SubscriberBuilder<
        Def,
        StartAt<SubscribeOptions, Position>,
        (Workers, Failures, Fixed, Batch),
        Codec,
    >
where
    Def: Declared,
{
    fn jetstream(self, stream: impl Into<String>) -> Self {
        self.map_source(|source| source.map_inner(|inner| inner.jetstream(stream)))
    }

    fn durable(self, name: impl Into<String>) -> Self {
        self.map_source(|source| source.map_inner(|inner| inner.durable(name)))
    }
}

Настройки издателя в терминах вашего брокера

Сторона публикации устроена так же. Политику точка монтирования называет через .out(marker, policy): маркер Reply - для того, что возвращает обработчик с publish("dest"), маркер слота - для слота Out. Хук над политикой этой позиции - MapPublisher:

use ruststream::runtime::MapPublisher;

pub trait NatsPublish {
    fn stream(self, name: impl Into<String>) -> Self;
    fn expect_last_sequence(self, seq: u64) -> Self;
}

impl<T: MapPublisher<Policy = Publish>> NatsPublish for T {
    fn stream(self, name: impl Into<String>) -> Self {
        self.map_publisher(|policy| policy.stream(name))
    }

    fn expect_last_sequence(self, seq: u64) -> Self {
        self.map_publisher(|policy| policy.expect_last_sequence(seq))
    }
}

В сервисе это читается так:

b.include(confirm).out_reply(Publish).stream("ORDERS");
b.include(mirror).out(Audit, Publish).stream("AUDIT").build();

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

map_publisher заменяет политику политикой того же типа, а другой тип политики означает другой режим публикации, и место ему в самом вызове .out(marker, policy). Уже настроенное значение можно передать и прямо туда: .out_reply(Publish::default().stream("ORDERS")).

Настройки на сообщение в билдере публикации

Место вызова меняет поле вашего Publisher::Options шагом, который вы добавляете к билдеру публикации. Издателя при этом ничто не оборачивает, поэтому публикация по-прежнему проходит через запись места монтирования, с её кодеком и её преобразованиями.

Частей четыре: тип настроек, где все поля необязательные; политика, которая хранит умолчания; живой издатель, сводящий одно с другим; и трейт-расширение над PublishBuilder, ограниченное типом настроек. Именно это ограничение не пускает ваши шаги на билдер над издателем другого брокера:

tests/publish_options.rs
/// The broker's per-message settings. Every field optional: what a call leaves unset keeps what
/// the policy fixed. `Debug` and `PartialEq` are not part of the contract, they are what a test
/// naming this type needs to assert on it.
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
struct PriorityOptions {
    priority: Option<u8>,
    ttl: Option<u8>,
}

/// The publish policy: pure declaration, constructible anywhere, and the place the defaults are
/// configured.
#[derive(Debug, Clone, Copy, Default)]
struct PriorityPublish {
    priority: u8,
    ttl: u8,
}

impl PriorityPublish {
    fn priority(mut self, priority: u8) -> Self {
        self.priority = priority;
        self
    }

    fn ttl(mut self, ttl: u8) -> Self {
        self.ttl = ttl;
        self
    }
}

/// The live publisher: the connection, plus the defaults the policy carried.
struct PriorityPublisher {
    inner: MemoryPublisher,
    default_priority: u8,
    default_ttl: u8,
}

impl PublishPolicy<ConnectedMemoryBroker> for PriorityPublish {
    type Live = PriorityPublisher;

    async fn pair(self, connected: &ConnectedMemoryBroker) -> Result<Self::Live, PairError> {
        Ok(PriorityPublisher {
            inner: Publish.pair(connected).await?,
            default_priority: self.priority,
            default_ttl: self.ttl,
        })
    }
}

impl Publisher for PriorityPublisher {
    // Forwarded to the bus underneath, which keeps what it is handed.
    type Payload = Take;
    type Error = MemoryError;
    type Options = PriorityOptions;

    async fn publish(
        &self,
        msg: OutgoingMessage<'_, BytesMut>,
        options: Option<&Self::Options>,
    ) -> Result<(), Self::Error> {
        let priority = options
            .and_then(|options| options.priority)
            .unwrap_or(self.default_priority);
        let ttl = options
            .and_then(|options| options.ttl)
            .unwrap_or(self.default_ttl);
        // A real broker hands the resolved values to its client as the protocol fields they are.
        // The in-memory bus has no such field, so this one puts them where a test can read them
        // back.
        let mut headers = msg.headers().clone();
        headers.insert("priority", priority.to_string());
        headers.insert("ttl", ttl.to_string());
        let stamped = OutgoingMessage::new(msg.name(), msg.payload()).with_headers(headers);
        self.inner.publish(stamped, None).await
    }
}

impl TransactionalPublisher for PriorityPublisher {
    async fn begin_transaction(&self) -> Result<(), Self::Error> {
        self.inner.begin_transaction().await
    }

    async fn commit(&self) -> Result<(), Self::Error> {
        self.inner.commit().await
    }

    async fn abort(&self) -> Result<(), Self::Error> {
        self.inner.abort().await
    }
}

/// The steps the broker puts in its prelude. The bound on the sink's options type is what keeps
/// them off a builder over any other broker's publisher.
trait PriorityPublishSteps {
    /// Sends this one message at `priority`, whatever the mount site's default is.
    #[must_use]
    fn priority(self, priority: u8) -> Self;

    /// Sends this one message with `ttl`, whatever the mount site's default is.
    #[must_use]
    fn ttl(self, ttl: u8) -> Self;
}

impl<Sink, Body, Enc, Hdrs, Dest> PriorityPublishSteps
    for PublishBuilder<Sink, Body, Enc, Hdrs, Dest>
where
    Sink: PublishSink<Options = PriorityOptions>,
{
    fn priority(mut self, priority: u8) -> Self {
        self.options_mut()
            .get_or_insert_with(PriorityOptions::default)
            .priority = Some(priority);
        self
    }

    fn ttl(mut self, ttl: u8) -> Self {
        self.options_mut()
            .get_or_insert_with(PriorityOptions::default)
            .ttl = Some(ttl);
        self
    }
}

Брокерная половина одинакова на пути с макросами и на ручном пути. Экспортируйте трейт-расширение из своей прелюдии рядом с псевдонимами политик.

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

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

Трейты-совместимости

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

Трейт Для брокеров, которые умеют
BatchSubscriber получать сообщения пакетами
TransactionalPublisher вызывать begin / commit / abort вокруг публикаций на дескрипторе издателя
OwnedTransactions / Transaction держать на одном дескрипторе сколько угодно транзакций сразу, каждую со своим буфером
RequestReply делать нативный request-reply
Partitioned ставить ключ партиционирования на исходящие сообщения
Seekable / Seeker перемещать живую подписку по воспроизводимому логу
Positioned сообщать позицию доставки в логе
DescribeServer сообщать ServerSpec для AsyncAPI

Seekable выдаёт дескриптор Seeker до того, как stream заимствует подписчика, поэтому работающую подписку можно переместить извне цикла диспетчеризации.

Позиции принадлежат брокеру: конструкторы в духе KafkaPosition вы объявляете на своём типе. Позиция, снятая с доставленного сообщения через Positioned::position, закрепляет контракт: перемотка на неё повторно доставит ровно это сообщение. У сконструированных позиций семантика та, которую описывает ваш тип позиции.

Опишите, на что распространяется одна перемотка (на экземпляр консьюмера или на общий курсор группы), и сбрасывайте весь учёт ack, который перемотка делает недействительным.

Чтобы тела обработчиков могли перематывать, положите позицию доставки и дескриптор перемотки полями своего контекста доставки и опубликуйте для них ключи ContextField; образец - MemoryContext in-memory брокера с ключами Position и SeekHandle. Пакетные формы берут дескриптор перемотки из пакетного контекста (см. ниже), в котором позиции нет.

Описание сервера из DescribeServer сообщает хост и порт, к которым подключаются клиенты. Учётные данные в него не попадают: документ создаётся для публикации. Брокер, настроенный по URL, строит описание через ServerSpec::from_url - этот конструктор убирает из URL имя пользователя и пароль. Если срезать у URL только схему и передать остаток дальше, учётные данные останутся в документе: это и есть ошибка, которую закрыл from_url, и она успела попасть в несколько крейтов брокеров. Брокер с несколькими адресами собирает их из ServerSpec::host_from_url.

Эти трейты - и есть словарь, которым пишет тело обработчика. Тело ограничивает свой слот нужной совместимостью (Out<impl TransactionalPublisher, Journal> или, на ручном пути, where W: TransactionalPublisher), а не вашим типом, и точка монтирования один раз, во время компиляции, сверяет с этим ограничением живую форму привязанной политики.

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

Пакеты: BatchSubscriber

Обработчик, принимающий &[T], потребляет пакет, а точка монтирования называет одно число - размер пакета. Рантайм передаёт его прямо в BatchSubscriber::batches(size). Пакет, который отдаёт ваш подписчик, тело и видит: рантайм не делит и не склеивает пакеты, поэтому в пакете никогда не больше size сообщений, а меньше - столько, сколько нашлось у транспорта.

Переведите size в понятие своего клиента: XREADGROUP COUNT, pull-пакет JetStream, лимит poll в Kafka. Всё остальное о том, как складывается пакет (таймаут блокировки, группа консьюмеров, окно предвыборки), остаётся вашим словарём и настраивается на источнике подписки через расширяющий трейт. Сервис тогда пишет b.include(handler.batch(nonzero!(6)).block(Duration::from_secs(5))): сначала слово ядра, потом ваше.

Ставьте совместимость на каждого подписчика, которого может открыть монтирование, а не только на того, которого открывает ваш дескриптор. #[subscriber("topic")] идёт через Subscribe, поэтому тело &[T] в этой форме просит BatchSubscriber у Subscribe::Subscriber.

У крейта, который поставил совместимость только на подписчика своего дескриптора, форма со строковым литералом не компилируется. Если это один тип, делать нечего; если типы разные, нужны оба.

Если транспорт отдаёт по одному сообщению, всё равно реализуйте совместимость и соберите пакеты на клиенте через BufferedSubscriber ядра: его batches соблюдает названный размер. Размер вы не выбираете, а срок, по которому закрывается неполный пакет, выбираете вы, и константой он быть не обязан.

Выведите этот срок на свой дескриптор подписки (.max_wait(Duration::from_millis(25))) и передавайте обёртке в момент подписки, чтобы сервис настраивал его на каждой подписке отдельно. Умолчание в 10 мс рассчитано на внутрипроцессную шину: при обращении по сети оно закрывает большинство пакетов на одной доставке, поэтому крейты брокеров, которые выводят этот срок опцией дескриптора, останавливаются на 10-50 мс.

Всё остальное в подписчике проходит сквозь обёртку без изменений:

tests/batch_subscriber.rs
/// What a broker crate writes when its transport has no batches of its own: the subscriber it
/// already has, wrapped in the core's client-side buffer, and `BatchSubscriber` delegated to it.
/// The deadline that closes a partial batch is the broker's own choice; the batch size is not -
/// it arrives per subscription, as the argument of `batches`.
struct TrickleSubscriber(BufferedSubscriber<MemorySubscriber<Retaining>>);

impl TrickleSubscriber {
    fn new(inner: MemorySubscriber<Retaining>) -> Self {
        Self(BufferedSubscriber::new(inner).max_wait(Duration::from_millis(5)))
    }
}

impl Subscriber for TrickleSubscriber {
    type Message = <MemorySubscriber<Retaining> as Subscriber>::Message;
    type Error = <MemorySubscriber<Retaining> as Subscriber>::Error;

    fn stream(&mut self) -> impl Stream<Item = Result<Self::Message, Self::Error>> + Send + '_ {
        self.0.stream()
    }
}

impl BatchSubscriber for TrickleSubscriber {
    type Batch = Vec<<MemorySubscriber<Retaining> as Subscriber>::Message>;

    fn batches(
        &mut self,
        size: NonZeroUsize,
    ) -> impl Stream<Item = Result<Self::Batch, Self::Error>> + Send + '_ {
        self.0.batches(size)
    }
}

/// Buffering does not move the subscription, so every other capability reaches through the
/// wrapper unchanged - here the seeker, which is what lets a batch subscription open at a
/// position even where the batches are assembled on the client.
impl Seekable for TrickleSubscriber {
    type Seeker = <MemorySubscriber<Retaining> as Seekable>::Seeker;

    fn seeker(&self) -> Self::Seeker {
        self.0.seeker()
    }
}

/// The broker's own subscription descriptor, opening the batching subscriber above.
#[derive(Clone)]
struct Trickle {
    name: &'static str,
}

impl SubscriptionSource<ConnectedMemoryBroker<Retaining>> for Trickle {
    type Subscriber = TrickleSubscriber;
    type Copies = AddressedCopies;

    fn name(&self) -> &str {
        self.name
    }

    async fn subscribe(
        self,
        connected: &ConnectedMemoryBroker<Retaining>,
    ) -> Result<TrickleSubscriber, MemoryError> {
        Ok(TrickleSubscriber::new(
            Subscribe::subscribe(connected, self.name).await?,
        ))
    }
}

// A descriptor that addresses its own copies says where they go, and one in-memory subject is
// both ends of the bus, so the name answers for itself with no lookup in between.
impl RedeliveryAddressed<ConnectedMemoryBroker<Retaining>> for Trickle {
    fn redelivery_address(
        &self,
        _connected: &ConnectedMemoryBroker<Retaining>,
    ) -> impl Future<Output = Result<RedeliveryAddress, MemoryError>> + Send {
        ready(Ok(RedeliveryAddress::new(self.name)))
    }
}

В точке монтирования не видно, какой из двух путей вы выбрали: сервис называет размер пакета и получает пакеты.

Отказаться от совместимости - тоже законный ответ там, где пакетирование ломает гарантию транспорта. Показательный случай - ROUTER в ZeroMQ: он отвечает каждому пиру на его собственный reply-to, а на весь пакет приходится один PublishContext, поэтому ответы всего пакета пришли бы на адрес одного пира. Напишите об этом в документации крейта: тело &[T] тогда не компилируется на этом транспорте.

Контракт проверяет пакетный набор conformance: он подписывается с размером меньше прогона и не пропускает брокер, у которого пакеты приходят длиннее. В harness::run_suite этот набор не входит: наборы совместимостей вы вызываете сами, по одному на каждую реализованную совместимость.

Прелюдия, которую поставляет ваш крейт

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

  1. pub use ruststream::prelude::*;, чтобы один glob обслуживал весь файл;
  2. вашу собственную поверхность, которую называет сервис: брокер, его источник подписки, его Config, его ошибку, ключи ContextField, которые читает тело;
  3. ваши политики публикации под теми едиными именами, которыми пользуются все брокеры, - Publish, а где они есть, ещё TransactionalPublish и Request (pub use crate::KafkaTransactionalPublish as TransactionalPublish;). Добавьте манифестом трейты совместимостей, которые вы реализуете на живых значениях, чтобы glob, называющий политики, приносил и их операции.

Под этими тремя именами прелюдия ядра не экспортирует ничего, поэтому точка монтирования читается одинаково на любом брокере, а смена брокера меняет только glob. Никогда не назначайте политике имя трейта ядра (Publisher, TransactionalPublisher, OwnedTransactions, RequestReply) и не реэкспортируйте под ним что-то другое: тело, которое подключило обе прелюдии, должно и дальше разрешать эти имена в трейты ядра.

Манифест - это то, что добавляет ваш glob: трейты потребляющей стороны, до которых тело добирается через ваш брокер (Positioned, Seeker, Transaction и подобные им). Четыре совместимости издателя уже лежат в прелюдии ядра, так что их повторный реэкспорт ничего не меняет.

Не включайте трейт, метод которого столкнётся с методом ядра по умолчанию (на практике это Partitioned::partition_key против IncomingMessage::partition_key): сервис, которому он нужен, импортирует его явно. BatchSubscriber не место ни в каком манифесте: его вызывает сам фреймворк, и ни одно тело не пишет его как ограничение.

Разобранный пример - ruststream::memory::prelude.

Расширение словаря слотов Out

Параметр обработчика Out<impl X, Marker> принимает любой X, реализованный живым значением за слотом; поверх этого ядро делегирует собственный набор совместимостей (Publisher, TransactionalPublisher, OwnedTransactions, RequestReply). Когда живое значение умеет больше - или вовсе не издатель (кэш продюсеров по партициям, шардирующий роутер), - объявите собственную трейт-совместимость и реализуйте её для живого значения.

Тело обработчика держит не само значение, а запись арены Slot<Marker, W, E, Pipe, Body> - прозрачное окно в него. Автоматическое разыменование пропускает через это окно вызов метода, но не ограничение по трейту: вспомогательная функция fn issue<L: Lanes>(lanes: &L) отвергает такую запись с ошибкой E0277.

Добавьте рядом со своим трейтом один blanket-impl - impl<M, W: Lanes, E, Pipe, Body> Lanes for Slot<M, W, E, Pipe, Body> с делегированием через Deref записи, - и функции и тела, обобщённые по совместимости, принимают запись как есть. Конкретный тип в коде приложения по-прежнему не появляется:

tests/out_slots.rs
// A paired value that is NOT a publisher: a lane router in the shape of a broker's
// per-partition producer cache. The capability is broker-defined; the core knows nothing
// about it.
#[derive(Clone)]
struct LaneRouter {
    publisher: MemoryPublisher,
}

/// The broker-defined capability: pick a destination lane for a shard.
trait ShardLanes {
    fn lane(&self, shard: u64) -> (&MemoryPublisher, &'static str);
}

impl ShardLanes for LaneRouter {
    fn lane(&self, shard: u64) -> (&MemoryPublisher, &'static str) {
        let dest = if shard.is_multiple_of(2) {
            "slots.lane.even"
        } else {
            "slots.lane.odd"
        };
        (&self.publisher, dest)
    }
}

// Grafted onto the arena entry once, for every marker, delegating through the entry's
// transparent `Deref`: this is how a broker crate extends the slot vocabulary with its own
// traits. A handler body holds the entry, so without this impl the capability is reachable by
// autoderef for a method call but never satisfies a trait bound.
impl<M, W: ShardLanes, E, Pipe, Body> ShardLanes for Slot<M, W, E, Pipe, Body> {
    fn lane(&self, shard: u64) -> (&MemoryPublisher, &'static str) {
        (**self).lane(shard)
    }
}

/// The bound the graft buys: a helper generic over the capability, not over the concrete live
/// type, takes the entry a handler body holds.
async fn sent<L: ShardLanes + Sync>(lanes: &L, event: &Event) -> bool {
    let (publisher, dest) = lanes.lane(event.id);
    publisher.message(event).to(dest).publish().await.is_ok()
}

/// The policy half: pure declaration pairing into the router, like a broker's
/// `per_partition()` policy pairs into its producer cache. No `Clone`: resolution consumes it.
struct LanePolicy;

impl PublishPolicy<ConnectedMemoryBroker> for LanePolicy {
    type Live = LaneRouter;

    async fn pair(self, connected: &ConnectedMemoryBroker) -> Result<Self::Live, PairError> {
        Ok(LaneRouter {
            publisher: Publish.pair(connected).await?,
        })
    }
}

/// The handler bounds its slot with the broker-defined capability, not a core one.
#[subscriber("slots.sharded")]
async fn route_shard(event: &Event, Out(lanes): Out<impl ShardLanes>) -> HandlerOutcome {
    if sent(lanes, event).await {
        HandlerOutcome::ack()
    } else {
        HandlerOutcome::retry()
    }
}
tests/manual_out_slots.rs
// A paired value that is NOT a publisher: a lane router in the shape of a broker's
// per-partition producer cache. The capability is broker-defined; the core knows nothing
// about it.
#[derive(Clone)]
struct LaneRouter {
    publisher: MemoryPublisher,
}

/// The broker-defined capability: pick a destination lane for a shard.
trait ShardLanes {
    fn lane(&self, shard: u64) -> (&MemoryPublisher, &'static str);
}

impl ShardLanes for LaneRouter {
    fn lane(&self, shard: u64) -> (&MemoryPublisher, &'static str) {
        let dest = if shard.is_multiple_of(2) {
            "slots.lane.even"
        } else {
            "slots.lane.odd"
        };
        (&self.publisher, dest)
    }
}

// Grafted onto the arena entry once, for every marker, delegating through the entry's
// transparent `Deref`: this is how a broker crate extends the slot vocabulary with its own
// traits. A body holds the entry, so without this impl the capability is reachable by autoderef
// for a method call but never satisfies a trait bound.
impl<M, W: ShardLanes, E, Pipe, Body> ShardLanes for Slot<M, W, E, Pipe, Body> {
    fn lane(&self, shard: u64) -> (&MemoryPublisher, &'static str) {
        (**self).lane(shard)
    }
}

/// The bound the graft buys: a helper generic over the capability, not over the concrete live
/// type, takes the entry a body holds.
async fn sent<L: ShardLanes + Sync>(lanes: &L, event: &Event) -> bool {
    let (publisher, dest) = lanes.lane(event.id);
    publisher.message(event).to(dest).publish().await.is_ok()
}

/// The policy half: pure declaration pairing into the router, like a broker's
/// `per_partition()` policy pairs into its producer cache. No `Clone`: resolution consumes it.
struct LanePolicy;

impl PublishPolicy<ConnectedMemoryBroker> for LanePolicy {
    type Live = LaneRouter;

    async fn pair(self, connected: &ConnectedMemoryBroker) -> Result<Self::Live, PairError> {
        Ok(LaneRouter {
            publisher: Publish.pair(connected).await?,
        })
    }
}

/// The body leaves the wired live value generic and bounds it with the broker-defined
/// capability, exactly as the attribute's `Out<impl ShardLanes>` does.
struct RouteShard;

struct Lanes;

impl OutSlot for Lanes {
    const NAME: &'static str = "Lanes";
    type Destination = Reads;
}

impl<L> Handle<Event, (), Outs<(L,)>> for RouteShard
where
    L: OutEntry<Lanes, Wire: ShardLanes>,
{
    async fn handle(
        &self,
        event: &Event,
        outs: &Outs<(L,)>,
        _ctx: &mut Context<'_>,
    ) -> Result<(), HandlerOutcome> {
        if sent(outs.get(Lanes), event).await {
            Ok(())
        } else {
            Err(HandlerOutcome::retry())
        }
    }
}

Форму трейта задаёт место отправки, и таких форм две.

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

Совместимость в форме шага задаёт одну настройку на сообщение и заканчивается одной публикацией: ключ упорядочивания, приоритет, QoS. Это вообще не трейт-совместимость: это поле вашего Publisher::Options, которое меняют шагом билдера публикации. См. настройки на сообщение в билдере публикации.

Прелюдия вашего крейта

Два файла импортируют разное, и это разделение сохраняет переносимость сервиса. Тело обработчика импортирует ruststream::prelude::* и ничего вашего: внедрённый слот оно ограничивает нужной трейт-совместимостью ядра - Out<impl Publisher>, Out<impl TransactionalPublisher>, Out<impl OwnedTransactions>, Out<impl RequestReply>, - то есть говорит, что ему нужно от издателя, и никогда - какой брокер это даёт.

Файл монтирования импортирует вашу прелюдию, потому что брокера называют именно там.

Единственное исключение - настройка на сообщение: место её вызова находится в теле. Тело, которое её меняет, импортирует вашу прелюдию ради шага и называет ваш тип настроек в ограничении (Out<impl Publisher<Options = MqttOptions>, Telemetry>). Такое тело привязано к вашему брокеру и говорит об этом своей сигнатурой.

Значит, ваша прелюдия - единственный ваш импорт, который пишет сервис, и её устройство входит в контракт. Псевдонимы политик (NatsPublish as Publish, KafkaTransactionalPublish as TransactionalPublish, LapinRequest as Request) делают файл монтирования одинаковым на любом брокере, а смену брокера сводят к смене импорта.

Ваша половина правила об именах: явный ре-экспорт перекрывает glob без единого предупреждения, поэтому имя, названное как трейт ядра, отнимает этот трейт у каждого сервиса, который пишет glob, и ошибка возникает в файле сервиса, а не в вашем.

Закрепите обе половины пробой за собственным glob-импортом: ограничение, которое пишет тело, обязано по-прежнему разрешаться в трейт ядра, а имя точки монтирования - оставаться вашей политикой.

// in your crate, behind your own prelude glob
use crate::prelude::*;

// A capability bound a body states: the core trait, not something of yours.
fn _p<T: Publisher>() {}

// A mount-site name: your policy, constructible with no connection in sight.
fn _q() {
    let _: Publish = Publish::default();
}

Контекст доставки и ключи Ctx

Брокер с нативными метаданными доставки (партиция, смещение, номер в стриме) отдаёт их типизированным контекстом доставки: это структура с #[non_exhaustive], которую называет подписчик, плюс типы-ключи ContextField. По ключу обработчик привязывает отдельное поле параметром через экстрактор Ctx<K>. Ключи - пустые структуры, а на пути доставки нет ни type-map, ни выделений в куче.

/// Per-delivery context of this broker.
#[non_exhaustive]
#[derive(Debug, Clone)]
pub struct MyContext {
    pub partition: i32,
}

/// `Ctx<Partition>` in a handler binds the delivery's partition.
#[derive(Debug, Default, Clone, Copy)]
pub struct Partition;

impl ContextField for Partition {
    type Context = MyContext;
    type Value = i32;
    fn read(self, src: &MyContext) -> i32 {
        src.partition
    }
}

В наброске ключ читает Copy-скаляр, у которого владение и заимствование неразличимы. Позицию, которая не Copy (идентификатор сообщения в Pulsar, шард Kinesis вместе со строкой последовательности), читают заимствованием: Field::Value<'a> обобщён по времени жизни источника, поэтому ключ отдаёт &'a MessageId, и тело, которое читает его через ctx.context(..), ничего не копирует.

Владеющим и 'static обязано быть только ContextField::Value - значение за экстрактором Ctx<K>, потому что значения экстракторов связываются до запуска тела, и этот ключ клонирует то, что вернул заимствующий. Обычно один ключ реализует оба трейта, по форме на каждый.

Брокер, у которого своих полей доставки нет, берёт ().

У пакетных подписок контекст свой, потому что пакет охватывает много доставок. Соберите вторую структуру из того, что общее для всей подписки (дескриптор перемотки, имя стрима, группа консьюмеров), реализуйте на ней BuildBatchContext и опубликуйте ключи Field, чтобы пакетное тело читало её через ctx.context(..). Одно значение на пакет рантайм строит по первой доставке пакета.

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

Образец - MemoryBatchContext in-memory брокера: дескриптор перемотки подписки лежит под тем же ключом SeekHandle, который публикует контекст доставки. Брокер, которому нечего дать на уровне подписки, не реализует ничего и оставляет пакеты на умолчании ().

Middleware на асинхронных краях

Интеграциям, которым нужен асинхронный I/O вокруг кодирования и декодирования (schema registry, конверт поверх формата передачи), не место в Codec: кодек ядра синхронный, а обработчикам стоит оставаться на кодеке по умолчанию.

Ставьте такие интеграции на асинхронные края. Входящие полезные нагрузки перекодируйте на пути доставки подписки, до того как их увидит кодек, а исходящие оборачивайте слоем PublishLayer из ядра, добавленным на всё приложение через RustStream::publish_layer. Слой публикации асинхронный и может вернуть ошибку, а Outgoing::payload_mut существует ровно для упаковки в конверт.

Биндинги протокола

В генерируемом документе AsyncAPI есть место для того, что знает только ваш брокер: долговечность очереди RabbitMQ, группа потребителей Kafka, QoS у MQTT. Спецификация называет это биндингами, а заполняет их ваш дескриптор.

tests/asyncapi.rs
/// A descriptor of the shape a broker crate ships: it reads its own private fields and says
/// what the protocol calls them.
#[derive(Clone)]
struct RabbitQueue {
    name: &'static str,
    durable: bool,
}

impl<C: Subscribe> SubscriptionSource<C> for RabbitQueue {
    type Subscriber = C::Subscriber;
    type Copies = AddressedCopies;

    fn name(&self) -> &str {
        self.name
    }

    async fn subscribe(self, connected: &C) -> Result<Self::Subscriber, C::Error> {
        connected.subscribe(self.name).await
    }

    fn channel_bindings(&self) -> Bindings {
        let body = AmqpChannel {
            is: "queue",
            queue: AmqpQueue {
                name: self.name,
                durable: self.durable,
            },
        };
        Binding::new("amqp", "0.3.0", &body)
            .map(|binding| Bindings::new().with(binding))
            .unwrap_or_default()
    }

    fn operation_bindings(&self) -> Bindings {
        Binding::new("amqp", "0.3.0", &AmqpOperation { ack: true })
            .map(|binding| Bindings::new().with(binding))
            .unwrap_or_default()
    }

    fn message_bindings(&self) -> Bindings {
        let body = AmqpMessage {
            message_type: "order",
        };
        Binding::new("amqp", "0.3.0", &body)
            .map(|binding| Bindings::new().with(binding))
            .unwrap_or_default()
    }
}

Binding::new(protocol, version, &body) сериализует тело один раз и сам вписывает bindingVersion, так что биндинг без версии вы не отправите. Ключ протокола ядро сверяет с закрытым списком спецификации, и неизвестный ключ возвращается ошибкой, а не уезжает в документ, который ни один инструмент не прочитает. Bindings пуст по умолчанию: дескриптор, который ничего не говорит, ничего в документе не меняет.

Уровень сервера - это поле, а не метод: сервер описывается один раз на брокера. ServerSpec::new(host, protocol).bindings(..) в вашей реализации DescribeServer.

Что уместно в биндинге, ограничивают три правила.

Значение считается по одному дескриптору. Документ собирается до соединения, поэтому настоящее число партиций топика Kafka, топик за подпиской Pub/Sub и ARN очереди SQS оттуда взяться не могут.

Учётные данные в биндинг не попадают - по той же причине, что и в DescribeServer. Проверка лежит в conformance::harness: настройте брокера и дескриптор с заведомо известным паролем и запустите сканирование.

tests/conformance_self.rs
/// The in-memory broker has no network address and no binding, so nothing it describes can leak a
/// password. A broker configured from a URL runs this with the password it was configured with.
#[cfg(feature = "asyncapi")]
#[test]
fn memory_broker_describes_without_credentials() {
    harness::describes_without_credentials(
        &MemoryBroker::new(),
        &MemorySource::new("orders"),
        "hunter2",
    );
}

Протокол, для которого биндинга в спецификации нет, едет в Binding::extension("x-kinesis", &body). Ключи протоколов - закрытый список, поэтому у ZeroMQ, Kinesis и файлового транспорта законного ключа нет; расширение x- стоит на том же уровне и bindingVersion не несёт.

Биндинги приходят от дескриптора, поэтому подписка, открытая по голому имени, не несёт ни одного: описывать там нечего. Брокеру, которому нужны биндинги, нужен тип SubscriptionSource.

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

tests/asyncapi.rs
/// A publish policy of the shape a broker crate ships: an SNS topic is named by its binding's
/// required `name`, and the hook is handed the destination the mount site resolved.
#[derive(Clone, Copy, Default)]
struct TopicPublish;

impl PublishPolicy<ConnectedMemoryBroker> for TopicPublish {
    type Live = MemoryPublisher;

    fn pair(
        self,
        connected: &ConnectedMemoryBroker,
    ) -> impl Future<Output = Result<Self::Live, PairError>> {
        live(connected)
    }

    fn channel_bindings(&self, channel: &str) -> Bindings {
        let body = SnsChannel {
            name: channel.to_owned(),
        };
        one("sns", "0.1.0", &body)
    }

    fn operation_bindings(&self, channel: &str) -> Bindings {
        let body = SnsOperation {
            topic: SnsChannel {
                name: channel.to_owned(),
            },
        };
        one("sns", "0.1.0", &body)
    }

    fn message_bindings(&self, _channel: &str) -> Bindings {
        let body = SnsMessage {
            message_type: "progress",
        };
        one("sns", "0.1.0", &body)
    }
}

Каждый метод получает адресата, которого разрешила точка монтирования. У ответа это собственное имя типа ответа или клауза publish("dest") регистрации, у записи слота - её имя, у недоставленного сообщения - объявление dead_letter("dlq"). Топик SNS и очередь SQS названы обязательным полем name своего биндинга, и это имя берётся отсюда: политика держит настройки вашего брокера и никогда - адресата. Там, где преобразование называет адресата на каждой доставке, канал не сообщает адреса, и метод получает запасное имя точки монтирования.

У дескриптора такого параметра нет: он знает подписку, которую описывает.

Три правила здесь те же. Не переносится одно: своей операции send у ответа нет, поэтому operation_bindings политики ответа до документа не доходит. У слота и у адресата недоставленных она есть.

Четвёртый метод принадлежит только ответу. reply_address_location говорит, откуда клиент читает адрес ответа, если ваш брокер маршрутизирует ответы через заголовок reply-to:

tests/asyncapi.rs
/// The reply's own policy: it answers where a client reads the address of an answer.
#[derive(Clone, Copy, Default)]
struct ReplyToPublish;

impl PublishPolicy<ConnectedMemoryBroker> for ReplyToPublish {
    type Live = MemoryPublisher;

    fn pair(
        self,
        connected: &ConnectedMemoryBroker,
    ) -> impl Future<Output = Result<Self::Live, PairError>> {
        live(connected)
    }

    fn channel_bindings(&self, channel: &str) -> Bindings {
        let body = NatsChannel {
            subject: channel.to_owned(),
        };
        one("nats", "0.1.0", &body)
    }

    fn reply_address_location(&self) -> Option<&'static str> {
        Some("$message.header#/reply-to")
    }
}

Тогда документ показывает канал ответа с address: null, а выражение кладёт в reply.address.location операции receive. Читается оно только там, где точка монтирования поставила преобразование, называющее адресата на каждой доставке; иначе ответ уходит на объявленное имя, и документ говорит именно его.

У протокола, которого нет в списке спецификации, нет биндинга и на стороне публикации. Брокер в памяти - как раз такой случай: ключа memory не существует, поэтому MemoryPublish молчит, а не выдумывает ключ.

Хуки закрыты фичей ядра asyncapi. Пробросьте её из своего крейта:

[features]
asyncapi = ["ruststream/asyncapi"]

и поставьте #[cfg(feature = "asyncapi")] на каждый метод, который заполняете.

Конфигурация и умолчания

Config принадлежит вашему крейту: специфичной для брокера конфигурации в ядре нет. Если у поля нет разумного умолчания, не реализуйте Default: пусть пользователь задаст значение явно, а не получит умолчание, которое потом сломается.

Ошибки

Возьмите thiserror и одно перечисление ошибок на весь крейт, с вариантами по источнику. Публичные перечисления ошибок помечайте #[non_exhaustive]. anyhow в библиотечном крейте не используйте никогда.

Поддержка тестирования

Поставляйте внутрипроцессный транспорт, который реализует TestableBroker на подключённой форме, под фичей testing. Регистрируйте его через register_testable_broker! именно для подключённого типа: обвязка подключает каждый брокер, прежде чем достать его транспорт. Тогда пользователи смогут писать модульные тесты обработчиков на вашем брокере с обвязкой TestApp.

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

Не имитируйте в транспорте специфичную для брокера семантику (durable-курсоры, таймеры повторной доставки, смещения, маршрутизацию в dead-letter): её проверяют сквозными тестами на настоящем сервере.

Эталон - собственная реализация in-memory брокера (на ConnectedMemoryBroker):

src/memory/mod.rs
// The harness drives the connected form: TestApp connects every registered broker before it
// recovers the in-process transport, and run_suite scenarios receive connected brokers.
#[cfg(feature = "testing")]
impl<Log: LogMode> crate::testing::TestableBroker for ConnectedMemoryBroker<Log> {
    fn install_coordinator(&self, coordinator: Coordinator) {
        self.state.install_coordinator(coordinator);
    }

    fn inject(&self, message: OutgoingMessage<'_>) {
        let (name, payload, headers) = message.into_parts();
        // Injecting into a shut-down bus is a harness bug (both run_suite and TestApp drive
        // the bus strictly before shutdown), so fail loudly instead of losing the message.
        self.state
            .fanout(name, Bytes::copy_from_slice(payload), headers)
            .expect("inject on a shut-down broker: drive the harness before shutdown");
    }

    /// What the broker holds under `name`: everything published there on a discarding broker
    /// (which records for the length of a harness run), and the retained window on a retaining
    /// one, so an assertion never claims more than the broker keeps.
    fn published(&self, name: &str) -> Vec<RawMessage> {
        self.state
            .log
            .lock()
            .expect("memory broker mutex poisoned")
            .name(name)
            .map(|log| log.messages(name))
            .unwrap_or_default()
    }
}

// One registration per log mode: the harness recovers a broker by its concrete type, and the
// two modes are two types.
#[cfg(feature = "testing")]
crate::register_testable_broker!(ConnectedMemoryBroker<Discarding>);
#[cfg(feature = "testing")]
crate::register_testable_broker!(ConnectedMemoryBroker<Retaining>);

Транспорт вызывает Coordinator::enqueued на каждую постановку в очередь подписчика и Coordinator::consumed, когда доставка завершена или отброшена, - по ним обвязка определяет, что реакция закончилась. Отложенные повторные доставки транспорт направляет через Coordinator::schedule_redelivery.

Один такой тип работает и с TestApp, и с набором conformance. Пользовательская сторона описана в разделе Тестирование, а Conformance показывает, как подтвердить реализацию через run_suite и проверку жизненного цикла lifecycle.

Как написать транспорт, которому можно верить

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

Прогоняйте наборы ядра не только по серверу, но и по внутрипроцессному транспорту. Наборы написаны против трейтов и не различают, кто отвечает - настоящий брокер или внутрипроцессный транспорт. Достаточно одного #[tokio::test]:

tests/conformance_self.rs
/// 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;
}

Первым запускайте lifecycle. Он проходит new -> connect -> подписку -> публикацию -> ack -> shutdown, а потом задаёт вопрос, который внутрипроцессному транспорту почти никогда не задают: вернёт ли ошибку издатель, созданный до остановки, если обратиться к нему после? Настоящий клиент отвечает «нет соединения». Внутрипроцессный транспорт, у которого публикация сводится к отправке в канал, ошибки не вернёт: он примет сообщение.

tests/conformance_self.rs
// `make_source` / `make_publisher` must stay closures: their bounds are higher-ranked
// (`Fn(&str) -> _` / `Fn(&B) -> _`), so a bare method path - which binds one concrete lifetime -
// would not type-check.
#[allow(clippy::redundant_closure, clippy::redundant_closure_for_method_calls)]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn memory_broker_passes_lifecycle() {
    harness::lifecycle(
        MemoryBroker::new,
        |name| MemorySource::new(name),
        |broker| broker.publisher(),
    )
    .await;
}

Так же добавьте наборы capabilities::* - по одному на каждую реализованную возможность.

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

Завершайте доставку так же, как её завершает транспорт. Настоящий ack возвращает AckError::Unsupported там, где транспорт не подтверждает доставку: при отправке без подтверждений и на уровне качества at-most-once. Внутрипроцессный транспорт возвращает то же самое. Если ответить Ok(()), лишь бы набор прошёл, обработчик с HandlerOutcome::retry() пройдёт тесты и потеряет сообщение на настоящем брокере. Честный ответ наборы принимают.

Повторяйте то, что делает клиент, и не подделывайте то, что делает брокер. Граница проходит не по трудозатратам, а по тому, на какой стороне работает механизм. Конкурирующие потребители, распределение по группе, корреляция и маршрутизация ответов, буферизация до коммита - это работа клиента или уровня маршрутизации, и внутри процесса она повторяется точно. Атомарность в кластере, fencing, тайм-ауты на стороне брокера, exactly-once - работа брокера, и внутри процесса получается выдумка.

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

У каждого пробела - комментарий: какое утверждение этот пробел делает несостоятельным. Пишите не «возможности нет», а какому тесту читатель больше не может верить и что закрывает эту проверку на самом деле:

// No fencing: a second producer claiming the same transactional id is not rejected here, so a
// test cannot assert the first one is fenced out. `capabilities::transactions` against a real
// server is what covers that.

Закрепляйте пробел ещё и тестом. Комментарий устареет в тот день, когда кто-нибудь «починит» транспорт и научит его маршрутизировать то, что он намеренно не маршрутизирует. Тест, который утверждает, что обработчик не вызван, в этот день упадёт и объяснит себя сам.

Монтируйте внутрипроцессный транспорт так же, как настоящий брокер. Ваши источники подписки и политики публикации должны работать с ним без изменений - тогда сервис проверяет тот файл маршрутов, который поставляет. Если ради запуска теста пользователь меняет OrdersStream на что-то другое, монтирование остаётся непроверенным.