From 20757203d997d08d922a5f6f39ee90cff887cf8d Mon Sep 17 00:00:00 2001 From: hectorzhao Date: Thu, 9 Jul 2026 16:01:10 +0800 Subject: [PATCH] fix: submit cmpp2 packets for cmpp2 channels --- docs/testing-progress.md | 1 + gateway/internal/upstream/manager.go | 134 +++++++++++++----- .../internal/upstream/submit_packet_test.go | 88 ++++++++++++ 3 files changed, 190 insertions(+), 33 deletions(-) create mode 100644 gateway/internal/upstream/submit_packet_test.go diff --git a/docs/testing-progress.md b/docs/testing-progress.md index dc8ca6d..a952ddc 100644 --- a/docs/testing-progress.md +++ b/docs/testing-progress.md @@ -7,6 +7,7 @@ - 测试口径同步:`TC-ADMIN-003` 增加通道测试短信闭环要求,必须能从页面/API 发起真实测试短信,短信记录页面可查询到对应记录,Gateway submit worker 按通道真实 CMPP 配置消费发送。 - 已执行:`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 通道复制默认停用与真实连接池状态回写修复 diff --git a/gateway/internal/upstream/manager.go b/gateway/internal/upstream/manager.go index c0f900d..718577e 100644 --- a/gateway/internal/upstream/manager.go +++ b/gateway/internal/upstream/manager.go @@ -390,8 +390,10 @@ type connection struct { } type submitPartResponse struct { - rsp *cmpp.Cmpp3SubmitRspPkt - err error + seqID uint32 + msgID uint64 + result uint32 + err error } func (c *connection) matches(config queue.UpstreamConfig) bool { @@ -419,25 +421,7 @@ func (c *connection) ensureConnected() (bool, error) { func (c *connection) submitPart(ctx context.Context, cmd queue.SubmitCommand, part submitPart) (uint32, string, queue.SubmitResult, error) { rspCh := make(chan submitPartResponse, 1) - pkt := &cmpp.Cmpp3SubmitReqPkt{ - PkTotal: part.PkTotal, - PkNumber: part.PkNumber, - TpUdhi: part.TpUdhi, - RegisteredDelivery: uint8(cmd.CMPP.RegisteredDelivery), - MsgLevel: 1, - ServiceId: cmd.CMPP.ServiceID, - FeeUserType: uint8(defaultInt(cmd.CMPP.FeeUserType, 2)), - FeeTerminalId: cmd.PhoneNumber, - MsgFmt: uint8(cmd.CMPP.MsgFmt), - MsgSrc: c.config.Account, - FeeType: defaultString(cmd.CMPP.FeeType, "02"), - FeeCode: defaultString(cmd.CMPP.FeeCode, "0"), - SrcId: cmd.CMPP.SrcID, - DestUsrTl: 1, - DestTerminalId: []string{cmd.PhoneNumber}, - MsgLength: uint8(len(part.MsgContent)), - MsgContent: part.MsgContent, - } + pkt := c.submitRequestPacket(cmd, part) c.sendMu.Lock() seq, err := c.client.SendReqPkt(pkt) @@ -468,29 +452,106 @@ func (c *connection) submitPart(ctx context.Context, cmd queue.SubmitCommand, pa result := submitResult(cmd, seq, "", "timeout", "CONNECTION_LOST", rsp.err.Error()) return seq, "", result, rsp.err } - if rsp.rsp == nil { - err := fmt.Errorf("submit response is empty") - result := submitResult(cmd, seq, "", "timeout", "EMPTY_SUBMIT_RESPONSE", err.Error()) - return seq, "", result, err - } - gatewayMessageID := fmt.Sprint(rsp.rsp.MsgId) + gatewayMessageID := fmt.Sprint(rsp.msgID) status := "accepted" errorCode := "" errorMessage := "" - if rsp.rsp.Result != 0 { + if rsp.result != 0 { status = "rejected" - errorCode = fmt.Sprint(rsp.rsp.Result) - errorMessage = fmt.Sprintf("upstream submit rejected with result %d", rsp.rsp.Result) + errorCode = fmt.Sprint(rsp.result) + errorMessage = fmt.Sprintf("upstream submit rejected with result %d", rsp.result) } - if rsp.rsp.Result == 0 { + if rsp.result == 0 { c.mu.Lock() - c.tracker[rsp.rsp.MsgId] = cmd + c.tracker[rsp.msgID] = cmd c.mu.Unlock() } return seq, gatewayMessageID, submitResult(cmd, seq, gatewayMessageID, status, errorCode, errorMessage), nil } } +func (c *connection) submitRequestPacket(cmd queue.SubmitCommand, part submitPart) cmpp.Packer { + base := submitRequestFields{ + PkTotal: part.PkTotal, + PkNumber: part.PkNumber, + TpUdhi: part.TpUdhi, + RegisteredDelivery: uint8(cmd.CMPP.RegisteredDelivery), + MsgLevel: 1, + ServiceId: cmd.CMPP.ServiceID, + FeeUserType: uint8(defaultInt(cmd.CMPP.FeeUserType, 2)), + FeeTerminalId: cmd.PhoneNumber, + MsgFmt: uint8(cmd.CMPP.MsgFmt), + MsgSrc: c.config.Account, + FeeType: defaultString(cmd.CMPP.FeeType, "02"), + FeeCode: defaultString(cmd.CMPP.FeeCode, "0"), + SrcId: cmd.CMPP.SrcID, + DestUsrTl: 1, + DestTerminalId: []string{cmd.PhoneNumber}, + MsgLength: uint8(len(part.MsgContent)), + MsgContent: part.MsgContent, + } + if protocolVersion(c.config.CMPPVersion) == cmpp.V20 { + return &cmpp.Cmpp2SubmitReqPkt{ + PkTotal: base.PkTotal, + PkNumber: base.PkNumber, + RegisteredDelivery: base.RegisteredDelivery, + MsgLevel: base.MsgLevel, + ServiceId: base.ServiceId, + FeeUserType: base.FeeUserType, + FeeTerminalId: base.FeeTerminalId, + TpUdhi: base.TpUdhi, + MsgFmt: base.MsgFmt, + MsgSrc: base.MsgSrc, + FeeType: base.FeeType, + FeeCode: base.FeeCode, + SrcId: base.SrcId, + DestUsrTl: base.DestUsrTl, + DestTerminalId: base.DestTerminalId, + MsgLength: base.MsgLength, + MsgContent: base.MsgContent, + } + } + return &cmpp.Cmpp3SubmitReqPkt{ + PkTotal: base.PkTotal, + PkNumber: base.PkNumber, + RegisteredDelivery: base.RegisteredDelivery, + MsgLevel: base.MsgLevel, + ServiceId: base.ServiceId, + FeeUserType: base.FeeUserType, + FeeTerminalId: base.FeeTerminalId, + TpUdhi: base.TpUdhi, + MsgFmt: base.MsgFmt, + MsgSrc: base.MsgSrc, + FeeType: base.FeeType, + FeeCode: base.FeeCode, + SrcId: base.SrcId, + DestUsrTl: base.DestUsrTl, + DestTerminalId: base.DestTerminalId, + MsgLength: base.MsgLength, + MsgContent: base.MsgContent, + } +} + +type submitRequestFields struct { + PkTotal uint8 + PkNumber uint8 + RegisteredDelivery uint8 + MsgLevel uint8 + ServiceId string + FeeUserType uint8 + FeeTerminalId string + TpUdhi uint8 + MsgFmt uint8 + MsgSrc string + FeeType string + FeeCode string + SrcId string + DestUsrTl uint8 + DestTerminalId []string + MsgLength uint8 + MsgContent string +} + func (c *connection) tryAcquireWindow() bool { if c.window == nil { c.window = make(chan struct{}, defaultWindowSize) @@ -530,12 +591,19 @@ func (c *connection) readLoop() { continue } switch p := pkt.(type) { + case *cmpp.Cmpp2SubmitRspPkt: + c.mu.Lock() + ch := c.pending[p.SeqId] + c.mu.Unlock() + if ch != nil { + ch <- submitPartResponse{seqID: p.SeqId, msgID: p.MsgId, result: uint32(p.Result)} + } case *cmpp.Cmpp3SubmitRspPkt: c.mu.Lock() ch := c.pending[p.SeqId] c.mu.Unlock() if ch != nil { - ch <- submitPartResponse{rsp: p} + ch <- submitPartResponse{seqID: p.SeqId, msgID: p.MsgId, result: p.Result} } case *cmpp.Cmpp3DeliverReqPkt: c.handleDeliver(p) diff --git a/gateway/internal/upstream/submit_packet_test.go b/gateway/internal/upstream/submit_packet_test.go new file mode 100644 index 0000000..e06fdb4 --- /dev/null +++ b/gateway/internal/upstream/submit_packet_test.go @@ -0,0 +1,88 @@ +package upstream + +import ( + "testing" + "time" + + "cmpp-platform/gateway/internal/queue" + + cmpp "github.com/bigwhite/gocmpp" +) + +func TestSubmitRequestPacketUsesCMPP2PacketForCMPP20Channel(t *testing.T) { + conn := &connection{config: queue.UpstreamConfig{CMPPVersion: "2.0", Account: "ljcs02"}} + cmd := submitCommandForPacketTest("2.0") + parts, err := splitSubmitContent(cmd.CMPP.MsgFmt, cmd.Content) + if err != nil { + t.Fatalf("splitSubmitContent returned error: %v", err) + } + + packet, ok := conn.submitRequestPacket(cmd, parts[0]).(*cmpp.Cmpp2SubmitReqPkt) + if !ok { + t.Fatalf("expected Cmpp2SubmitReqPkt, got %T", packet) + } + if packet.MsgFmt != 8 { + t.Fatalf("MsgFmt = %d, want 8", packet.MsgFmt) + } + if packet.MsgSrc != "ljcs02" { + t.Fatalf("MsgSrc = %q, want ljcs02", packet.MsgSrc) + } + if packet.SrcId != "1069999999" { + t.Fatalf("SrcId = %q, want 1069999999", packet.SrcId) + } + if packet.MsgLength != uint8(len(parts[0].MsgContent)) { + t.Fatalf("MsgLength = %d, want %d", packet.MsgLength, len(parts[0].MsgContent)) + } + if len(packet.MsgContent)%2 != 0 { + t.Fatalf("CMPP2 UCS2 payload length must be even, got %d", len(packet.MsgContent)) + } + if packet.DestTerminalId[0] != cmd.PhoneNumber { + t.Fatalf("DestTerminalId = %q, want %q", packet.DestTerminalId[0], cmd.PhoneNumber) + } +} + +func TestSubmitRequestPacketUsesCMPP3PacketForCMPP30Channel(t *testing.T) { + conn := &connection{config: queue.UpstreamConfig{CMPPVersion: "3.0", Account: "ljcs02"}} + cmd := submitCommandForPacketTest("3.0") + parts, err := splitSubmitContent(cmd.CMPP.MsgFmt, cmd.Content) + if err != nil { + t.Fatalf("splitSubmitContent returned error: %v", err) + } + + packet, ok := conn.submitRequestPacket(cmd, parts[0]).(*cmpp.Cmpp3SubmitReqPkt) + if !ok { + t.Fatalf("expected Cmpp3SubmitReqPkt, got %T", packet) + } + if packet.MsgFmt != 8 { + t.Fatalf("MsgFmt = %d, want 8", packet.MsgFmt) + } + if packet.MsgLength != uint8(len(parts[0].MsgContent)) { + t.Fatalf("MsgLength = %d, want %d", packet.MsgLength, len(parts[0].MsgContent)) + } +} + +func submitCommandForPacketTest(cmppVersion string) queue.SubmitCommand { + return queue.SubmitCommand{ + Envelope: queue.Envelope{ + SchemaVersion: queue.SchemaVersion, + MessageType: queue.MessageTypeSubmitCommand, + TraceID: "trace-1", + MessageID: "message-1", + ChannelID: "channel-1", + CreatedAt: time.Now(), + }, + SubmitID: "submit-1", + PhoneNumber: "18821203795", + Content: "【安徽航天信息】您的验证码是070926,有效时间30分钟。", + CMPP: queue.CMPP{ + ServiceID: "SMS", + SrcID: "1069999999", + RegisteredDelivery: 1, + MsgFmt: 8, + }, + Upstream: queue.UpstreamConfig{ + Account: "ljcs02", + CMPPVersion: cmppVersion, + }, + } +}