NATS message system protocol codec and async client for MoonBit
Dependencies
| 模块 | 状态 | 说明 |
|---|---|---|
| 主题校验与通配符匹配 | ✅ 完成 | * 单层、> 多层通配 |
| 客户端命令编码 | ✅ 完成 | 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 个场景 |
moon add sa2360/moonnats// 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"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}")
}@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> 生成器 │
└─────────────────────────────────────────────┘moon check # 类型检查
moon test # 默认后端跑测试
moon test --target native # native 后端(需要 C 编译器)
moon test --target wasm-gc # wasm-gc 后端
moon run cmd/main # 协议层往返演示
moon fmt # 格式化(提交前必须跑)scoop install gccsetx MOON_CC "C:\Users\<你>\scoop\apps\gcc\current\bin\gcc.exe"#ifndef _CRT_RAND_S
#define _CRT_RAND_S
#endif这是 moon 运行时的上游问题,重装工具链后需要重新打一次。pub(all) enum ClientCommand {
Connect(String)
Pub(Publish)
Sub(Subscribe)
Unsub(Unsubscribe)
Ping
Pong
} derive(Eq, Debug)pub(all) enum ServerOp {
Info(String)
Msg(DeliveredMsg)
HMsg(DeliveredHMsg)
Ping
Pong
Ok
ErrOp(String)
} derive(Eq, Debug)fn is_valid_subject(subject : String, wildcards? : Bool) -> Boolfn subject_matches(subscription : String, subject : String) -> BoolInstall
Download zipNATS message system protocol codec and async client for MoonBit
Dependencies