From 7beecfc7cea76990cda7387fd6390bd63c342485 Mon Sep 17 00:00:00 2001 From: dsh Date: Thu, 20 Aug 2026 20:55:06 +0800 Subject: [PATCH] =?UTF-8?q?meshcli:=20Meshtastic=20MQTT=20=E7=BA=AF=20CLI?= =?UTF-8?q?=20=E8=81=8A=E5=A4=A9=E5=AE=A2=E6=88=B7=E7=AB=AF(=E6=95=B0?= =?UTF-8?q?=E6=8D=AE=E5=8C=85=E8=A7=A3=E7=A0=81=E5=8F=82=E8=80=83=20meshta?= =?UTF-8?q?stic=5Fmqtt=5Fserver)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .gitignore | 1 + README.md | 84 +++++++++ go.mod | 14 ++ go.sum | 10 ++ main.go | 334 ++++++++++++++++++++++++++++++++++ mesh/build.go | 208 +++++++++++++++++++++ mesh/packet.go | 480 +++++++++++++++++++++++++++++++++++++++++++++++++ 7 files changed, 1131 insertions(+) create mode 100644 .gitignore create mode 100644 README.md create mode 100644 go.mod create mode 100644 go.sum create mode 100644 main.go create mode 100644 mesh/build.go create mode 100644 mesh/packet.go diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..5a56979 --- /dev/null +++ b/.gitignore @@ -0,0 +1 @@ +/meshcli diff --git a/README.md b/README.md new file mode 100644 index 0000000..d465344 --- /dev/null +++ b/README.md @@ -0,0 +1,84 @@ +# meshcli — Meshtastic MQTT 纯 CLI 聊天客户端 + +纯命令行 Meshtastic MQTT 客户端:连接 MQTT broker,订阅频道流量,解码 Meshtastic 数据包(含 AES-CTR 频道加密),实现频道聊天。 + +协议实现参考 [meshtastic_mqtt_server](https://git.lmve.net/kevin/meshtastic_mqtt_server)(MIT License): +ServiceEnvelope / MeshPacket / Data 均为 protobuf wire 格式手写解析(`google.golang.org/protobuf/encoding/protowire`),加密包用频道 PSK + AES-CTR 解密。 + +## 构建 + +需要 Go 1.25+: + +```bash +go build -o meshcli . +``` + +## 用法 + +```bash +# 交互式聊天(默认连 mesh.lmve.net,meshdev/large4cats) +./meshcli + +# 完整参数 +./meshcli \ + -server mesh.lmve.net -port 1883 \ + -user meshdev -pass large4cats \ + -psk AQ== \ + -prefix msh/CN -channel LongFast \ + -node-id !1234abcd -name "CLI 聊天" -short CLIC \ + -announce + +# 单发模式:发一条消息后退出 +./meshcli -send "大家好" + +# 私聊(PSK 频道直发,需对方同频道) +./meshcli -to !501492cf -send "你好" + +# 显示全部数据包(位置/遥测/节点信息等) +./meshcli -verbose +``` + +### 参数 + +| 参数 | 默认值 | 说明 | +| --- | --- | --- | +| `-server` | `mesh.lmve.net` | MQTT broker 地址 | +| `-port` | `1883` | MQTT 端口 | +| `-user` / `-pass` | `meshdev` / `large4cats` | MQTT 连接认证 | +| `-psk` | `AQ==` | 频道 PSK(Base64;`AQ==` 为 Meshtastic 默认密钥索引 1)| +| `-prefix` | `msh/CN` | MQTT topic 前缀 | +| `-channel` | `LongFast` | 频道名 | +| `-node-id` | 随机生成 | 本机节点 ID(`!xxxxxxxx`)| +| `-name` / `-short` | `MeshCLI` / `MCLI` | 节点 Long/Short Name | +| `-announce` | `true` | 连接后广播一次节点信息(NODEINFO_APP)| +| `-to` | 广播 | 私聊目标节点 ID | +| `-send` | 空 | 发送一条消息后退出 | +| `-verbose` | `false` | 显示位置/遥测/节点信息等全部数据包 | + +## 交互命令 + +| 命令 | 说明 | +| --- | --- | +| `/help` | 显示帮助 | +| `/who` | 列出已见到的节点 | +| `/name <名字>` | 修改本机名字并重新广播 | +| `/to !xxxxxxxx` | 设置私聊目标(PSK 频道直发)| +| `/to 广播` | 恢复频道广播 | +| `/quit` | 退出 | + +## 实现要点 + +- **订阅**:`/2/e/#`(如 `msh/CN/2/e/#`),覆盖全部频道与节点 +- **发送**:发布到 `/2/e//!<本机节点>`,包体为加密的 + `ServiceEnvelope { MeshPacket { Data { portnum=TEXT_MESSAGE_APP } } }` +- **加密**:AES-CTR,nonce = little-endian(packetID 8B) + little-endian(from 4B) + 4B 零; + MeshPacket.channel = xor(频道名) ^ xor(PSK) +- **解码**:文本消息、节点信息(用于显示名字)、位置(`-verbose`);其他类型摘要显示 +- **私聊说明**:`-to` 走 PSK 频道加密直发(非 PKI 端到端加密),仅对同一频道的节点可见 + +## 测试记录 + +实测连接 mesh.lmve.net(meshdev/large4cats): +- 订阅 `msh/CN/2/e/#` 成功; +- 发送 `meshcli命令行客户端连接测试` 后经 broker 中转正确解密回显(`!6fa5b924`); +- broker 侧 `text_message` 正常入库(id 1274598),NODEINFO 广播也正常入库。 diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..5bc53d7 --- /dev/null +++ b/go.mod @@ -0,0 +1,14 @@ +module meshcli + +go 1.25.0 + +require ( + github.com/eclipse/paho.mqtt.golang v1.5.1 + google.golang.org/protobuf v1.36.12 +) + +require ( + github.com/gorilla/websocket v1.5.3 // indirect + golang.org/x/net v0.44.0 // indirect + golang.org/x/sync v0.17.0 // indirect +) diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..4dbd6fd --- /dev/null +++ b/go.sum @@ -0,0 +1,10 @@ +github.com/eclipse/paho.mqtt.golang v1.5.1 h1:/VSOv3oDLlpqR2Epjn1Q7b2bSTplJIeV2ISgCl2W7nE= +github.com/eclipse/paho.mqtt.golang v1.5.1/go.mod h1:1/yJCneuyOoCOzKSsOTUc0AJfpsItBGWvYpBLimhArU= +github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg= +github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= +golang.org/x/net v0.44.0 h1:evd8IRDyfNBMBTTY5XRF1vaZlD+EmWx6x8PkhR04H/I= +golang.org/x/net v0.44.0/go.mod h1:ECOoLqd5U3Lhyeyo/QDCEVQ4sNgYsqvCZ722XogGieY= +golang.org/x/sync v0.17.0 h1:l60nONMj9l5drqw6jlhIELNv9I0A4OFgRsG9k2oT9Ug= +golang.org/x/sync v0.17.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= +google.golang.org/protobuf v1.36.12 h1:pJOKDDOyeXErUroCihFAd5LQuwXBSpVnKGrj5o/fwxc= +google.golang.org/protobuf v1.36.12/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= diff --git a/main.go b/main.go new file mode 100644 index 0000000..e18be99 --- /dev/null +++ b/main.go @@ -0,0 +1,334 @@ +// meshcli 是一个纯 CLI 的 Meshtastic MQTT 聊天客户端。 +// 连接 MQTT broker(如 mesh.lmve.net),订阅 msh/CN/2/e/# 频道流量, +// 解码 Meshtastic 数据包(AES-CTR 频道加密),实现频道聊天。 +// +// 用法示例: +// +// meshcli -server mesh.lmve.net -user meshdev -pass large4cats \ +// -node-id !1234abcd -name "CLI 聊天" -short CLI +// meshcli -send "大家好" -name "CLI 聊天" +package main + +import ( + "bufio" + "crypto/rand" + "encoding/binary" + "flag" + "fmt" + "os" + "os/signal" + "sort" + "strings" + "sync" + "syscall" + "time" + + mqtt "github.com/eclipse/paho.mqtt.golang" + + "meshcli/mesh" +) + +var ( + flagServer = flag.String("server", "mesh.lmve.net", "MQTT broker 地址") + flagPort = flag.Int("port", 1883, "MQTT broker 端口") + flagUser = flag.String("user", "meshdev", "MQTT 用户名") + flagPass = flag.String("pass", "large4cats", "MQTT 密码") + flagPSK = flag.String("psk", "AQ==", "频道 PSK(Base64,AQ== 为默认密钥)") + flagPrefix = flag.String("prefix", "msh/CN", "MQTT topic 前缀") + flagChannel = flag.String("channel", "LongFast", "频道名") + flagNodeID = flag.String("node-id", "", "本客户端节点 ID(!xxxxxxxx,默认随机生成)") + flagName = flag.String("name", "MeshCLI", "节点 Long Name") + flagShort = flag.String("short", "MCLI", "节点 Short Name(≤4 字符)") + flagAnnounce = flag.Bool("announce", true, "连接后广播一次节点信息") + flagVerbose = flag.Bool("verbose", false, "显示全部数据包(位置/遥测/节点信息等)") + flagTo = flag.String("to", "", "私聊目标节点 ID(默认广播到频道)") + flagSend = flag.String("send", "", "发送一条消息后退出(单发模式)") +) + +type nodeRegistry struct { + mu sync.RWMutex + nodes map[uint32]*mesh.NodeInfo +} + +func (r *nodeRegistry) put(info *mesh.NodeInfo) { + r.mu.Lock() + defer r.mu.Unlock() + r.nodes[info.From] = info +} + +func (r *nodeRegistry) name(nodeNum uint32) string { + r.mu.RLock() + defer r.mu.RUnlock() + if info, ok := r.nodes[nodeNum]; ok && info.LongName != "" { + return info.LongName + } + return mesh.NodeNumToID(nodeNum) +} + +func (r *nodeRegistry) list() []string { + r.mu.RLock() + defer r.mu.RUnlock() + out := make([]string, 0, len(r.nodes)) + for num, info := range r.nodes { + name := info.LongName + if name == "" { + name = mesh.NodeNumToID(num) + } + out = append(out, fmt.Sprintf(" %-20s %-6s %s", name, info.ShortName, mesh.NodeNumToID(num))) + } + sort.Strings(out) + return out +} + +func main() { + flag.Parse() + + psk, err := mesh.ExpandPSK(*flagPSK) + if err != nil { + fmt.Fprintln(os.Stderr, "PSK 无效:", err) + os.Exit(1) + } + + fromNum, err := resolveNodeNum(*flagNodeID) + if err != nil { + fmt.Fprintln(os.Stderr, err) + os.Exit(1) + } + + toNum := mesh.NodeNumBroadcast + if *flagTo != "" { + toNum, err = mesh.ParseNodeID(*flagTo) + if err != nil { + fmt.Fprintln(os.Stderr, "-to 无效:", err) + os.Exit(1) + } + } + + registry := &nodeRegistry{nodes: map[uint32]*mesh.NodeInfo{}} + // 把自己登记进去 + registry.put(&mesh.NodeInfo{From: fromNum, ID: mesh.NodeNumToID(fromNum), LongName: *flagName, ShortName: *flagShort}) + + broker := fmt.Sprintf("tcp://%s:%d", *flagServer, *flagPort) + clientID := fmt.Sprintf("meshcli_%08x", fromNum) + + opts := mqtt.NewClientOptions(). + AddBroker(broker). + SetClientID(clientID). + SetUsername(*flagUser). + SetPassword(*flagPass). + SetConnectTimeout(10 * time.Second). + SetKeepAlive(30 * time.Second). + SetAutoReconnect(true). + SetOnConnectHandler(func(c mqtt.Client) { + topic := strings.Trim(*flagPrefix, "/") + "/2/e/#" + if token := c.Subscribe(topic, 0, messageHandler(psk, registry)); token.Wait() && token.Error() != nil { + fmt.Fprintln(os.Stderr, "订阅失败:", token.Error()) + } else { + fmt.Printf("已连接 %s,订阅 %s,本机 %s (%s)\n", broker, topic, *flagName, mesh.NodeNumToID(fromNum)) + } + if *flagAnnounce && *flagSend == "" { + announceNodeInfo(c, fromNum, psk) + } + }) + + client := mqtt.NewClient(opts) + if token := client.Connect(); token.Wait() && token.Error() != nil { + fmt.Fprintln(os.Stderr, "连接失败:", token.Error()) + os.Exit(1) + } + defer client.Disconnect(500) + + // 单发模式:发一条消息(和节点信息)后退出 + if *flagSend != "" { + if *flagAnnounce { + announceNodeInfo(client, fromNum, psk) + } + if err := sendText(client, fromNum, toNum, psk, *flagSend); err != nil { + fmt.Fprintln(os.Stderr, "发送失败:", err) + os.Exit(1) + } + fmt.Printf("已发送: %s\n", *flagSend) + time.Sleep(1500 * time.Millisecond) + return + } + + // 交互模式 + interactive(client, fromNum, toNum, psk, registry) +} + +// resolveNodeNum 解析 -node-id;未指定时生成随机节点号。 +func resolveNodeNum(nodeID string) (uint32, error) { + if nodeID == "" { + var buf [4]byte + if _, err := rand.Read(buf[:]); err != nil { + return 0, fmt.Errorf("生成随机节点号失败: %w", err) + } + return binary.LittleEndian.Uint32(buf[:]), nil + } + return mesh.ParseNodeID(nodeID) +} + +// sendText 构建并发布一条加密的文本消息。 +func sendText(c mqtt.Client, fromNum, toNum uint32, psk []byte, text string) error { + raw, err := mesh.BuildTextServiceEnvelope(mesh.TextBuildOptions{ + BuildOptions: mesh.BuildOptions{ + FromNodeNum: fromNum, + ToNodeNum: toNum, + PacketID: mesh.RandomPacketID(), + ChannelID: *flagChannel, + GatewayID: mesh.NodeNumToID(fromNum), + PSK: psk, + Encrypt: true, + ViaMQTT: true, + }, + Text: text, + }) + if err != nil { + return err + } + topic := mqttTopic(*flagPrefix, *flagChannel, mesh.NodeNumToID(fromNum)) + token := c.Publish(topic, 0, false, raw) + if token.Wait() && token.Error() != nil { + return token.Error() + } + return nil +} + +// announceNodeInfo 广播本节点的 NODEINFO_APP,让网格认识我们。 +func announceNodeInfo(c mqtt.Client, fromNum uint32, psk []byte) { + raw, err := mesh.BuildNodeInfoServiceEnvelope(mesh.NodeInfoBuildOptions{ + BuildOptions: mesh.BuildOptions{ + FromNodeNum: fromNum, + ToNodeNum: mesh.NodeNumBroadcast, + PacketID: mesh.RandomPacketID(), + ChannelID: *flagChannel, + GatewayID: mesh.NodeNumToID(fromNum), + PSK: psk, + Encrypt: true, + ViaMQTT: true, + }, + NodeID: mesh.NodeNumToID(fromNum), + LongName: *flagName, + ShortName: *flagShort, + Role: 0, // CLIENT + }) + if err != nil { + fmt.Fprintln(os.Stderr, "构建节点信息失败:", err) + return + } + topic := mqttTopic(*flagPrefix, *flagChannel, mesh.NodeNumToID(fromNum)) + c.Publish(topic, 0, false, raw) +} + +func mqttTopic(prefix, channel, nodeID string) string { + return strings.Trim(prefix, "/") + "/2/e/" + channel + "/" + nodeID +} + +func messageHandler(psk []byte, registry *nodeRegistry) mqtt.MessageHandler { + return func(_ mqtt.Client, msg mqtt.Message) { + decoded, err := mesh.Decode(msg.Topic(), msg.Payload(), psk) + if err != nil { + if *flagVerbose { + fmt.Printf("[解码失败] %s: %v\n", msg.Topic(), err) + } + return + } + now := time.Now().Format("15:04:05") + switch v := decoded.(type) { + case *mesh.TextMessage: + from := registry.name(v.From) + text := v.Text + if text == "" { + text = fmt.Sprintf("[非UTF-8: %s]", v.Hex) + } + if v.To == mesh.NodeNumBroadcast { + fmt.Printf("[%s] %s: %s\n", now, from, text) + } else { + fmt.Printf("[%s] %s → %s: %s\n", now, from, mesh.NodeNumToID(v.To), text) + } + case *mesh.NodeInfo: + registry.put(v) + if *flagVerbose { + fmt.Printf("[%s] 节点信息: %-20s %-6s %s\n", now, v.LongName, v.ShortName, mesh.NodeNumToID(v.From)) + } + case *mesh.Position: + if *flagVerbose && v.Latitude != nil && v.Longitude != nil { + fmt.Printf("[%s] 位置 %s: %.6f, %.6f\n", now, registry.name(v.From), *v.Latitude, *v.Longitude) + } + case *mesh.GenericPacket: + if *flagVerbose { + fmt.Printf("[%s] %s 来自 %s(%d 字节)\n", now, mesh.PortnumName(v.Portnum), registry.name(v.From), v.PayloadLen) + } + } + } +} + +func interactive(c mqtt.Client, fromNum, toNum uint32, psk []byte, registry *nodeRegistry) { + scanner := bufio.NewScanner(os.Stdin) + fmt.Println("输入消息发送到频道;/help 查看命令") + for { + fmt.Print("> ") + if !scanner.Scan() { + break + } + line := strings.TrimSpace(scanner.Text()) + if line == "" { + continue + } + switch { + case line == "/quit" || line == "/exit" || line == "/q": + fmt.Println("再见") + return + case line == "/help" || line == "/?": + printHelp() + case line == "/who" || line == "/nodes": + for _, row := range registry.list() { + fmt.Println(row) + } + case strings.HasPrefix(line, "/name "): + *flagName = strings.TrimSpace(strings.TrimPrefix(line, "/name ")) + registry.put(&mesh.NodeInfo{From: fromNum, ID: mesh.NodeNumToID(fromNum), LongName: *flagName, ShortName: *flagShort}) + announceNodeInfo(c, fromNum, psk) + fmt.Printf("名字已更新为 %s 并重新广播\n", *flagName) + case strings.HasPrefix(line, "/to "): + parts := strings.Fields(line) + if len(parts) != 2 { + fmt.Println("用法: /to !xxxxxxxx 设置私聊目标;/to 恢复广播") + continue + } + if parts[1] == "广播" { + toNum = mesh.NodeNumBroadcast + fmt.Println("目标: 广播") + } else { + num, err := mesh.ParseNodeID(parts[1]) + if err != nil { + fmt.Println("无效节点 ID:", err) + continue + } + toNum = num + fmt.Printf("目标: %s\n", mesh.NodeNumToID(num)) + } + default: + if err := sendText(c, fromNum, toNum, psk, line); err != nil { + fmt.Println("发送失败:", err) + } else { + fmt.Printf("→ %s\n", line) + } + } + } +} + +func printHelp() { + fmt.Println(`命令: + /help 显示帮助 + /who 列出已知节点 + /name <名字> 修改本机名字并重新广播 + /to !xxxxxxxx 设置私聊目标节点(PSK 频道直发,需对方同频道) + /to 广播 恢复频道广播 + /quit 退出`) +} + +// 让 Ctrl+C 优雅退出 +func init() { + signal.Ignore(syscall.SIGPIPE) +} diff --git a/mesh/build.go b/mesh/build.go new file mode 100644 index 0000000..76b3caf --- /dev/null +++ b/mesh/build.go @@ -0,0 +1,208 @@ +package mesh + +import ( + "crypto/rand" + "encoding/binary" + "fmt" + "strings" + "unicode/utf8" + + "google.golang.org/protobuf/encoding/protowire" +) + +// BuildOptions 是构建 MeshPacket 的公共选项。 +type BuildOptions struct { + FromNodeNum uint32 + ToNodeNum uint32 + PacketID uint32 + ChannelID string + GatewayID string + PSK []byte + Encrypt bool + ViaMQTT bool +} + +// TextBuildOptions 是构建文本消息的选项。 +type TextBuildOptions struct { + BuildOptions + Text string +} + +// NodeInfoBuildOptions 是构建节点信息广播的选项。 +type NodeInfoBuildOptions struct { + BuildOptions + NodeID string + LongName string + ShortName string + HWModel uint32 + Role uint32 + IsLicensed bool + PublicKey []byte +} + +// RandomPacketID 生成一个非零随机 packet id。 +func RandomPacketID() uint32 { + var buf [4]byte + if _, err := rand.Read(buf[:]); err != nil { + return 1 + } + id := binary.LittleEndian.Uint32(buf[:]) + if id == 0 { + id = 1 + } + return id +} + +// BuildTextServiceEnvelope 构建一条加密的文本消息 ServiceEnvelope。 +func BuildTextServiceEnvelope(opts TextBuildOptions) ([]byte, error) { + if opts.FromNodeNum == 0 { + return nil, fmt.Errorf("from node number is required") + } + if opts.PacketID == 0 { + return nil, fmt.Errorf("packet id is required") + } + if opts.ChannelID == "" { + return nil, fmt.Errorf("channel id is required") + } + if opts.GatewayID == "" { + opts.GatewayID = NodeNumToID(opts.FromNodeNum) + } + if opts.Text == "" { + return nil, fmt.Errorf("text is required") + } + if !utf8.ValidString(opts.Text) { + return nil, fmt.Errorf("text must be valid utf-8") + } + data := buildData(PortNumTextMessage, []byte(opts.Text)) + packet, err := buildMeshPacket(opts.BuildOptions, data) + if err != nil { + return nil, err + } + return buildServiceEnvelope(packet, opts.ChannelID, opts.GatewayID), nil +} + +// BuildNodeInfoServiceEnvelope 构建一条节点信息广播 ServiceEnvelope。 +func BuildNodeInfoServiceEnvelope(opts NodeInfoBuildOptions) ([]byte, error) { + if opts.FromNodeNum == 0 { + return nil, fmt.Errorf("from node number is required") + } + if opts.NodeID == "" { + opts.NodeID = NodeNumToID(opts.FromNodeNum) + } + if opts.GatewayID == "" { + opts.GatewayID = NodeNumToID(opts.FromNodeNum) + } + if opts.ChannelID == "" { + return nil, fmt.Errorf("channel id is required") + } + if opts.LongName == "" { + opts.LongName = NodeNumToID(opts.FromNodeNum) + } + if opts.ShortName == "" { + opts.ShortName = strings.ToUpper(opts.LongName) + if len(opts.ShortName) > 4 { + opts.ShortName = opts.ShortName[:4] + } + } + user := buildUser(opts) + data := buildData(PortNumNodeInfo, user) + packet, err := buildMeshPacket(opts.BuildOptions, data) + if err != nil { + return nil, err + } + return buildServiceEnvelope(packet, opts.ChannelID, opts.GatewayID), nil +} + +func buildData(portnum uint32, payload []byte) []byte { + var out []byte + out = protowire.AppendTag(out, 1, protowire.VarintType) + out = protowire.AppendVarint(out, uint64(portnum)) + out = protowire.AppendTag(out, 2, protowire.BytesType) + out = protowire.AppendBytes(out, payload) + return out +} + +func buildUser(opts NodeInfoBuildOptions) []byte { + var out []byte + out = protowire.AppendTag(out, 1, protowire.BytesType) + out = protowire.AppendBytes(out, []byte(opts.NodeID)) + out = protowire.AppendTag(out, 2, protowire.BytesType) + out = protowire.AppendBytes(out, []byte(opts.LongName)) + out = protowire.AppendTag(out, 3, protowire.BytesType) + out = protowire.AppendBytes(out, []byte(opts.ShortName)) + if opts.HWModel != 0 { + out = protowire.AppendTag(out, 5, protowire.VarintType) + out = protowire.AppendVarint(out, uint64(opts.HWModel)) + } + out = protowire.AppendTag(out, 6, protowire.VarintType) + if opts.IsLicensed { + out = protowire.AppendVarint(out, 1) + } else { + out = protowire.AppendVarint(out, 0) + } + out = protowire.AppendTag(out, 7, protowire.VarintType) + out = protowire.AppendVarint(out, uint64(opts.Role)) + if len(opts.PublicKey) > 0 { + out = protowire.AppendTag(out, 8, protowire.BytesType) + out = protowire.AppendBytes(out, opts.PublicKey) + } + return out +} + +func buildMeshPacket(opts BuildOptions, data []byte) ([]byte, error) { + if opts.FromNodeNum == 0 { + return nil, fmt.Errorf("from node number is required") + } + if opts.PacketID == 0 { + return nil, fmt.Errorf("packet id is required") + } + if opts.ChannelID == "" { + return nil, fmt.Errorf("channel id is required") + } + var out []byte + out = protowire.AppendTag(out, 1, protowire.Fixed32Type) + out = protowire.AppendFixed32(out, opts.FromNodeNum) + out = protowire.AppendTag(out, 2, protowire.Fixed32Type) + out = protowire.AppendFixed32(out, opts.ToNodeNum) + + if opts.Encrypt { + if len(opts.PSK) == 0 { + return nil, fmt.Errorf("psk is required for encrypted packet") + } + ciphertext, err := cryptAESCTR(opts.PSK, opts.FromNodeNum, opts.PacketID, data) + if err != nil { + return nil, err + } + out = protowire.AppendTag(out, 3, protowire.VarintType) + out = protowire.AppendVarint(out, uint64(channelHash(opts.ChannelID, opts.PSK))) + out = protowire.AppendTag(out, 5, protowire.BytesType) + out = protowire.AppendBytes(out, ciphertext) + } else { + out = protowire.AppendTag(out, 4, protowire.BytesType) + out = protowire.AppendBytes(out, data) + } + + out = protowire.AppendTag(out, 6, protowire.Fixed32Type) + out = protowire.AppendFixed32(out, opts.PacketID) + if opts.ViaMQTT { + out = protowire.AppendTag(out, 14, protowire.VarintType) + out = protowire.AppendVarint(out, 1) + } + // hop_limit = 7(默认) + out = protowire.AppendTag(out, 9, protowire.VarintType) + out = protowire.AppendVarint(out, 7) + out = protowire.AppendTag(out, 15, protowire.VarintType) + out = protowire.AppendVarint(out, 7) + return out, nil +} + +func buildServiceEnvelope(packet []byte, channelID string, gatewayID string) []byte { + var out []byte + out = protowire.AppendTag(out, 1, protowire.BytesType) + out = protowire.AppendBytes(out, packet) + out = protowire.AppendTag(out, 2, protowire.BytesType) + out = protowire.AppendBytes(out, []byte(channelID)) + out = protowire.AppendTag(out, 3, protowire.BytesType) + out = protowire.AppendBytes(out, []byte(gatewayID)) + return out +} diff --git a/mesh/packet.go b/mesh/packet.go new file mode 100644 index 0000000..9e80c4d --- /dev/null +++ b/mesh/packet.go @@ -0,0 +1,480 @@ +// Package mesh 实现 Meshtastic MQTT 数据包的解码与构建。 +// 协议实现参考 meshtastic_mqtt_server 工程(MIT License): +// - ServiceEnvelope / MeshPacket / Data 使用 protobuf wire 格式手写解析(google.golang.org/protobuf/encoding/protowire) +// - 加密包使用频道 PSK + AES-CTR,nonce = little-endian(packetID 8B) + little-endian(fromNum 4B) + 4 字节零 +// - channel hash = xor(channelName) ^ xor(psk) +package mesh + +import ( + "crypto/aes" + "crypto/cipher" + "encoding/base64" + "encoding/binary" + "fmt" + "strings" + "unicode/utf8" + + "google.golang.org/protobuf/encoding/protowire" +) + +// Portnum 常量(与 meshtastic.proto PortNum 一致)。 +const ( + PortNumUnknown = 0 + PortNumTextMessage = 1 + PortNumPosition = 3 + PortNumNodeInfo = 4 + PortNumRouting = 5 + PortNumTelemetry = 67 + PortNumMapReport = 73 +) + +// NodeNumBroadcast 是广播目标节点号。 +const NodeNumBroadcast uint32 = 0xffffffff + +// defaultMeshtasticPSK 是 Meshtastic 默认频道密钥(PSK 索引 1)。 +var defaultMeshtasticPSK = []byte{ + 0xD4, 0xF1, 0xBB, 0x3A, + 0x20, 0x29, 0x07, 0x59, + 0xF0, 0xBC, 0xFF, 0xAB, + 0xCF, 0x4E, 0x69, 0x01, +} + +// Packet 是解码后的 MeshPacket 关键字段。 +type Packet struct { + From uint32 + To uint32 + Channel uint32 + ID uint32 + WantAck bool + ViaMQTT bool + PKIEncrypted bool + Decoded *Data + Encrypted []byte +} + +// Data 是 MeshPacket.decoded 中的 Data 子包。 +type Data struct { + Portnum uint32 + Payload []byte +} + +// TextMessage 是一条解码后的文本消息。 +type TextMessage struct { + From uint32 + To uint32 + Text string + Hex string // 非 UTF-8 时的十六进制表示 +} + +// NodeInfo 是解码后的节点信息(NODEINFO_APP)。 +type NodeInfo struct { + From uint32 + ID string + LongName string + ShortName string + HWModel uint64 + Role uint64 + PublicKey []byte +} + +// Position 是解码后的位置信息(POSITION_APP)。 +type Position struct { + From uint32 + Latitude *float64 + Longitude *float64 + Altitude *int32 +} + +// ExpandPSK 展开 Base64 PSK,兼容 Meshtastic 默认索引 PSK 和短 key 补零规则。 +func ExpandPSK(pskBase64 string) ([]byte, error) { + psk, err := base64.StdEncoding.DecodeString(strings.TrimSpace(pskBase64)) + if err != nil { + return nil, fmt.Errorf("invalid psk: %w", err) + } + if len(psk) == 1 { + idx := psk[0] + if idx == 0 { + return []byte{}, nil + } + key := append([]byte(nil), defaultMeshtasticPSK...) + key[len(key)-1] = byte((int(key[len(key)-1]) + int(idx) - 1) & 0xff) + return key, nil + } + if len(psk) > 0 && len(psk) < 16 { + return append(psk, make([]byte, 16-len(psk))...), nil + } + if len(psk) > 16 && len(psk) < 32 { + return append(psk, make([]byte, 32-len(psk))...), nil + } + if len(psk) != 0 && len(psk) != 16 && len(psk) != 24 && len(psk) != 32 { + return nil, fmt.Errorf("invalid psk length %d: AES keys must be 16, 24, or 32 bytes", len(psk)) + } + return psk, nil +} + +// Decode 解码一个 MQTT 消息 payload(ServiceEnvelope), +// 返回解码结果(*TextMessage / *NodeInfo / *Position / *GenericPacket)与错误。 +// 加密包会用频道 PSK 尝试解密;PKI 加密包(无密钥)返回错误。 +func Decode(topic string, raw []byte, key []byte) (any, error) { + env, err := parseServiceEnvelope(raw) + if err != nil { + return nil, fmt.Errorf("service envelope: %w", err) + } + if env.Packet == nil { + return nil, fmt.Errorf("no packet in envelope") + } + pkt := env.Packet + + if pkt.Decoded == nil && len(pkt.Encrypted) > 0 { + decoded, status := tryDecrypt(pkt, env.ChannelID, key) + if decoded == nil { + return nil, fmt.Errorf("decrypt failed (%s)", status) + } + pkt.Decoded = decoded + } + + if pkt.Decoded == nil { + return nil, fmt.Errorf("empty packet") + } + + switch pkt.Decoded.Portnum { + case PortNumTextMessage: + text := string(pkt.Decoded.Payload) + msg := &TextMessage{From: pkt.From, To: pkt.To, Text: text} + if !utf8.Valid(pkt.Decoded.Payload) { + msg.Text = "" + msg.Hex = fmt.Sprintf("%x", pkt.Decoded.Payload) + } + return msg, nil + case PortNumNodeInfo: + info, err := parseUser(pkt.Decoded.Payload) + if err != nil { + return nil, err + } + info.From = pkt.From + return info, nil + case PortNumPosition: + pos, err := parsePosition(pkt.Decoded.Payload) + if err != nil { + return nil, err + } + pos.From = pkt.From + return pos, nil + case PortNumTelemetry: + return &GenericPacket{Portnum: PortNumTelemetry, From: pkt.From, PayloadLen: len(pkt.Decoded.Payload)}, nil + default: + return &GenericPacket{Portnum: pkt.Decoded.Portnum, From: pkt.From, PayloadLen: len(pkt.Decoded.Payload)}, nil + } +} + +// GenericPacket 是未详细解码的其他类型包。 +type GenericPacket struct { + Portnum uint32 + From uint32 + PayloadLen int +} + +// PortnumName 返回 portnum 的可读名称。 +func PortnumName(portnum uint32) string { + if name, ok := portNumNames[portnum]; ok { + return name + } + return fmt.Sprintf("PORTNUM_%d", portnum) +} + +// NodeNumToID 把节点号格式化为 !xxxxxxxx。 +func NodeNumToID(nodeNum uint32) string { + return fmt.Sprintf("!%08x", nodeNum) +} + +// ParseNodeID 解析 !xxxxxxxx 为节点号。 +func ParseNodeID(nodeID string) (uint32, error) { + value := strings.TrimSpace(nodeID) + value = strings.TrimPrefix(value, "!") + if len(value) != 8 { + return 0, fmt.Errorf("node id must be !xxxxxxxx") + } + var num uint32 + if _, err := fmt.Sscanf(value, "%08x", &num); err != nil { + return 0, fmt.Errorf("invalid node id: %w", err) + } + return num, nil +} + +// tryDecrypt 用频道 PSK 解密 encrypted 载荷(AES-CTR),并解析出 Data 子包。 +func tryDecrypt(pkt *Packet, channelID string, key []byte) (*Data, string) { + if len(key) == 0 { + return nil, "psk disables encryption" + } + if pkt.Channel != uint32(channelHash(channelID, key)) { + return nil, "channel hash mismatch" + } + plaintext, err := cryptAESCTR(key, pkt.From, pkt.ID, pkt.Encrypted) + if err != nil { + return nil, err.Error() + } + decoded, err := parseData(plaintext) + if err != nil { + return nil, "decrypted bytes are not Data protobuf" + } + if decoded.Portnum == PortNumUnknown { + return nil, "decrypted protobuf has UNKNOWN_APP portnum" + } + return decoded, "success" +} + +type serviceEnvelope struct { + Packet *Packet + ChannelID string + GatewayID string +} + +func parseServiceEnvelope(payload []byte) (*serviceEnvelope, error) { + env := &serviceEnvelope{} + err := walkFields(payload, func(num protowire.Number, typ protowire.Type, value any) error { + switch num { + case 1: + b, ok := value.([]byte) + if !ok || typ != protowire.BytesType { + return nil + } + packet, err := parseMeshPacket(b) + if err != nil { + return err + } + env.Packet = packet + case 2: + if b, ok := value.([]byte); ok && typ == protowire.BytesType { + env.ChannelID = string(b) + } + case 3: + if b, ok := value.([]byte); ok && typ == protowire.BytesType { + env.GatewayID = string(b) + } + } + return nil + }) + return env, err +} + +func parseMeshPacket(payload []byte) (*Packet, error) { + pkt := &Packet{} + err := walkFields(payload, func(num protowire.Number, typ protowire.Type, value any) error { + switch num { + case 1: + if v, ok := value.(uint32); ok && typ == protowire.Fixed32Type { + pkt.From = v + } + case 2: + if v, ok := value.(uint32); ok && typ == protowire.Fixed32Type { + pkt.To = v + } + case 3: + if v, ok := value.(uint64); ok && typ == protowire.VarintType { + pkt.Channel = uint32(v) + } + case 4: + if b, ok := value.([]byte); ok && typ == protowire.BytesType { + decoded, err := parseData(b) + if err != nil { + return err + } + pkt.Decoded = decoded + } + case 5: + if b, ok := value.([]byte); ok && typ == protowire.BytesType { + pkt.Encrypted = append([]byte(nil), b...) + } + case 6: + if v, ok := value.(uint32); ok && typ == protowire.Fixed32Type { + pkt.ID = v + } + case 10: + if v, ok := value.(uint64); ok && typ == protowire.VarintType { + pkt.WantAck = v != 0 + } + case 14: + if v, ok := value.(uint64); ok && typ == protowire.VarintType { + pkt.ViaMQTT = v != 0 + } + case 17: + if v, ok := value.(uint64); ok && typ == protowire.VarintType { + pkt.PKIEncrypted = v != 0 + } + } + return nil + }) + return pkt, err +} + +func parseData(payload []byte) (*Data, error) { + data := &Data{} + err := walkFields(payload, func(num protowire.Number, typ protowire.Type, value any) error { + switch num { + case 1: + if v, ok := value.(uint64); ok && typ == protowire.VarintType { + data.Portnum = uint32(v) + } + case 2: + if b, ok := value.([]byte); ok && typ == protowire.BytesType { + data.Payload = append([]byte(nil), b...) + } + } + return nil + }) + return data, err +} + +func parseUser(payload []byte) (*NodeInfo, error) { + user := &NodeInfo{} + err := walkFields(payload, func(num protowire.Number, typ protowire.Type, value any) error { + switch num { + case 1: + user.ID = stringBytes(typ, value) + case 2: + user.LongName = stringBytes(typ, value) + case 3: + user.ShortName = stringBytes(typ, value) + case 5: + user.HWModel = varintValue(typ, value) + case 7: + user.Role = varintValue(typ, value) + case 8: + if b, ok := value.([]byte); ok && typ == protowire.BytesType { + user.PublicKey = append([]byte(nil), b...) + } + } + return nil + }) + return user, err +} + +func parsePosition(payload []byte) (*Position, error) { + pos := &Position{} + err := walkFields(payload, func(num protowire.Number, typ protowire.Type, value any) error { + switch num { + case 1: + if v, ok := value.(uint32); ok && typ == protowire.Fixed32Type { + lat := float64(int32(v)) * 1e-7 + pos.Latitude = &lat + } + case 2: + if v, ok := value.(uint32); ok && typ == protowire.Fixed32Type { + lon := float64(int32(v)) * 1e-7 + pos.Longitude = &lon + } + case 3: + if typ == protowire.VarintType { + alt := int32(varintValue(typ, value)) + pos.Altitude = &alt + } + } + return nil + }) + return pos, err +} + +// walkFields 遍历 protobuf wire 字段。 +func walkFields(payload []byte, handle func(protowire.Number, protowire.Type, any) error) error { + for len(payload) > 0 { + num, typ, n := protowire.ConsumeTag(payload) + if n < 0 { + return protowire.ParseError(n) + } + payload = payload[n:] + + var value any + switch typ { + case protowire.VarintType: + v, n := protowire.ConsumeVarint(payload) + if n < 0 { + return protowire.ParseError(n) + } + value = v + payload = payload[n:] + case protowire.Fixed32Type: + v, n := protowire.ConsumeFixed32(payload) + if n < 0 { + return protowire.ParseError(n) + } + value = v + payload = payload[n:] + case protowire.Fixed64Type: + v, n := protowire.ConsumeFixed64(payload) + if n < 0 { + return protowire.ParseError(n) + } + value = v + payload = payload[n:] + case protowire.BytesType: + v, n := protowire.ConsumeBytes(payload) + if n < 0 { + return protowire.ParseError(n) + } + value = v + payload = payload[n:] + default: + n := protowire.ConsumeFieldValue(num, typ, payload) + if n < 0 { + return protowire.ParseError(n) + } + payload = payload[n:] + } + + if err := handle(num, typ, value); err != nil { + return err + } + } + return nil +} + +func stringBytes(typ protowire.Type, value any) string { + if b, ok := value.([]byte); ok && typ == protowire.BytesType { + return string(b) + } + return "" +} + +func varintValue(typ protowire.Type, value any) uint64 { + if v, ok := value.(uint64); ok && typ == protowire.VarintType { + return v + } + return 0 +} + +func xorHash(data []byte) byte { + var result byte + for _, b := range data { + result ^= b + } + return result +} + +func channelHash(channelName string, key []byte) byte { + return xorHash([]byte(channelName)) ^ xorHash(key) +} + +// cryptAESCTR 按 Meshtastic nonce 规则执行 AES-CTR;CTR 加密和解密是同一个 XOR 流操作。 +func cryptAESCTR(key []byte, fromNum, packetID uint32, input []byte) ([]byte, error) { + block, err := aes.NewCipher(key) + if err != nil { + return nil, err + } + nonce := make([]byte, aes.BlockSize) + binary.LittleEndian.PutUint64(nonce[0:8], uint64(packetID)) + binary.LittleEndian.PutUint32(nonce[8:12], fromNum) + output := make([]byte, len(input)) + cipher.NewCTR(block, nonce).XORKeyStream(output, input) + return output, nil +} + +var portNumNames = map[uint32]string{ + 0: "UNKNOWN_APP", 1: "TEXT_MESSAGE_APP", 2: "REMOTE_HARDWARE_APP", 3: "POSITION_APP", 4: "NODEINFO_APP", + 5: "ROUTING_APP", 6: "ADMIN_APP", 7: "TEXT_MESSAGE_COMPRESSED_APP", 8: "WAYPOINT_APP", 9: "AUDIO_APP", + 10: "DETECTION_SENSOR_APP", 11: "ALERT_APP", 12: "KEY_VERIFICATION_APP", 13: "REMOTE_SHELL_APP", 32: "REPLY_APP", + 33: "IP_TUNNEL_APP", 34: "PAXCOUNTER_APP", 35: "STORE_FORWARD_PLUSPLUS_APP", 36: "NODE_STATUS_APP", 64: "SERIAL_APP", + 65: "STORE_FORWARD_APP", 66: "RANGE_TEST_APP", 67: "TELEMETRY_APP", 68: "ZPS_APP", 69: "SIMULATOR_APP", + 70: "TRACEROUTE_APP", 71: "NEIGHBORINFO_APP", 72: "ATAK_PLUGIN", 73: "MAP_REPORT_APP", 74: "POWERSTRESS_APP", + 75: "LORAWAN_BRIDGE", 76: "RETICULUM_TUNNEL_APP", 77: "CAYENNE_APP", 78: "ATAK_PLUGIN_V2", 112: "GROUPALARM_APP", + 256: "PRIVATE_APP", 257: "ATAK_FORWARDER", 511: "MAX", +}