For AI agents: the complete documentation index is available at https://a3s-lab.github.io/Boot/llms.txt, the full documentation bundle is available at https://a3s-lab.github.io/Boot/llms-full.txt, and this page is available as Markdown at https://a3s-lab.github.io/Boot/protocols/microservices.md.
  • 简体中文
  • v0.2.0
  • 消息服务与 Transport

    Boot 的 microservice 模型把 handler 与具体 broker 解耦,但不会抹平协议差异。MessagePatternDefinition 描述 request-response 或 event-only pattern,MessageTransport 负责启动、停止、dispatch 与 client 契约。

    消息 Controller

    use a3s_boot::{message_body, message_controller, message_pattern, Result};
    
    #[derive(Debug)]
    struct MathController;
    
    #[message_controller]
    impl MathController {
        #[message_pattern("sum")]
        async fn sum(&self, #[message_body] values: Vec<i64>) -> Result<i64> {
            Ok(values.into_iter().sum())
        }
    }

    Module 使用 message_controllers = [MathController] 注册宏 Controller。也可以直接创建 MessagePatternDefinition。

    Request-response pattern 返回实现 IntoTransportReply 的结果。Event pattern 表达没有业务 reply 的单向处理。具体 transport 仍可能产生协议级 acknowledgement 或 error envelope。

    运行模式

    let mut service = BootFactory::create_microservice(
        AppModule,
        TcpTransport::new(TcpTransportOptions::new("127.0.0.1:4000")),
    )?;
    service.listen().await?;

    应用也可以在 HTTP shell 上连接一个或多个 microservice,形成 hybrid application。Provider-only worker 则可以直接注入 transport client,而不监听 HTTP。

    选择 factory 的同步或 async 变体取决于 Provider 构建是否包含 async factory,不取决于 handler 是否 async。

    可用实现

    TransportFeature主要 wire 模型
    InProcessTransport无同进程 dispatch,适合测试与 worker 内通信
    TcpTransporttcp-transportnewline-delimited JSON frame
    RedisTransportredis-transportPub/Sub channel
    NatsTransportnats-transportrequest/reply 与 event subject
    MqttTransportmqtt-transportrequest/reply 与 event topic,显式 QoS
    RabbitMqTransportrabbitmq-transportrequest/reply 与 event queue
    KafkaTransportkafka-transportrequest/reply 与 event topic
    GrpcTransportgrpc-transportunary request/reply 与 event call

    Feature 只加入实现和依赖。Broker 地址、TLS、credential、topic 或 queue 创建、retention、replication 与监控仍由部署配置。

    共享执行能力

    每次 dispatch 创建 TransportContext 和 scope。管线支持:

    • payload Pipe 与类型化 validation
    • TransportGuard
    • around TransportInterceptor
    • TransportExceptionFilter
    • application、Controller 和 pattern metadata
    • scoped Provider 与 Provider-backed handler

    宏 #[use_guard]、#[use_pipe]、#[use_interceptor]、#[use_filter] 与 #[metadata] 可放在消息 Controller 或 pattern。全局组件使用 use_global_transport_* builder。

    Transport error envelope 使用与 HTTP BootError 相同的 status 与 error kind 映射,但仍是消息 reply,不应被当成真实 HTTP response。

    Pattern 与 schema

    Pattern name 是协议契约。为跨服务消息选择稳定、带领域前缀的名称,例如 billing.invoice.create。Payload 使用显式 version 字段或新的 pattern 演进,避免在已有 consumer 不知情时改变字段含义。

    Event handler 应接受可能重复的交付。Request-response client 必须配置 timeout,并区分业务错误、远端 transport error、连接失败与超时。

    不同协议不能假装相同

    MessageTransport 统一调用表面,不统一 durability:

    • Redis Pub/Sub 不等同于持久队列。
    • MQTT QoS 不自动提供业务幂等性。
    • Kafka offset、consumer group 与 partition ordering 仍属于 Kafka 策略。
    • RabbitMQ ack、requeue 与 dead-letter 配置影响交付。
    • NATS Core 与其他 NATS persistence 模式不能混为一谈。
    • TCP frame 需要应用自己定义连接恢复与 request correlation 生命周期。

    在选择 transport 前记录 delivery、ordering、replay、backpressure、maximum payload 与 failure recovery 要求。

    测试与停机

    先用 InProcessTransport 验证 pattern、scope、Guard、Interceptor、Pipe、Filter 和 reply。再为实际 transport 写 integration test,覆盖连接、序列化、timeout、duplicate、broker restart 与 shutdown。

    Graceful shutdown 应停止接收新消息,有界等待 active handler,然后按协议提交或释放未完成交付。不要在进程结束时默认认为消息已被可靠确认。