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

RustStream

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

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

RustStream подписывает Rust-сервис на потоки событий и публикует в них сообщения. Сервис при этом не привязан к одному брокеру сообщений. Ядро - это трейты и рантайм с роутером. Вместе с ядром поставляются кодеки, генерация AsyncAPI, метрики Prometheus и набор проверок conformance для авторов брокеров.

Фреймворк определяют два архитектурных обязательства:

  1. Полноценный интерфейс для сторонних брокеров. Ядро содержит только трейты и типы, ни одной зависимости от брокера. Каждый брокер - самостоятельный крейт. Соблюдение контракта проверяет набор conformance.
  2. Конфигурация брокера остаётся в его крейте. В ядре нет ни настроек, ни умолчаний, привязанных к конкретному брокеру. Каждый крейт брокера объявляет свой Config. Поэтому изменение на стороне брокера затрагивает только его крейт, а не фреймворк.
examples/quickstart.rs
//! The landing-page example: a one-handler service with no runtime boilerplate.
//!
//! ```text
//! cargo run --example quickstart --features macros,memory,json -- run
//! ```

use ruststream::memory::prelude::*;
use serde::Deserialize;

#[derive(Debug, Deserialize)]
struct Order {
    id: u64,
}

#[subscriber("orders")]
async fn handle(order: &Order) -> HandlerOutcome {
    println!("got order {}", order.id);
    HandlerOutcome::ack()
}

#[ruststream::app]
fn app() -> RustStream {
    RustStream::new(AppInfo::new("orders", "0.1.0")).with_broker(MemoryBroker::new(), |b| {
        b.include(handle);
    })
}
examples/manual/quickstart.rs
//! The landing-page example written without the `macros` feature: the handler is a named type
//! with an `impl Handle`, mounted with the `subscriber` constructor, and `main` is hand-written.
//!
//! ```text
//! cargo run --example manual_quickstart --no-default-features --features memory,json
//! ```

use std::error::Error;
use std::future::{Future, ready};

use ruststream::memory::prelude::*;
use serde::Deserialize;

#[derive(Debug, Deserialize, schemars::JsonSchema)]
struct Order {
    id: u64,
}

/// The handler: `#[subscriber("orders")]` generates this type and this impl. Every axis of the
/// form - the reply, the injections, the broker context, the application state - is a defaulted
/// parameter of `Handle`, so a plain body names none of them.
struct Receive;

impl Handle<Order> for Receive {
    fn handle(
        &self,
        order: &Order,
        _outs: &(),
        _ctx: &mut Context<'_>,
    ) -> impl Future<Output = Result<(), HandlerOutcome>> {
        println!("got order {}", order.id);
        ready(Ok(()))
    }
}

fn app() -> RustStream {
    RustStream::new(AppInfo::new("orders", "0.1.0")).with_broker(MemoryBroker::new(), |b| {
        b.include(subscriber("orders", Receive).build());
    })
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> {
    app().run().await?;
    Ok(())
}

#[ruststream::app] генерирует main со всем шаблонным кодом рантайма. Поэтому cargo run -- run запускает сервис, а cargo run -- asyncapi gen печатает его AsyncAPI-документ.

Принципы устройства

  • Полностью асинхронный, на tokio. В публичном API нет блокирующих вызовов.
  • Обобщённое ядро, никакого dyn в контракте. Контракт построен на ассоциированных типах и нативном async fn in trait. Стирание типов, если оно нужно сервису, выполняет рантайм.
  • Подписчики - это Stream, а не колбэки. Обратное давление обеспечивает сам Stream. Колбэки надстраивает рантайм.
  • Ack поглощает self. Второй ack - ошибка компиляции.
  • Трейты-совместимости для необязательных возможностей. BatchSubscriber, TransactionalPublisher, RequestReply, Partitioned и Seekable не входят в обязательный интерфейс.

Куда идти дальше

Что входит в этот репозиторий

Этот сайт документирует ruststream - ядро, не зависящее от брокера. Конкретные брокеры (NATS, Kafka, RabbitMQ, Redis, MQTT и другие) поставляются отдельными крейтами. Каждый такой крейт подключает ruststream с crates.io.

Справочник по Rust API опубликован на docs.rs - см. Справочник API.