From dc7d015219c591e11f0c62bb583f2800fd7047e5 Mon Sep 17 00:00:00 2001 From: hectorzhao Date: Thu, 9 Jul 2026 16:14:11 +0800 Subject: [PATCH] fix: handle cmpp2 deliver receipts --- docs/testing-progress.md | 1 + gateway/internal/upstream/deliver_test.go | 88 +++++++++++++++++++++++ gateway/internal/upstream/manager.go | 75 ++++++++++++++----- 3 files changed, 146 insertions(+), 18 deletions(-) create mode 100644 gateway/internal/upstream/deliver_test.go diff --git a/docs/testing-progress.md b/docs/testing-progress.md index a952ddc..44f614b 100644 --- a/docs/testing-progress.md +++ b/docs/testing-progress.md @@ -8,6 +8,7 @@ - 已执行:`npm --prefix api test -- channels.service.spec.ts --runInBand`、`npm --prefix api run build`、`npm run build`。待生产部署后用指定号码做一次真实发送验证,并回查 DB/短信记录。 - 2026-07-09 追加:按产品边界收窄通道测试短信,`SmsMessageRecord/SmsSubmitRecord/SmsReceiptRecord/SmsMessageSegmentAudit` 支持 `tenantId/batchTaskId` 为空;通道测试只记录短信、提交结果和回执,不进入企业账务、发送任务进度、客户侧 Deliver 推送或业务补发。 - 2026-07-09 追加:生产验证短信进入 Gateway 后出现 `SUBMIT_TIMEOUT/context deadline exceeded`,上游平台疑似内容乱码;排查确认 API、数据库和 Redis Stream 中中文内容正常,`SubmitCommand.cmpp.msgFmt=8`,根因为 Gateway 在 CMPP 2.0 通道上仍固定构造 `Cmpp3SubmitReqPkt`。已修复为按通道 `cmppVersion` 分别构造 `Cmpp2SubmitReqPkt`/`Cmpp3SubmitReqPkt`,并同时处理 CMPP 2.0/3.0 SubmitResp。新增 `submit_packet_test.go` 覆盖 CMPP 2.0 中文 UCS2 Submit 包类型和长度。 +- 2026-07-09 追加:生产验证 Submit 已 accepted 后未见 `SmsReceiptRecord`,结合上游为 CMPP 2.0 排查 Gateway readLoop,确认只处理 `Cmpp3DeliverReqPkt`,CMPP 2.0 回执 Deliver 即使到达连接也不会 ACK 或进入 receipt 事件链路。已修复为同时处理 `Cmpp2DeliverReqPkt`/`Cmpp3DeliverReqPkt`,分别返回 `Cmpp2DeliverRspPkt`/`Cmpp3DeliverRspPkt`,并统一 Deliver 解码、回执和上行处理。新增 `deliver_test.go` 覆盖 CMPP 2.0 `DELIVRD` 回执写入事件链路。 ## 2026-07-09 通道复制默认停用与真实连接池状态回写修复 diff --git a/gateway/internal/upstream/deliver_test.go b/gateway/internal/upstream/deliver_test.go new file mode 100644 index 0000000..b13d44f --- /dev/null +++ b/gateway/internal/upstream/deliver_test.go @@ -0,0 +1,88 @@ +package upstream + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "testing" + "time" + + "cmpp-platform/gateway/internal/queue" + + cmpp "github.com/bigwhite/gocmpp" +) + +func TestHandleCMPP2DeliverReceiptPostsReceiptEvent(t *testing.T) { + events := make(chan queue.ReceiptEvent, 1) + api := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/gateway/events/receipt" { + t.Fatalf("unexpected path: %s", r.URL.Path) + } + var event queue.ReceiptEvent + if err := json.NewDecoder(r.Body).Decode(&event); err != nil { + t.Fatalf("decode receipt event: %v", err) + } + events <- event + w.WriteHeader(http.StatusOK) + })) + defer api.Close() + + receipt := cmpp.CmppReceiptPkt{ + MsgId: 8412634832294102675, + Stat: "DELIVRD", + SubmitTime: "2607091558", + DoneTime: "2607091559", + DestTerminalId: "18821203795", + SmscSequence: 1, + } + payload, err := receipt.Pack() + if err != nil { + t.Fatalf("pack receipt: %v", err) + } + + conn := &connection{ + channelID: "channel-1", + apiBaseURL: api.URL, + httpClient: api.Client(), + tracker: map[uint64]queue.SubmitCommand{ + receipt.MsgId: { + Envelope: queue.Envelope{ + SchemaVersion: queue.SchemaVersion, + MessageType: queue.MessageTypeSubmitCommand, + TraceID: "trace-1", + MessageID: "message-1", + ChannelID: "channel-1", + CreatedAt: time.Now(), + }, + SubmitID: "submit-1", + }, + }, + } + conn.handleDeliver(deliverPacketFromCMPP2(&cmpp.Cmpp2DeliverReqPkt{ + SeqId: 7, + MsgId: 999, + RegisterDelivery: 1, + MsgContent: string(payload), + })) + + select { + case event := <-events: + if event.MessageID != "message-1" || event.ChannelID != "channel-1" || event.TraceID != "trace-1" { + t.Fatalf("unexpected envelope: %+v", event.Envelope) + } + if event.SequenceID != 7 { + t.Fatalf("SequenceID = %d, want 7", event.SequenceID) + } + if event.GatewayMessageID != "8412634832294102675" { + t.Fatalf("GatewayMessageID = %q", event.GatewayMessageID) + } + if event.PhoneNumber != "18821203795" { + t.Fatalf("PhoneNumber = %q", event.PhoneNumber) + } + if event.ReceiptStatus != "delivered" || event.RawStatus != "DELIVRD" { + t.Fatalf("unexpected receipt status: %+v", event) + } + case <-time.After(time.Second): + t.Fatal("timed out waiting for receipt event") + } +} diff --git a/gateway/internal/upstream/manager.go b/gateway/internal/upstream/manager.go index 718577e..3c46ac6 100644 --- a/gateway/internal/upstream/manager.go +++ b/gateway/internal/upstream/manager.go @@ -605,25 +605,64 @@ func (c *connection) readLoop() { if ch != nil { ch <- submitPartResponse{seqID: p.SeqId, msgID: p.MsgId, result: p.Result} } + case *cmpp.Cmpp2DeliverReqPkt: + _ = c.client.SendRspPkt(&cmpp.Cmpp2DeliverRspPkt{MsgId: p.MsgId, Result: 0}, p.SeqId) + c.handleDeliver(deliverPacketFromCMPP2(p)) case *cmpp.Cmpp3DeliverReqPkt: - c.handleDeliver(p) + _ = c.client.SendRspPkt(&cmpp.Cmpp3DeliverRspPkt{MsgId: p.MsgId, Result: 0}, p.SeqId) + c.handleDeliver(deliverPacketFromCMPP3(p)) case *cmpp.CmppActiveTestReqPkt: _ = c.client.SendRspPkt(&cmpp.CmppActiveTestRspPkt{}, p.SeqId) } } } -func (c *connection) handleDeliver(pkt *cmpp.Cmpp3DeliverReqPkt) { - _ = c.client.SendRspPkt(&cmpp.Cmpp3DeliverRspPkt{MsgId: pkt.MsgId, Result: 0}, pkt.SeqId) +type deliverPacket struct { + seqID uint32 + msgID uint64 + destID string + tpUdhi uint8 + msgFmt uint8 + srcTerminalID string + registerDelivery uint8 + msgContent string +} - if pkt.RegisterDelivery == 1 { +func deliverPacketFromCMPP2(pkt *cmpp.Cmpp2DeliverReqPkt) deliverPacket { + return deliverPacket{ + seqID: pkt.SeqId, + msgID: pkt.MsgId, + destID: pkt.DestId, + tpUdhi: pkt.TpUdhi, + msgFmt: pkt.MsgFmt, + srcTerminalID: pkt.SrcTerminalId, + registerDelivery: pkt.RegisterDelivery, + msgContent: pkt.MsgContent, + } +} + +func deliverPacketFromCMPP3(pkt *cmpp.Cmpp3DeliverReqPkt) deliverPacket { + return deliverPacket{ + seqID: pkt.SeqId, + msgID: pkt.MsgId, + destID: pkt.DestId, + tpUdhi: pkt.TpUdhi, + msgFmt: pkt.MsgFmt, + srcTerminalID: pkt.SrcTerminalId, + registerDelivery: pkt.RegisterDelivery, + msgContent: pkt.MsgContent, + } +} + +func (c *connection) handleDeliver(pkt deliverPacket) { + if pkt.registerDelivery == 1 { var receipt cmpp.CmppReceiptPkt - if err := receipt.Unpack([]byte(pkt.MsgContent)); err != nil { + if err := receipt.Unpack([]byte(pkt.msgContent)); err != nil { return } cmd, ok := c.commandFor(receipt.MsgId) if !ok { - cmd, ok = c.commandFor(pkt.MsgId) + cmd, ok = c.commandFor(pkt.msgID) } traceID := fmt.Sprintf("receipt-%d", receipt.MsgId) messageID := fmt.Sprintf("receipt-%d", receipt.MsgId) @@ -642,7 +681,7 @@ func (c *connection) handleDeliver(pkt *cmpp.Cmpp3DeliverReqPkt) { ChannelID: channelID, CreatedAt: time.Now().UTC(), }, - SequenceID: pkt.SeqId, + SequenceID: pkt.seqID, GatewayMessageID: fmt.Sprint(receipt.MsgId), PhoneNumber: strings.TrimSpace(receipt.DestTerminalId), ReceiptStatus: receiptStatus(receipt.Stat), @@ -660,7 +699,7 @@ func (c *connection) handleDeliver(pkt *cmpp.Cmpp3DeliverReqPkt) { if !complete { return } - cmd, _ := c.commandFor(pkt.MsgId) + cmd, _ := c.commandFor(pkt.msgID) event := queue.UplinkEvent{ Envelope: queue.Envelope{ SchemaVersion: queue.SchemaVersion, @@ -670,32 +709,32 @@ func (c *connection) handleDeliver(pkt *cmpp.Cmpp3DeliverReqPkt) { ChannelID: c.channelID, CreatedAt: time.Now().UTC(), }, - SequenceID: pkt.SeqId, - PhoneNumber: strings.TrimSpace(pkt.SrcTerminalId), - DestID: strings.TrimSpace(pkt.DestId), + SequenceID: pkt.seqID, + PhoneNumber: strings.TrimSpace(pkt.srcTerminalID), + DestID: strings.TrimSpace(pkt.destID), Content: content, ReceivedAt: time.Now().UTC(), } _ = postJSON(context.Background(), c.httpClient, c.apiBaseURL, "/gateway/events/uplink", event) } -func (c *connection) decodeUplinkContent(pkt *cmpp.Cmpp3DeliverReqPkt) (string, bool, error) { - if pkt.TpUdhi != 1 { - content, err := decodeContent(pkt.MsgFmt, pkt.MsgContent) +func (c *connection) decodeUplinkContent(pkt deliverPacket) (string, bool, error) { + if pkt.tpUdhi != 1 { + content, err := decodeContent(pkt.msgFmt, pkt.msgContent) return content, true, err } - ref, total, number, payload, ok := parseConcatSegment(pkt.MsgContent) + ref, total, number, payload, ok := parseConcatSegment(pkt.msgContent) if !ok { - content, err := decodeContent(pkt.MsgFmt, pkt.MsgContent) + content, err := decodeContent(pkt.msgFmt, pkt.msgContent) return content, true, err } - key := fmt.Sprintf("%s:%s:%s:%d:%d", c.channelID, strings.TrimSpace(pkt.SrcTerminalId), strings.TrimSpace(pkt.DestId), ref, total) + key := fmt.Sprintf("%s:%s:%s:%d:%d", c.channelID, strings.TrimSpace(pkt.srcTerminalID), strings.TrimSpace(pkt.destID), ref, total) c.mu.Lock() if c.longUplink == nil { c.longUplink = make(map[string]*longUplinkAssembly) } pruneLongUplinkAssemblies(c.longUplink, time.Now(), 10*time.Minute) - content, complete, err := assembleLongUplink(c.longUplink, key, pkt.MsgFmt, total, number, payload) + content, complete, err := assembleLongUplink(c.longUplink, key, pkt.msgFmt, total, number, payload) c.mu.Unlock() return content, complete, err }