Files

328 lines
7.5 KiB
Go

// 用 MQTT 向真实 nixmsg 提交一条未来定时消息,供后台「投递记录」页看到未完成记录。
// 管理接口只读 messages,不能造记录;不改服务器业务代码。
package main
import (
"bytes"
"encoding/json"
"flag"
"fmt"
"io"
"os"
"sync"
"time"
"git.asio.asia/nixevol/NixMsg/internal/protocol"
"git.asio.asia/nixevol/NixMsg/test/harness"
"github.com/mochi-mqtt/server/v2/packets"
)
func main() {
base := flag.String("base", "", "http base of nixmsg, e.g. http://127.0.0.1:12345")
fromID := flag.String("from", "", "sender endpoint id")
fromPass := flag.String("from-pass", "", "sender login password")
toID := flag.String("to", "", "receiver endpoint id")
flag.Parse()
if *base == "" || *fromID == "" || *fromPass == "" || *toID == "" {
fmt.Fprintln(os.Stderr, "usage: seedpending -base URL -from ID -from-pass PWD -to ID")
os.Exit(2)
}
if err := seed(*base, *fromID, *fromPass, *toID); err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
fmt.Println("ok")
}
func seed(httpBase, fromID, fromPass, toID string) error {
mc, err := harness.DialMQTTWebSocket(httpBase, 10*time.Second)
if err != nil {
return fmt.Errorf("dial mqtt: %w", err)
}
defer func() { _ = mc.Close() }()
s := &session{mc: mc, endpointID: fromID, pktID: 10, done: make(chan struct{})}
if helloErr := s.connectSubscribeHello(fromPass); helloErr != nil {
return helloErr
}
go s.readLoop()
defer s.close()
sendAt := time.Now().Add(2 * time.Hour).UnixMilli()
frame := map[string]any{
"v": 1, "type": "send", "rid": "w4seed1", "id": "w4-e2e-scheduled",
"to": map[string]any{"kind": "endpoint", "id": toID},
"body": map[string]any{"enc": "utf8", "data": "w4-e2e-pending"},
"send_at_ms": sendAt,
}
resp, reqErr := s.request(frame, 15*time.Second)
if reqErr != nil {
return reqErr
}
if !resp.ok {
return fmt.Errorf("send failed: %+v", resp.raw)
}
return nil
}
type session struct {
mc harness.MQTTClient
endpointID string
pktID uint16
mu sync.Mutex
inbox []map[string]any
closed bool
done chan struct{}
}
type appResp struct {
ok bool
raw map[string]any
}
func (s *session) close() {
s.mu.Lock()
if s.closed {
s.mu.Unlock()
return
}
s.closed = true
s.mu.Unlock()
_ = s.mc.Close()
select {
case <-s.done:
case <-time.After(3 * time.Second):
}
}
func (s *session) nextPkt() uint16 {
s.pktID++
if s.pktID == 0 {
s.pktID = 1
}
return s.pktID
}
func (s *session) connectSubscribeHello(password string) error {
pk := packets.Packet{
FixedHeader: packets.FixedHeader{Type: packets.Connect},
ProtocolVersion: 5,
Connect: packets.ConnectParams{
ProtocolName: []byte("MQTT"),
Clean: true,
ClientIdentifier: s.endpointID,
Keepalive: 30,
UsernameFlag: true,
Username: []byte(s.endpointID),
PasswordFlag: true,
Password: []byte(password),
},
}
var buf bytes.Buffer
if err := pk.ConnectEncode(&buf); err != nil {
return err
}
if err := s.mc.Send(buf.Bytes()); err != nil {
return err
}
ack, err := s.mc.Recv()
if err != nil {
return err
}
if len(ack) < 4 || ack[0]>>4 != packets.Connack || ack[3] != 0 {
return fmt.Errorf("connack %x", ack)
}
sub := packets.Packet{
FixedHeader: packets.FixedHeader{Type: packets.Subscribe, Qos: 1},
ProtocolVersion: 5,
PacketID: s.nextPkt(),
Filters: packets.Subscriptions{
{Filter: "nix/c/" + s.endpointID + "/down", Qos: 1},
},
}
buf.Reset()
if err := sub.SubscribeEncode(&buf); err != nil {
return err
}
if err := s.mc.Send(buf.Bytes()); err != nil {
return err
}
if _, err := s.mc.Recv(); err != nil {
return err
}
hello, _ := protocol.Marshal(protocol.Hello{V: protocol.Version, Type: protocol.TypeHello, RID: "h0"})
if err := s.publishRaw(hello); err != nil {
return err
}
deadline := time.Now().Add(10 * time.Second)
for time.Now().Before(deadline) {
raw, err := s.mc.Recv()
if err != nil {
return err
}
m := s.handlePacket(raw)
if m == nil {
continue
}
if m["type"] == "resp" && m["ok"] == true {
return nil
}
if m["type"] == "resp" {
return fmt.Errorf("hello failed: %v", m)
}
s.push(m)
}
return fmt.Errorf("hello timeout")
}
func (s *session) publishRaw(payload []byte) error {
pub := packets.Packet{
FixedHeader: packets.FixedHeader{Type: packets.Publish, Qos: 1},
ProtocolVersion: 5,
TopicName: "nix/c/" + s.endpointID + "/up",
PacketID: s.nextPkt(),
Payload: payload,
}
var buf bytes.Buffer
if err := pub.PublishEncode(&buf); err != nil {
return err
}
return s.mc.Send(buf.Bytes())
}
func (s *session) readLoop() {
defer close(s.done)
for {
raw, err := s.mc.Recv()
if err != nil {
return
}
m := s.handlePacket(raw)
if m != nil {
s.push(m)
}
}
}
func (s *session) handlePacket(raw []byte) map[string]any {
if len(raw) < 2 {
return nil
}
typ := raw[0] >> 4
qos := (raw[0] >> 1) & 0x3
switch typ {
case packets.Puback, packets.Pingresp, packets.Suback:
return nil
case packets.Publish:
payload, err := decodePublishPayload(raw)
if err != nil {
return nil
}
if qos == 1 {
rem, n, _ := decodeRemainingLength(raw[1:])
body := raw[1+n:]
pk := packets.Packet{ProtocolVersion: 5, FixedHeader: packets.FixedHeader{Type: packets.Publish, Remaining: rem, Qos: qos}}
if decErr := pk.PublishDecode(body); decErr == nil && pk.PacketID != 0 {
ack := packets.Packet{
FixedHeader: packets.FixedHeader{Type: packets.Puback},
ProtocolVersion: 5,
PacketID: pk.PacketID,
}
var buf bytes.Buffer
if encErr := ack.PubackEncode(&buf); encErr == nil {
_ = s.mc.Send(buf.Bytes())
}
}
}
var m map[string]any
if json.Unmarshal(payload, &m) != nil {
return nil
}
return m
default:
return nil
}
}
func (s *session) push(m map[string]any) {
s.mu.Lock()
s.inbox = append(s.inbox, m)
s.mu.Unlock()
}
func (s *session) request(frame map[string]any, timeout time.Duration) (appResp, error) {
rid, _ := frame["rid"].(string)
payload, err := protocol.Marshal(frame)
if err != nil {
return appResp{}, err
}
if err := s.publishRaw(payload); err != nil {
return appResp{}, err
}
deadline := time.Now().Add(timeout)
for time.Now().Before(deadline) {
m := s.takeMatching(func(x map[string]any) bool {
return x["type"] == "resp" && x["rid"] == rid
})
if m != nil {
return appResp{ok: m["ok"] == true, raw: m}, nil
}
time.Sleep(5 * time.Millisecond)
}
return appResp{}, fmt.Errorf("timeout waiting resp rid=%s", rid)
}
func (s *session) takeMatching(pred func(map[string]any) bool) map[string]any {
s.mu.Lock()
defer s.mu.Unlock()
for i, m := range s.inbox {
if pred(m) {
s.inbox = append(s.inbox[:i], s.inbox[i+1:]...)
return m
}
}
return nil
}
func decodePublishPayload(raw []byte) ([]byte, error) {
if len(raw) < 2 {
return nil, io.ErrUnexpectedEOF
}
rem, n, err := decodeRemainingLength(raw[1:])
if err != nil {
return nil, err
}
body := raw[1+n:]
if len(body) != rem {
return nil, io.ErrUnexpectedEOF
}
pk := packets.Packet{
ProtocolVersion: 5,
FixedHeader: packets.FixedHeader{
Type: packets.Publish,
Remaining: rem,
Qos: (raw[0] >> 1) & 0x3,
},
}
if err := pk.PublishDecode(body); err != nil {
return nil, err
}
return pk.Payload, nil
}
func decodeRemainingLength(b []byte) (value int, n int, err error) {
var mul uint32 = 1
var v uint32
for i := 0; i < len(b) && i < 4; i++ {
v += uint32(b[i]&127) * mul
n++
if b[i]&128 == 0 {
return int(v), n, nil
}
mul *= 128
}
return 0, 0, io.ErrUnexpectedEOF
}