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

Memory

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

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

MemoryBroker за фичей memory - полноценный брокер внутри процесса. Он подходит, когда очередь нужна в рамках одного приложения, а не в сети. На нём построен шаблон cargo generate по умолчанию (templates/memory), поэтому свежий проект запускается без внешних зависимостей.

ruststream = { version = "0.7", features = ["macros", "memory", "json"] }
use ruststream::memory::MemoryBroker;

let broker = MemoryBroker::new();

Сколько он хранит

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

Сервису с перемоткой нужна история, и он называет её объём:

examples/seek.rs
// Replaying reads what the broker kept, so this one keeps the last 64 entries of every log.
// `MemoryBroker::new()` keeps nothing, and a mount that seeks does not compile on it.
let broker = MemoryBroker::retaining(Retention::Messages(nonzero!(64)));

Retention ограничивает одну тему: Messages(n) хранит n свежих сообщений каждой темы, Bytes(n) - свежие тела, которые помещаются в n байт, а MessagesAndBytes { .. } применяет оба предела сразу. Брокер, публикующий в тысячу тем, держит столько на каждую из них. Самое свежее сообщение остаётся всегда, поэтому тело шире предела по байтам хранится в одиночку, а не пропадает в момент публикации.

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

Прелюдия, которую импортирует точка монтирования

ruststream::memory::prelude - glob этого брокера, устроенный как прелюдия любого брокерного крейта. Он реэкспортирует прелюдию ядра, затем поверхность самого брокера (MemoryBroker, MemorySource, MemoryError, MemoryPosition, Retention с режимами журнала Discarding / Retaining и ключи контекста MemoryContext / MemoryBatchContext / Position / SeekHandle), затем политики публикации под едиными именами Publish, TransactionalPublish и Request. Все три - псевдонимы MemoryPublish и MemoryRequest. Издатель этого брокера реализует обе разновидности транзакций, поэтому TransactionalPublish здесь - та же политика, что и Publish; у брокера с отдельной транзакционной конфигурацией это имя указывает на другую политику.

use ruststream::memory::prelude::*;

Тот же glob вводит в область видимости трейты-совместимости TransactionalPublisher, OwnedTransactions, Transaction, RequestReply, Positioned и Seeker, поэтому их операции доступны там же, где и политики. Partitioned в него не входит: в области видимости он делает msg.partition_key() неоднозначным с одноимённым методом IncomingMessage. Сервис, который читает ключи партиционирования, импортирует Partitioned сам.

Тело обработчика оставляет use ruststream::prelude::*;: оно называет совместимости, а не политики, и не знает, какой брокер его выполняет. Файлу, где лежат и тело, и точка монтирования, хватает одного брокерного glob.

Семантика

  • Имя темы сравнивается целиком. Подписка на тему orders получает сообщения, опубликованные в orders.
  • Доставка всем подписчикам. Каждый подписчик темы получает каждое сообщение, опубликованное после подписки.
  • Ack ничего не делает, а nack(requeue: true) доставляет ту же полезную нагрузку тому же подписчику заново.
  • retry_after брокер отрабатывает сам. Доставка возвращается тому же подписчику, когда истечёт задержка, и ничего не публикуется заново.
  • Доставки считаются. Каждая доставка сообщает, сколько раз брокер отдавал это сообщение подписчику, считая первую доставку. Поэтому предел max_attempts(n) расходуется на этих повторных доставках, а доставка, которая его израсходовала, уходит в назначение dead_letter(name).
  • Совместное владение. MemoryBroker - дескриптор с подсчётом ссылок: все его копии работают с одним и тем же состоянием, поэтому копия, которую держит тест, видит всё, что публикует приложение.

Обработчик, middleware и декодирование работают здесь так же, как с сетевым брокером: рантайм диспетчеризует сообщения одним и тем же путём.

Совместимости

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

  • Запрос и ответ. broker.requester() даёт MemoryRequester: его request публикует сообщение и указывает в заголовке reply-to уникальную внутрипроцессную тему для ответа, а завершается первым сообщением, которое туда доставили. Отвечающая сторона читает reply-to из запроса и публикует ответ в эту тему. Запрос, на который никто не ответил, возвращает ошибку RequestError::Timeout. Политика MemoryRequest конструирует MemoryRequester, поэтому слот с ограничением Out<impl RequestReply, ..> вы привязываете к MemoryRequest.
  • Пакеты. MemorySubscriber реализует BatchSubscriber: пакет - это первая пришедшая доставка и всё, что уже лежит в буфере, но не больше размера, заданного при регистрации обработчика через batch(n). Неполный пакет доставляется сразу.
  • Транзакции. Политика MemoryPublish конструирует MemoryPublisher, который реализует обе разновидности транзакций, поэтому слот или связывание с ограничением TransactionalPublisher или OwnedTransactions вы привязываете к MemoryPublish. Публикации внутри области транзакции буферизуются: commit доставляет их всем подписчикам разом в порядке публикации, abort отбрасывает. Каждая владеющая транзакция буферизует сама по себе, а копии дескриптора издателя транзакцию не разделяют. Нарушение порядка вызовов на самом издателе возвращает MemoryError: второй begin_transaction при открытой транзакции - TransactionBusy, и открытая транзакция остаётся нетронутой; commit или abort без транзакции - NoTransaction.
  • Ключи партиционирования. MemoryMessage реализует Partitioned и читает ключ из заголовка partition-key (memory::PARTITION_KEY_HEADER).
  • Перемотка журнала. У хранящего брокера MemorySubscriber реализует Seekable поверх журнала каждой темы: получите MemorySeeker до начала чтения, а затем вызовите seek на позиции MemoryPosition - снятой с доставленного сообщения через Positioned::position (тогда это же сообщение доставляется заново) или построенной (MemoryPosition::start() / sequence(n) / end()). Номера последовательности абсолютны и продолжают называть то же сообщение, пока предел хранения вытесняет более старые. start() - самое старое из сохранённых сообщений, end() - конец журнала, за всем опубликованным. Перемотка вперёд пропускает доставки, стоявшие в очереди до цели. Номер, уже вытесненный пределом, возвращает ошибку MemoryError::PositionEvicted с самой старой оставшейся позицией. Действует перемотка на один экземпляр подписчика, а через дескриптор уже остановленной шины возвращает ошибку MemoryError::ShutDown. Внутри приложения MemoryContext содержит позицию сообщения и MemorySeeker, а обработчик читает их по ключам Position и SeekHandle (см. Перемотку). Пакетный обработчик читает MemoryBatchContext: там есть SeekHandle, но нет Position, потому что пакет охватывает много доставок.
  • Остановка. MemoryBroker::connect(self) даёт ConnectedMemoryBroker, а его shutdown поглощает self и возвращает ClosedMemoryBroker, который сообщает, сколько регистраций подписчиков отброшено при остановке. После остановки публикация, фиксация транзакции и запрос через ранее выданные дескрипторы возвращают ошибку MemoryError::ShutDown или RequestError::ShutDown.

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

