moonnats

    NATS message system protocol codec and async client for MoonBit

    nats
    messaging
    protocol
    pubsub
    client
    Download zip
    Author
    Version
    0.1.1
    License
    Apache-2.0
    Last updated
    21 days ago
    Downloads
    9

    Dependencies

    #moonnats

    NATS message system protocol codec and async client for MoonBit.

    CI License mooncakes

    moonnats 为 MoonBit 提供 NATS 消息系统的完整协议实现: 一个零 I/O 依赖、可在任何后端(含 WASM)运行的协议编解码层,以及一个基于 moonbitlang/async 的异步客户端(已发布)。

    #项目背景

    NATS 是云原生生态广泛使用的轻量级消息系统,核心协议是基于 TCP 的文本协议, 常用于微服务通信、事件通知和 IoT 场景。本项目启动时 mooncakes.io 上没有任何 NATS 相关的包(Redis、Kafka 已有客户端,NATS 是空位),因此从零实现:

    • 协议层是纯函数式的字节流 ↔ 结构化消息转换,不碰网络,单测即全覆盖;
    • 客户端层负责连接、握手、订阅分发、心跳和 request-reply,只支持 native 后端。

    协议行为以 NATS 官方协议规范 为准,API 设计参考 nats.go(Apache-2.0), 不复制其实现代码。

    #当前状态

    模块状态说明
    主题校验与通配符匹配✅ 完成* 单层、> 多层通配
    客户端命令编码✅ 完成PUB / SUB / UNSUB / PING / CONNECT
    服务端帧流式解析✅ 完成TCP 粘包/半包安全,二进制 payload 安全
    INFO 握手解析✅ 完成JSON → ServerInfo
    HMSG 头部块解析✅ 完成NATS/1.0 版本行 + 键值对
    INBOX 主题生成✅ 完成request-reply 的地基
    异步客户端✅ 完成moonbitlang/async,native 后端
    CLI 工具✅ 完成moonnats_cli pub / sub,已对真实服务器端到端验证
    nats-server 集成测试✅ 完成CI 内对真实 nats-server 2.10 跑 10 个场景

    测试在 native 和 wasm-gc 双后端通过(84 个用例:协议层 42 + 客户端 27 + 集成 15),CI 在 Ubuntu 和 Windows 上运行完整矩阵,含真实 nats-server 集成测试和 CLI 端到端检查。

    明确不在首版范围内:JetStream、TLS、JWT/NKEY 认证、集群自动发现。

    #安装

    moon add sa2360/moonnats

    在包的 moon.pkg 中导入:

    // moon.pkg { "import": [ "sa2360/moonnats" ] }

    #使用

    #主题校验与通配符匹配

    // 发布主题不允许通配符
    @moonnats.is_valid_subject("foo.bar") // true
    @moonnats.is_valid_subject("foo.*") // false

    // 订阅主题允许 * 和 >
    @moonnats.is_valid_subject("foo.*", wildcards=true) // true
    @moonnats.is_valid_subject("foo.>", wildcards=true) // true
    @moonnats.is_valid_subject("foo.>.bar", wildcards=true) // false(> 必须在末尾)

    // 匹配:* 恰好一层,> 一层或多层
    @moonnats.subject_matches("foo.*", "foo.bar") // true
    @moonnats.subject_matches("foo.*", "foo.bar.baz") // false
    @moonnats.subject_matches("foo.>", "foo.bar.baz") // true
    @moonnats.subject_matches("foo.>", "foo") // false(> 至少吞一层)

    #编码客户端命令

    let cmd : @moonnats.ClientCommand = @moonnats.Pub(
    @moonnats.Publish::{ subject: "events", reply: None, payload: b"hi" },
    )
    let wire : Bytes = @moonnats.encode_command(cmd)
    // wire == b"PUB events 2\r\nhi\r\n"

    CONNECT 握手配置:

    let cfg = @moonnats.ConnectConfig::default(name=Some("my-service"))
    cfg.encode() // 完整 CONNECT 命令字节
    cfg.to_json() // 对应的 Json 值

    #解析服务端字节流

    let parser = @moonnats.Parser::new()
    parser.feed(chunk) // 每次 socket 读到什么都 feed 进去
    match parser.next_op() {
    @moonnats.NeedMore => () // 帧不完整,等下一次读
    @moonnats.Op(op, n) => // 一个完整 ServerOp,n 是消费的字节数
    match op {
    @moonnats.Msg(m) => handle(m.subject, m.payload)
    @moonnats.Ping => send_pong()
    _ => ()
    }
    @moonnats.Fail(reason) => log("protocol error: \{reason}")
    }

    长连接建议用 feed_and_compact,它会顺带回收已消费的缓冲区。

    #解析 INFO 和 HMSG 头部

    @moonnats.parse_server_info(json_text) // Result[ServerInfo, String]
    @moonnats.parse_header_block(headers_bytes) // Result[HeaderBlock, String]

    #架构

    ┌─────────────────────────────────────────────┐ │ 应用 / CLI(cmd/) │ ├─────────────────────────────────────────────┤ │ 客户端层(完成) │ │ 连接握手 · 订阅分发 · PING/PONG · req-rep │ │ 依赖 moonbitlang/async(仅 native) │ ├─────────────────────────────────────────────┤ │ 协议层(完成) │ │ protocol.mbt 消息数据模型 │ │ encode.mbt ClientCommand → Bytes │ │ parser.mbt Bytes → ServerOp(流式) │ │ connect.mbt CONNECT 配置与 JSON │ │ info.mbt INFO JSON → ServerInfo │ │ headers.mbt HMSG 头部块 → HeaderBlock │ │ subject.mbt 主题校验与通配符匹配 │ │ inbox.mbt _INBOX.<id>.<n> 生成器 │ └─────────────────────────────────────────────┘

    协议层设计要点:

    • MSG/HMSG 不能按行切:payload 是二进制且长度写在头里,解析器先解析头行、 按声明长度精确消费 payload 和终止 CRLF,剩余字节留给下一帧;
    • -ERR/INFO 的参数不能分词:错误描述和 JSON 都含空格,参数取整行原文;
    • 解析器是增量状态机:feed 进什么就缓冲什么,next_op 逐帧取出, 所有粘包/半包行为都有单测覆盖,不需要起服务器就能验证。

    #开发

    #环境要求

    • MoonBit 工具链(开发基于 0.1.20260904)
    • native 后端测试需要 C 编译器(Linux/macOS 自带;Windows 见下方说明)

    #常用命令

    moon check # 类型检查 moon test # 默认后端跑测试 moon test --target native # native 后端(需要 C 编译器) moon test --target wasm-gc # wasm-gc 后端 moon run cmd/main # 协议层往返演示 moon fmt # 格式化(提交前必须跑)

    #Windows(native 后端)注意事项

    1. msys64 自带的 gcc 可能是坏的(include 搜索路径为空)。用 scoop 装一个完好的 MinGW-w64:

      scoop install gcc

    2. 指定 moon 使用它(一次性配置):

      setx MOON_CC "C:\Users\<你>\scoop\apps\gcc\current\bin\gcc.exe"

    3. 当前 moon 运行时(0.1.20260904)在 mingw 下编译会报 rand_s 未声明—— windows.h 的包含链在 _CRT_RAND_S 定义前就拉进了 stdlib.h。 解决:编辑 ~/.moon/lib/runtime/env.c,把这段挪到文件顶部(任何 #include 之前,#include "moonbit.h" 之前):

      #ifndef _CRT_RAND_S #define _CRT_RAND_S #endif

      这是 moon 运行时的上游问题,重装工具链后需要重新打一次。

    #提交代码

    1. fork 仓库,从 master 拉分支;
    2. 改动保证 moon check 无错误、moon fmt 已跑、 native 和 wasm-gc 双后端测试全绿;
    3. 涉及协议行为的改动必须带测试用例(参考 parser_test.mbt 的写法);
    4. commit message 用 conventional 风格(feat: / fix: / docs: / test:), 一个提交做一件事;
    5. 发 PR,CI(双平台 × 双后端)通过后合并。

    #测试约定

    • 黑盒测试(*_test.mbt)放在包外视角,通过 @moonnats. 前缀访问公开 API, 公开行为都写在这里;
    • 白盒测试(*_wbtest.mbt)用于包内私有逻辑;
    • 协议测试的二进制输入用 ascii("...") 辅助函数构造,payload 里可以放心写 \r\n,测试的就是解析器对它的容忍度。

    #路线图

    #许可证

    Apache-2.0。协议行为遵循 NATS 官方协议规范; API 设计参考 nats.go(Apache-2.0), 未复制其实现代码。NATS 是 NATS.io 的商标,本项目是独立的社区实现,与 NATS.io 无隶属关系。

    ClientCommand

    pub(all) enum ClientCommand {
    Connect(String)
    Pub(Publish)
    Sub(Subscribe)
    Unsub(Unsubscribe)
    Ping
    Pong
    } derive(Eq,
    Debug
    )

    A client-to-server command.

    ConnectConfig

    pub(all) struct ConnectConfig {
    echo : Bool
    headers : Bool
    protocol : Int
    verbose : Bool
    no_echo : Bool
    user : String?
    pass : String?
    auth_token : String?
    name : String?
    lang : String
    version : String
    } derive(Eq,
    Debug
    )

    The JSON document sent as the CONNECT argument after receiving the server's INFO line. Field names and defaults follow the NATS protocol spec; unknown-to-the-server fields are ignored by it, and the server echoes the options it will honor in its next INFO.

    ConnectConfig::default

    fn ConnectConfig::default(name? : String?) -> ConnectConfig

    Defaults matching nats.go's behavior: echo on, headers on, protocol 1, verbose off, no_echo off.

    ConnectConfig::encode

    fn ConnectConfig::encode(self : ConnectConfig) -> Bytes

    The full CONNECT command bytes for this config.

    ConnectConfig::to_json

    fn ConnectConfig::to_json(self : ConnectConfig) -> Json

    Serialize to the JSON object expected after the CONNECT verb. Optional credentials are omitted entirely rather than sent as empty strings, as nats.go does.

    DeliveredHMsg

    pub(all) struct DeliveredHMsg {
    subject : String
    sid : String
    reply : String?
    headers : Bytes
    payload : Bytes
    } derive(Eq,
    Debug
    )

    A message delivered by the server via HMSG. The headers block is the raw NATS/1.0\r\nKey: Value\r\n\r\n prefix of the body; splitting it into key/value pairs is a separate concern.

    DeliveredMsg

    pub(all) struct DeliveredMsg {
    subject : String
    sid : String
    reply : String?
    payload : Bytes
    } derive(Eq,
    Debug
    )

    A message delivered by the server via MSG.
    pub(all) struct Header {
    name : String
    value : String
    } derive(Eq,
    Debug
    )

    One decoded header of an HMSG message.

    HeaderBlock

    pub(all) struct HeaderBlock {
    version_line : String
    headers : Array[Header]
    } derive(Eq,
    Debug
    )

    The headers block of an HMSG body: a NATS/1.0 line optionally followed by an inline status, then CRLF-separated Name: Value pairs ended by an empty line.

    NATS/1.0 100 Idle Heartbeat\r\nK: V\r\n\r\n

    InboxGen

    pub struct InboxGen {
    // private fields
    } derive(Eq,
    Debug
    )

    Generates unique reply subjects for request-reply calls, following the _INBOX.<id>.<counter> convention used by nats.go. One generator per connection: the prefix identifies the connection (the client supplies a random id at startup) and the counter disambiguates in-flight requests.

    InboxGen::new

    fn InboxGen::new(prefix? : String) -> InboxGen

    InboxGen::next

    fn InboxGen::next(self : InboxGen) -> String

    Next inbox subject; the counter restarts at 1.

    ParseResult

    pub(all) enum ParseResult {
    NeedMore
    Op(ServerOp, Int)
    Fail(String)
    } derive(Eq,
    Debug
    )

    Result of trying to parse one operation from the buffered bytes.

    Parser

    pub struct Parser {
    // private fields
    } derive(Eq,
    Debug
    )

    Incremental parser for the server side of the NATS wire protocol.

    Bytes are appended with feed as TCP delivers them — possibly split in the middle of a frame — and complete operations are pulled out with next_op. The parser is a state machine over one growable buffer; no allocation happens per fed chunk beyond the buffer growth itself.

    Parser::buffered

    fn Parser::buffered(self : Parser) -> Int

    Bytes currently buffered but not yet consumed.

    Parser::feed

    fn Parser::feed(self : Parser, data : Bytes) -> Unit

    Append bytes delivered by the transport.

    Parser::feed_and_compact

    fn Parser::feed_and_compact(self : Parser, data : Bytes) -> Unit

    Drop consumed bytes so the buffer cannot grow without bound on long-lived connections. Called by feed to piggyback compaction on the next write.

    Parser::new

    fn Parser::new() -> Parser

    Parser::next_op

    fn Parser::next_op(self : Parser) -> ParseResult

    Try to parse the next complete operation. Consumes it on success; leaves the buffer untouched on NeedMore.

    Publish

    pub(all) struct Publish {
    subject : String
    reply : String?
    payload : Bytes
    } derive(Eq,
    Debug
    )

    Arguments of PUB <subject> [reply-to] <#bytes>.

    ServerInfo

    pub(all) struct ServerInfo {
    server_id : String
    server_name : String
    version : String
    host : String
    port : Int
    auth_required : Bool
    headers : Bool
    max_payload : Int
    proto : Int
    } derive(Eq,
    Debug
    )

    The server information advertised in the INFO handshake line and on every cluster topology change. Fields the client acts on are decoded; the rest is available through raw.

    ServerOp

    pub(all) enum ServerOp {
    Info(String)
    Msg(DeliveredMsg)
    HMsg(DeliveredHMsg)
    Ping
    Pong
    Ok
    ErrOp(String)
    } derive(Eq,
    Debug
    )

    A server-to-client operation on the NATS wire protocol.

    The text protocol frames every operation as one line terminated by \r\n; delivered messages additionally carry a binary body whose length is declared in the header line.

    Subscribe

    pub(all) struct Subscribe {
    subject : String
    sid : String
    queue_group : String?
    } derive(Eq,
    Debug
    )

    Arguments of SUB <subject> [queue-group] <sid>.

    Unsubscribe

    pub(all) struct Unsubscribe {
    sid : String
    max_msgs : Int?
    } derive(Eq,
    Debug
    )

    Arguments of UNSUB <sid> [max-msgs].

    encode_command

    fn encode_command(cmd : ClientCommand) -> Bytes

    Encode a client command into the exact bytes to write to the wire.

    Every command is one \r\n-terminated line; Pub appends the payload followed by a second \r\n as the protocol requires. The header line is ASCII by protocol construction, so the line-to-bytes mapping is lossless; payload bytes are concatenated, never routed through a string.

    is_valid_subject

    fn is_valid_subject(subject : String, wildcards? : Bool) -> Bool

    Validate a NATS subject.

    A subject is dot-separated tokens of printable ASCII characters without spaces. Tokens must be non-empty, so leading/trailing dots and .. are invalid.

    When wildcards is false (publish subjects) neither * nor > may appear. When true, * must form an entire token and > must be the final token.

    parse_header_block

    fn parse_header_block(raw : Bytes) -> Result[HeaderBlock, String]

    Decode a raw headers block. Tolerates LF-only line endings, as nats-server itself produces CRLF but user headers may be re-encoded by intermediaries.

    parse_server_info

    fn parse_server_info(json : String) -> Result[ServerInfo, String]

    Decode the JSON argument of an INFO operation. Returns an error when the JSON is malformed or missing the fields the client cannot proceed without (server_id, version, host, port).

    ping_command

    fn ping_command() -> Bytes

    Keepalive PING command bytes.

    pong_command

    fn pong_command() -> Bytes

    Keepalive PONG command bytes (reply to a server PING).

    server_info_from_op

    fn server_info_from_op(op : ServerOp) -> Result[ServerInfo, String]

    Decode the INFO carried by a parsed ServerOp::Info.

    subject_matches

    fn subject_matches(subscription : String, subject : String) -> Bool

    Report whether a subscription subject (wildcards allowed) matches a concrete publish subject (wildcards forbidden).

    * matches exactly one token; a trailing > matches one or more remaining tokens, so foo.> does not match foo. Neither input is validated here; pair with is_valid_subject at the API boundary.