fix: submit cmpp2 packets for cmpp2 channels

This commit is contained in:
hectorzhao
2026-07-09 16:01:10 +08:00
parent 64bb3fb430
commit 20757203d9
3 changed files with 190 additions and 33 deletions
+1
View File
@@ -7,6 +7,7 @@
- 测试口径同步:`TC-ADMIN-003` 增加通道测试短信闭环要求,必须能从页面/API 发起真实测试短信,短信记录页面可查询到对应记录,Gateway submit worker 按通道真实 CMPP 配置消费发送。 - 测试口径同步:`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/短信记录。 - 已执行:`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 追加:按产品边界收窄通道测试短信,`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 通道复制默认停用与真实连接池状态回写修复 ## 2026-07-09 通道复制默认停用与真实连接池状态回写修复
+100 -32
View File
@@ -390,7 +390,9 @@ type connection struct {
} }
type submitPartResponse struct { type submitPartResponse struct {
rsp *cmpp.Cmpp3SubmitRspPkt seqID uint32
msgID uint64
result uint32
err error err error
} }
@@ -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) { func (c *connection) submitPart(ctx context.Context, cmd queue.SubmitCommand, part submitPart) (uint32, string, queue.SubmitResult, error) {
rspCh := make(chan submitPartResponse, 1) rspCh := make(chan submitPartResponse, 1)
pkt := &cmpp.Cmpp3SubmitReqPkt{ pkt := c.submitRequestPacket(cmd, part)
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,
}
c.sendMu.Lock() c.sendMu.Lock()
seq, err := c.client.SendReqPkt(pkt) 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()) result := submitResult(cmd, seq, "", "timeout", "CONNECTION_LOST", rsp.err.Error())
return seq, "", result, rsp.err return seq, "", result, rsp.err
} }
if rsp.rsp == nil { gatewayMessageID := fmt.Sprint(rsp.msgID)
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)
status := "accepted" status := "accepted"
errorCode := "" errorCode := ""
errorMessage := "" errorMessage := ""
if rsp.rsp.Result != 0 { if rsp.result != 0 {
status = "rejected" status = "rejected"
errorCode = fmt.Sprint(rsp.rsp.Result) errorCode = fmt.Sprint(rsp.result)
errorMessage = fmt.Sprintf("upstream submit rejected with result %d", rsp.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.mu.Lock()
c.tracker[rsp.rsp.MsgId] = cmd c.tracker[rsp.msgID] = cmd
c.mu.Unlock() c.mu.Unlock()
} }
return seq, gatewayMessageID, submitResult(cmd, seq, gatewayMessageID, status, errorCode, errorMessage), nil 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 { func (c *connection) tryAcquireWindow() bool {
if c.window == nil { if c.window == nil {
c.window = make(chan struct{}, defaultWindowSize) c.window = make(chan struct{}, defaultWindowSize)
@@ -530,12 +591,19 @@ func (c *connection) readLoop() {
continue continue
} }
switch p := pkt.(type) { 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: case *cmpp.Cmpp3SubmitRspPkt:
c.mu.Lock() c.mu.Lock()
ch := c.pending[p.SeqId] ch := c.pending[p.SeqId]
c.mu.Unlock() c.mu.Unlock()
if ch != nil { if ch != nil {
ch <- submitPartResponse{rsp: p} ch <- submitPartResponse{seqID: p.SeqId, msgID: p.MsgId, result: p.Result}
} }
case *cmpp.Cmpp3DeliverReqPkt: case *cmpp.Cmpp3DeliverReqPkt:
c.handleDeliver(p) c.handleDeliver(p)
@@ -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,
},
}
}