ConnectedMemoryBroker реализует Subscribe, поэтому #[subscriber("orders")] работает напрямую. Ту же подписку задаёт дескриптор MemorySource - в той же форме, в какой подписку описывает любой брокер. Из примера routed_service:

examples/routed_service/orders.rs
use ruststream::memory::prelude::*;

#[subscriber(MemorySource::new("orders"), publish("confirmations"))]
pub(crate) async fn confirm(
    order: &Order,
    ctx: &mut Context<'_, (), Repository>,
) -> Result<Confirmation, HandlerOutcome> {
    let repo = ctx.state();
    tracing::debug!(
        order = order.id,
        customer = %order.customer,
        item = %order.item,
        "confirming order"
    );
    match repo.record_order(order.id).await {
        Ok(()) => Ok(Confirmation {
            order_id: order.id,
            accepted: order.quantity > 0,
        }),
        Err(e) if e.is_transient() => {
            tracing::warn!(order = order.id, "store busy, asking for redelivery");
            Err(HandlerOutcome::retry())
        }
        Err(e) => {
            tracing::error!(order = order.id, error = %e, "dropping order");
            Err(HandlerOutcome::drop())
        }
    }
}
examples/manual/routed_service_orders.rs
use ruststream::memory::prelude::*;

struct Confirm;

// The state is named on the body, not on a definition: this one reads a `Repository`, so it is a
// `Handle` for that state alone and mounts only on an application that carries it.
impl Handle<Order, Confirmation, (), (), Repository> for Confirm {
    async fn handle(
        &self,
        order: &Order,
        _outs: &(),
        ctx: &mut Context<'_, (), Repository>,
    ) -> Result<Confirmation, HandlerOutcome> {
        let repo = ctx.state();
        tracing::debug!(
            order = order.id,
            customer = %order.customer,
            item = %order.item,
            "confirming order"
        );
        match repo.record_order(order.id).await {
            Ok(()) => Ok(Confirmation {
                order_id: order.id,
                accepted: order.quantity > 0,
            }),
            Err(e) if e.is_transient() => {
                tracing::warn!(order = order.id, "store busy, asking for redelivery");
                Err(HandlerOutcome::retry())
            }
            Err(e) => {
                tracing::error!(order = order.id, error = %e, "dropping order");
                Err(HandlerOutcome::drop())
            }
        }
    }
}

/// The mount, and the whole declaration the attribute's clauses carried: the broker's own
/// descriptor as the source, `.to(..)` for the reply channel, and `.describe(..)` for the sentence
/// the attribute lifts off the handler's doc comment. Who publishes the reply is wiring rather
/// than declaration, so it lives on the mount chain on both paths: `.out_reply(Publish)` names
/// the position the returned value leaves through and the policy that carries it, which pairs with
/// the connected broker at startup and encodes with the default codec. The definition says what it
/// replies with and where; the chain says who sends it.
fn confirm_route() -> impl RouterDef<MemoryBroker, Repository> {
    Router::<MemoryBroker>::new()
        .include(
            subscriber(MemorySource::new("orders"), Confirm)
                .reply()
                .to("confirmations")
                .describe("Confirms an order and replies on `confirmations`.")
                .build(),
        )
        .out_reply(Publish)
        .build()
}

Для тестов

Приложение на MemoryBroker вы проверяете обвязкой TestApp: соберите приложение, отдайте его в TestApp::start, публикуйте сообщения и проверяйте, какие сообщения получили и какие опубликовали обработчики. Полностью приём разобран в разделе Тестирование.

На время прогона обвязка записывает всё, что публикует сервис, поэтому проверки published::<T>(..) читают один и тот же список, на какой бы форме брокера приложение ни было собрано. Вне обвязки чтение журнала через TestableBroker::published показывает то, что брокер хранит: у хранящего - всё в пределах его ограничения, у брокера по умолчанию - ничего.