Memory¶
MemoryBroker за фичей memory - полноценный брокер внутри процесса. Он подходит, когда очередь
нужна в рамках одного приложения, а не в сети. На нём построен шаблон cargo generate по умолчанию
(templates/memory), поэтому свежий проект запускается без внешних зависимостей.
Сколько он хранит¶
Брокер из MemoryBroker::new() не хранит ничего. Сообщение живёт от публикации до чтения последним
подписчиком, поэтому долгоживущий сервис держит в памяти только ту работу, которую обработчики ещё
не сделали, сколько бы сообщений через него ни прошло.
Сервису с перемоткой нужна история, и он называет её объём:
// 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; у брокера с отдельной транзакционной
конфигурацией это имя указывает на другую политику.
Тот же 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:
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())
}
}
}
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 показывает то, что брокер хранит: у хранящего - всё в
пределах его ограничения, у брокера по умолчанию - ничего.