From 010ba3216889032a6160cdb14d8536b616ae7102 Mon Sep 17 00:00:00 2001 From: hectorzhao Date: Wed, 16 Sep 2026 19:08:42 +0800 Subject: [PATCH] fix: register supplier response waiter before reader can consume reply --- docs/phase-4-send-pipeline-redesign.md | 6 + docs/system-functional-test-cases.md | 4 + docs/testing-progress.md | 8 ++ .../upstream/immediate_response_test.go | 103 ++++++++++++++++++ gateway/internal/upstream/submit.go | 11 +- 5 files changed, 128 insertions(+), 4 deletions(-) create mode 100644 gateway/internal/upstream/immediate_response_test.go diff --git a/docs/phase-4-send-pipeline-redesign.md b/docs/phase-4-send-pipeline-redesign.md index 7c5598e..2943d3b 100644 --- a/docs/phase-4-send-pipeline-redesign.md +++ b/docs/phase-4-send-pipeline-redesign.md @@ -346,3 +346,9 @@ GatewaySubmitOutbox - 收尾needs_review及等待超过300秒由监控采集器直接查询耐久工作/事件表,写入同一告警事件库,恢复仍保留、仅人工清除;不依赖应用发布工具安装Prometheus规则。过程计数仍提供低基数metrics,采集进程口径需区分。 - 两项迁移仅新增四张表及索引,无历史回填、无删除、无自动重发;全量106项已在新隔离数据库重放。工作时间使用UTC表达式,避免数据库会话时区与Prisma时间不一致。应用回退前必须盘点和接续未完成工作,旧代码不能消费这些新表。 - 本地验证包括真实PostgreSQL、两个独立OS进程、事务回滚、旧消费者挂起后接管、三段通知目标、非零账务和重试上限可见告警;网络投递与目标环境的重启/排空验收独立记录,不将路由隔离测试冒称整链路通过。 + +### 10.12 测试环境发现的快速SubmitResp竞态(2026-09-16) + +a350aca测试环境长短信验收发现两条消息首尝试分别仅写出2/4、1/3段,Gateway在已收到即时应答时仍等待60秒后误判SUBMIT_TIMEOUT;补发及最终账务/通知收尾正常,但不能据此视作无异常验收。只读代码证据:submitPart在SendReqPkt、异步日志启动之后才登记pending[seq],readLoop可能提前消费应答并因不存在等待者丢弃。新增真实TCP回归在修改前分别于CMPP2.0第206次、3.0第7次复现。 + +最小修复保持现有协议、接口、存储和补发策略:使用与heartbeat一致的mu→sendMu锁序,将连接有效性检查、写包和登记pending置于同一临界区,响应读取须等登记完成。网络失败仍走原连接关闭/失败处理。锁内不得执行日志、业务回调或数据库操作。验证两种协议各500次即时应答、Gateway全量test/vet及测试环境新的长短信样本;原异常证据保留,不将旧样本改成无补发成功。 diff --git a/docs/system-functional-test-cases.md b/docs/system-functional-test-cases.md index 08b337c..47d3279 100644 --- a/docs/system-functional-test-cases.md +++ b/docs/system-functional-test-cases.md @@ -5606,3 +5606,7 @@ CLIENT-0914-01~07 的模板样式/顺序、文档归属与检索、中文状 | TC-OPS-0916-05 | 收尾第12次失败进入needs_review,自动出现在耐久告警;恢复后不自动消失 | 真实PG通过 | RC-01/02/04/08/09/10/11目前仅部分本地证据:三段齐段、重复与双进程、事务回滚和过期接管、通知目标数与非零账务、迁移重放;尚不能标整条矩阵通过。RC-03/05/06/07/12的完整组合、网络故障/进程重启、修复前后性能对比须独立补验。无真实运营商流量,未声明容量提升。 + +### TC-RC-20260916-13 Gateway即时SubmitResp + +真实TCP通道在收到每个Submit后立即回包,两种协议各500次;全部一次受理,不得丢应答后误报SUBMIT_TIMEOUT。目标环境补测多段消息,核对每段wire提交一次、唯一账单/最终业务通知;保留修复前2/4、1/3段异常与后继记录。本地修改前已复现,修复后重复5轮通过;目标环境补充发布待执行。 diff --git a/docs/testing-progress.md b/docs/testing-progress.md index 67b240c..13908e7 100644 --- a/docs/testing-progress.md +++ b/docs/testing-progress.md @@ -5091,3 +5091,11 @@ CUA本轮可用,实际后端文档三尺寸1600×1000/1366×768/390×844无页 前端全量35套163项通过(115.33秒),后端77套835项通过(171.189秒);类型检查通过。新增隔离PG全量106迁移重放和11组耐久工作/告警验证通过,包含独立OS进程并发及挂起旧消费者恢复,日志在.local-data/five-fixes-20260916。第12次失败转人工处理及数据库告警采集已补验证;日志中的故意故障注入不得当线上异常。最新采集改动需候选精确提交回归。 真实环境尚未部署;独立Playwright获用户允许,专用Edge等待用户完成登录。测试Outbox发布开关及独立进程已核验。为通过更名涉及文件的既有门禁,仅移除official-export未使用导入并格式化该文件及RechargeReceiptDialog,业务规则不变。保护项未纳入;开始准备精确暂存,尚未提交/推送。 + +### 2026-09-16 19:10 测试环境验收发现并修复Gateway即时应答竞态 + +五项应用实现a350aca及更名a0209f9已提交/推送,测试环境按标准工具部署a350aca(计划20260916T104026-a350aca88336-7f5f9fd9)。候选前端163项/API834项及类型/格式/lint/样式门禁通过;与工作区API835项相差1项来自受保护的既有metrics测试,未夹带。106迁移完成,预生产未操作。 + +隔离真实HTTP/CMPP/Gateway/PostgreSQL/Redis验收已观察2/3/4段、乱序及重复回执、失败后唯一后继、最终失败退款、断线CMPP恢复9份客户分段回执、不申请回执不创建CMPP投递、Webhook TLS/503/超时后自动恢复。HTTP页面真实202及报文展示通过;无签名且含未报备链接的3段通道测试真实受理。监控实际恢复告警留存,故障注入保留最后真实快照并恢复刷新。完整对账及三尺寸页面/人工清除收尾仍在进行,不宣称全部通过。 + +对账发现两条测试消息首尝试只提交2/4、1/3段后误判SUBMIT_TIMEOUT,各自动产生一个后继。已只读确认Gateway发送后登记pending的窗口;真实TCP在未修代码CMPP2.0第206次、3.0第7次复现,详见方案10.12。按同一mu→sendMu锁序保护写包和等待者登记;修复后两协议各500次×5轮共5000次即时应答通过,Gateway全量go test ./...与go vet ./...通过。需重新提交/标准部署该补充修复并新增目标环境样本验证;原异常保留,不将异常重试样本计为无重复发送通过。 diff --git a/gateway/internal/upstream/immediate_response_test.go b/gateway/internal/upstream/immediate_response_test.go new file mode 100644 index 0000000..8c4b047 --- /dev/null +++ b/gateway/internal/upstream/immediate_response_test.go @@ -0,0 +1,103 @@ +package upstream + +import ( + "context" + "encoding/binary" + "io" + "net" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + "time" + + "cmpp-platform/gateway/internal/queue" + cmpp "github.com/bigwhite/gocmpp" +) + +func TestImmediateSupplierResponseIsNeverLost(t *testing.T) { + for _, version := range []string{"2.0", "3.0"} { + t.Run(version, func(t *testing.T) { + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + defer listener.Close() + var wire atomic.Int32 + done := make(chan struct{}) + go func() { + defer close(done) + socket, err := listener.Accept() + if err != nil { + return + } + defer socket.Close() + for { + header := make([]byte, 12) + if _, err = io.ReadFull(socket, header); err != nil { + return + } + n := binary.BigEndian.Uint32(header) + if n < 12 || n > 4096 { + return + } + body := make([]byte, int(n)-12) + if _, err = io.ReadFull(socket, body); err != nil { + return + } + kind := binary.BigEndian.Uint32(header[4:]) + var response []byte + switch kind { + case 1: + size := 18 + if version == "3.0" { + size = 21 + } + response = make([]byte, size) + response[size-1] = byte(protocolVersion(version)) + case 4: + wire.Add(1) + size := 9 + if version == "3.0" { + size = 12 + } + response = make([]byte, size) + binary.BigEndian.PutUint64(response, uint64(wire.Load())) + default: + return + } + binary.BigEndian.PutUint32(header, uint32(12+len(response))) + binary.BigEndian.PutUint32(header[4:], kind|0x80000000) + if _, err = socket.Write(append(header, response...)); err != nil { + return + } + } + }() + api := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(204) })) + defer api.Close() + client := cmpp.NewClient(protocolVersion(version)) + if err = client.Connect(listener.Addr().String(), "qa", "qa", time.Second); err != nil { + t.Fatal(err) + } + c := &connection{client: client, retired: true, config: queue.UpstreamConfig{CMPPVersion: version}, pending: make(map[uint32]chan submitPartResponse), tracker: make(map[uint64]queue.SubmitCommand), apiBaseURL: api.URL, httpClient: api.Client()} + defer func() { c.close(); <-done }() + go c.readLoop() + cmd := submitCommandForPacketTest(version) + parts, err := splitSubmitContent(cmd.CMPP.MsgFmt, cmd.Content) + if err != nil { + t.Fatal(err) + } + for i := 0; i < 500; i++ { + ctx, cancel := context.WithTimeout(context.Background(), 300*time.Millisecond) + _, _, result, err := c.submitPart(ctx, cmd, parts[0]) + cancel() + if err != nil || result.SubmitStatus != "accepted" { + t.Fatalf("response lost at wire submit %d: %+v / %v", i+1, result, err) + } + } + if wire.Load() != 500 { + t.Fatalf("wire submits = %d", wire.Load()) + } + }) + } +} diff --git a/gateway/internal/upstream/submit.go b/gateway/internal/upstream/submit.go index b7667a4..6e1be8d 100644 --- a/gateway/internal/upstream/submit.go +++ b/gateway/internal/upstream/submit.go @@ -138,8 +138,8 @@ func (c *connection) submitPart(ctx context.Context, cmd queue.SubmitCommand, pa c.mu.Lock() client := c.client closed := c.closed - c.mu.Unlock() if closed || client == nil { + c.mu.Unlock() err := fmt.Errorf("supplier connection is not available") result := submitResult(cmd, 0, "", "timeout", "CONNECTION_LOST", err.Error()) return 0, "", result, err @@ -153,6 +153,12 @@ func (c *connection) submitPart(ctx context.Context, cmd queue.SubmitCommand, pa wireSource = "gateway_write_complete" } c.sendMu.Unlock() + // The reader must not consume an immediate response before its waiter is + // registered. Keep the same mu -> sendMu lock order as the heartbeat path. + if err == nil { + c.pending[seq] = rspCh + } + c.mu.Unlock() if err != nil { c.emitProtocolLog(protocolLogEvent{ Protocol: "cmpp", @@ -195,9 +201,6 @@ func (c *connection) submitPart(ctx context.Context, cmd queue.SubmitCommand, pa }, }) - c.mu.Lock() - c.pending[seq] = rspCh - c.mu.Unlock() defer func() { c.mu.Lock() delete(c.pending, seq)