diff --git a/docs/codebase-modularization-roadmap.md b/docs/codebase-modularization-roadmap.md index 0e124b4..86db44c 100644 --- a/docs/codebase-modularization-roadmap.md +++ b/docs/codebase-modularization-roadmap.md @@ -1157,6 +1157,6 @@ global,不能错误归入client。`client-signature-*`、发送页、企业认 - V5 API数据库往返优化仍归`send-inbound-entry`、`send-gateway-submit`和`risk-review`现有边界:入口应用快照沿稳定Facade显式传递,单条快速入队只接受已持久化ID和优先级,通用批量入队保留原查询与取消校验;默认规则并发单飞属于`RiskReviewService`内部完整性保障,不得在`SendChainService`新增第二套规则缓存,也不得把余额、频控状态或实际规则决策缓存进进程内存。 - 500条/秒第一阶段把接收与业务处理边界固定在`CmppInboundSubmissionInbox`:`gateway/internal/inbound`只生成连接域内稳定请求键,`send-inbound-entry`只负责最小校验、Inbox持久化和稳定响应;`api/src/send-worker.ts`在独立进程内领取并调用既有发送链。领取事务不得包住风控、计费、路由、Redis或供应商调用,业务幂等事实必须落PostgreSQL,不能依赖进程内缓存或localStorage。 - API进程角色固定为`CMPP_PROCESS_ROLE=api`,不启动发送Worker、Inbox Worker和周期扫描;`cmpp-send-worker`固定为`worker`并拥有自己的Prisma连接池和回环指标。后续扩容允许水平增加Worker实例,但不得复制HTTP控制器或绕过Inbox直接创建业务消息。 -- 进程隔离同时包含容量隔离:`PrismaService`按进程角色选择独立、有界的连接池上限;Gateway inbound持有与受限Submit窗口匹配的专用API HTTP Transport。连接复用、池上限和超时只属于传输/基础设施边界,不得渗入风控、计费、路由或消息状态机。 +- 进程隔离同时包含容量隔离:`PrismaService`按进程角色选择独立、有界的连接池上限;Gateway inbound分别持有与受限Submit窗口匹配的Submit API Transport和小型后台Transport,协议日志/回执流量不得占用Submit连接。连接复用、池上限和超时只属于传输/基础设施边界,不得渗入风控、计费、路由或消息状态机。 - `api/src/infrastructure-monitoring/`是运营端监控聚合与固定阈值应用边界:只消费代码白名单 PromQL,并通过版本化 PostgreSQL 单例、promtool 校验和原子规则热加载管理数值阈值;Exporter安装、端口隔离、固定规则模板和权限仍归`tools/monitoring/`治理。 - 活动告警已读也归该边界:Prometheus保留告警事实,Prisma仅持久化逐管理员、逐触发周期的阅读状态;全局布局只消费轻量未读汇总,不复制指纹、activeAt或用户隔离逻辑。 diff --git a/docs/contracts/inbound-r6-declarations.json b/docs/contracts/inbound-r6-declarations.json index d71980b..7f2df8d 100644 --- a/docs/contracts/inbound-r6-declarations.json +++ b/docs/contracts/inbound-r6-declarations.json @@ -360,7 +360,7 @@ "name": "Server", "kind": "type", "file": "server.go", - "sha256": "aa919bb151d551a47e1c3fc1191fdd90ee9757f58ef10a718000ff4e8961e4d3" + "sha256": "405e6bca394e2990f7fca11eda9ba45f1915a33f64e1321a4fa63f292e191b88" }, { "name": "boundedSubmitWindow", @@ -564,7 +564,7 @@ "name": "submit", "kind": "func", "file": "submit.go", - "sha256": "b541d2c1d7e5592a6b8ad213d0cdbb81fcc98702023f1d5a4d35207476767ae7" + "sha256": "84c47d59444bf475d416dc9d1cb556b4257b8bd6d453773e01215ab568dc4963" }, { "name": "submitRequest", @@ -612,7 +612,13 @@ "name": "post", "kind": "func", "file": "transport.go", - "sha256": "7982e33b65432936ef97e85357e662a1547c4a93a9852f145127b74ab50609c3" + "sha256": "e4debafbce512419600840f4b8e6ddf8b32ad54b674e3a86f42aea43ea195efc" + }, + { + "name": "postWithClient", + "kind": "func", + "file": "transport.go", + "sha256": "0d98d43b9d8ebdb9cf0e0adf726f7af70fd771b2fbfdf7f6c47ebb0c529a4116" }, { "name": "remoteIP", diff --git a/docs/first-version-development-requirements.md b/docs/first-version-development-requirements.md index 4f3222a..df61df7 100644 --- a/docs/first-version-development-requirements.md +++ b/docs/first-version-development-requirements.md @@ -2108,7 +2108,7 @@ - Inbox的`DateTime`列沿用平台UTC无时区存储口径,领取、租约和退避SQL必须显式使用UTC时钟比较,不能受数据库会话或宿主机Asia/Shanghai时区影响;日期型日配额预留写Prisma时必须传合法Date对象,不能把`YYYY-MM-DD`字符串当作DateTime。 - 优先应用和普通应用共享同一耐久Inbox,但领取顺序必须保留优先级并在同一优先级内FIFO;进入BullMQ后继续沿用priority=1、normal=100。快路径成功只代表平台已可靠接收,异步业务拒绝仍必须落真实消息/任务状态并按既有CMPP失败回执链路通知客户,不得伪装为供应商最终送达。 - 第一阶段的验收是入口可持续接收、Inbox不丢不重且最终可排空,并为后续500条/秒全链路扩容建立解耦边界;不能仅凭SubmitResp吞吐宣称完整500条/秒。压测必须同时报告SubmitResp成功率/延迟、Inbox pending/processing/最老等待、异步完成速率和排空时间,以及命令Stream、结果Outbox和数据库最终对账。 -- Gateway到本机API的HTTP连接池必须按`GATEWAY_CMPP_INBOUND_MAX_CONCURRENCY`设置相同的每主机连接上限与空闲连接上限,复用长连接并继续保留10秒请求超时;不得让Go默认每主机仅保留2条空闲连接在高窗口下制造连接抖动。API与Worker数据库连接池必须分别以`API_DB_POOL_MAX`和`API_WORKER_DB_POOL_MAX`显式有界,默认32/8;API槽位优先保障Inbox受理,Worker不得以默认大池抢占入口连接,所有上限之和必须低于PostgreSQL`max_connections`并为运维连接留余量。 +- Gateway到本机API的Submit专用HTTP连接池必须按`GATEWAY_CMPP_INBOUND_MAX_CONCURRENCY`设置相同的每主机连接上限与空闲连接上限,复用长连接并继续保留10秒请求超时;协议日志、回执恢复和连接状态使用独立的16连接后台池,不得与Submit竞争传输槽,也不得让Go默认每主机仅保留2条空闲连接在高窗口下制造连接抖动。API与Worker数据库连接池必须分别以`API_DB_POOL_MAX`和`API_WORKER_DB_POOL_MAX`显式有界,默认32/8;API槽位优先保障Inbox受理,Worker不得以默认大池抢占入口连接,所有上限之和必须低于PostgreSQL`max_connections`并为运维连接留余量。 - 系统监控标题说明必须明确标注数据来自 Prometheus;“服务关键指标”位于趋势/核心服务区域之后、活动告警之前,并提供统一的“告警阈值设置”入口。安全检测页不重复渲染大号标题,说明文字必须明确标注使用 Fail2ban。 - 告警阈值仅开放固定指标的警告/严重数值,不允许前端提交 PromQL、标签、文件路径或持续时间;必须满足警告值小于严重值。配置以 PostgreSQL 保存版本、期望值、生效值和应用状态,经 `promtool check rules` 校验、同目录原子替换和 Prometheus 热加载成功后才标记生效,失败保留上一生效规则并展示原因。 - 右上角预警中心增加“系统监控告警”,通过独立轻量接口统计 Prometheus 当前 firing/pending 告警及严重数,跳转系统监控活动告警区;任一预警域失败不得清空其他域。 diff --git a/docs/system-functional-test-cases.md b/docs/system-functional-test-cases.md index ee1fa02..ba292ef 100644 --- a/docs/system-functional-test-cases.md +++ b/docs/system-functional-test-cases.md @@ -4741,5 +4741,5 @@ npm run verify:phase8 | TC-CMPP-500-P1-008 | 优先级 | 在普通Inbox积压期间持续混入priority应用 | priority先领取且类内FIFO;普通队列最终可排空;同时记录两类等待分位数 | | TC-CMPP-500-P1-009 | 进程隔离 | 检查systemd、进程、连接和指标端口 | API角色不运行发送后台任务;`cmpp-send-worker`独立非root运行,指标仅监听127.0.0.1:9465 | | TC-CMPP-500-P1-010 | 阶梯压测与对账 | 隔离供应商环境按既定同口径阶梯执行,压后等待全链排空 | 逐档报告SubmitResp、Inbox、两条Stream、最终数据库计数和排空时间;任何丢响应、重复、错误或未排空均失败,不发送真实短信 | -| TC-CMPP-500-P1-011 | Gateway API连接复用 | 将入站全局窗口设为48并构造Gateway专用HTTP客户端 | 每主机最大连接与空闲连接均为48,HTTP总超时仍为10秒;高并发压测不得再出现由连接池抖动造成的10秒API超时 | +| TC-CMPP-500-P1-011 | Gateway API连接复用与流量隔离 | 将入站全局窗口设为48,后台协议日志持续写入,并构造Submit专用HTTP客户端 | Submit池每主机最大连接与空闲连接均为48、后台池独立16连接,HTTP总超时仍为10秒;后台日志/回执不得占用Submit连接或造成10秒API超时 | | TC-CMPP-500-P1-012 | API/Worker数据库池隔离 | API与Worker分别配置32/8连接并在Worker积压时持续提交 | 两进程使用各自有界连接池,API受理连接不被Worker抢占;总连接数不超过PostgreSQL上限且压后无`idle in transaction`泄漏 | diff --git a/docs/testing-progress.md b/docs/testing-progress.md index 9f3f871..5c5d515 100644 --- a/docs/testing-progress.md +++ b/docs/testing-progress.md @@ -3798,3 +3798,6 @@ git diff --check - 首轮50条/秒30秒测试写入1499条Inbox,但Worker因把`YYYY-MM-DD`字符串传给Prisma DateTime以及Asia/Shanghai会话中用`NOW()`比较UTC无时区租约而反复重领。发现后立即停止Worker,未清理Inbox、未重投客户提交。修复提交`53073461e9b991978a3063a691d454b68c20ced9`改为合法UTC Date对象并在领取/退避SQL显式使用`NOW() AT TIME ZONE 'UTC'`;修复前第二套恢复资产位于`/opt/cmpp-platform-backups/p1fix-20260820T175000-before-utc-fix`且全部校验通过。重启后原1499条Inbox无需客户重发即全部完成,日限/频控预留和短信主记录均1499条,双Stream最终0/0。 - UTC修复后的第二轮50条/秒生成1499帧,收到1492个SubmitResp,其中1483成功、9拒绝、7缺响应,P50/P95/P99=`519/3385/3909ms`,按零拒绝/零丢响应停止线仍失败。数据库时间窗内1487条新Inbox最终全部completed。Gateway累计阶段证明1489次API往返成功、10次恰好10秒超时;API自身1493个`POST /api/gateway/events/inbound/submit`全部在1秒内完成,说明未归因尾延迟位于Gateway到本机API的HTTP传输,而不是SubmitResp写Socket。 - 根因补充为Go默认Transport每主机只保留2条空闲连接,无法匹配64槽入站窗口;同时PrismaPg API/Worker仍使用默认连接池上限,无法为入口和重业务事务建立容量隔离。最小补丁新增Gateway专用API Transport,最大/空闲连接数均受`GATEWAY_CMPP_INBOUND_MAX_CONCURRENCY`约束并保留10秒总超时;Prisma按角色显式使用`API_DB_POOL_MAX=32`和`API_WORKER_DB_POOL_MAX=8`。该补丁不改变最小校验、Inbox幂等、风控、计费、路由或状态机,待提交、建立新恢复资产、部署并从50条/秒重新阶梯验证。 +- 容量补丁`6708f1f7c540401d1c92083407c3750693bcbb4b`发布前恢复资产位于`/opt/cmpp-platform-backups/p1pool-20260820T101000-before-capacity-isolation`,PostgreSQL 55MB、运行源码187MB、环境/systemd和发布包均经`pg_restore --list`、tar读取及SHA-256校验。测试环境配置API/Worker池32/8,部署后服务、9464/9465回环指标和数据库连接正常。 +- 补丁后50条/秒30秒生成1499帧,1499/1499成功、零拒绝、零连接错误,P50/P95/P99=`1391/3460/4205ms`,1499条Inbox全部completed,双Stream在注入结束约62秒后归零,入口停止线通过。100条/秒档只生成2719帧(实际90.6条/秒),虽2719/2719成功且P95=4277ms,但客户端窗口被压满1346次,未达到100条/秒目标并停止升档;此时API受理2719个请求全部小于250ms,而Gateway API往返平均3.305秒。 +- 100档同时存在大量供应商Submit、回执和通讯日志;`inbound.Server`此前让Submit、协议日志、回执恢复和连接事件共享同一64连接Transport,后台流量仍能占满Submit传输槽。第二个最小补丁将Submit池按入站窗口独立,后台池固定16连接;保留10秒超时和全部日志/回执功能,不靠删除遥测通过压测。该补丁待专项验证、提交和再次建立恢复资产后发布复测。 diff --git a/gateway/cmd/gateway/main.go b/gateway/cmd/gateway/main.go index b2df089..202754e 100644 --- a/gateway/cmd/gateway/main.go +++ b/gateway/cmd/gateway/main.go @@ -54,7 +54,8 @@ func main() { if err := (inbound.Server{ Addr: cmppAddr, APIBaseURL: apiBaseURL, - HTTPClient: inbound.NewAPIHTTPClient(inboundConcurrency), + HTTPClient: inbound.NewAPIHTTPClient(16), + SubmitHTTPClient: inbound.NewAPIHTTPClient(inboundConcurrency), PresenceStore: presenceStore, RecoveryStore: recoveryStore, GatewayInstanceID: getenv("GATEWAY_INSTANCE_ID", hostname()), diff --git a/gateway/internal/inbound/http_client_test.go b/gateway/internal/inbound/http_client_test.go index 3d33202..14d2bb7 100644 --- a/gateway/internal/inbound/http_client_test.go +++ b/gateway/internal/inbound/http_client_test.go @@ -1,7 +1,10 @@ package inbound import ( + "io" + "net" "net/http" + "strings" "testing" ) @@ -26,3 +29,32 @@ func TestNewAPIHTTPClientUsesSafeDefault(t *testing.T) { t.Fatalf("expected default max connections 64, got %d", transport.MaxConnsPerHost) } } + +func TestSubmitUsesDedicatedHTTPClient(t *testing.T) { + background := &http.Client{Transport: roundTripFunc(func(*http.Request) (*http.Response, error) { + t.Fatal("background client must not carry inbound Submit") + return nil, nil + })} + submit := &http.Client{Transport: roundTripFunc(func(request *http.Request) (*http.Response, error) { + if request.URL.Path != "/api/gateway/events/inbound/submit" { + t.Fatalf("unexpected submit path %s", request.URL.Path) + } + return &http.Response{ + StatusCode: http.StatusCreated, + Status: "201 Created", + Header: make(http.Header), + Body: io.NopCloser(strings.NewReader(`{"accepted":true,"messageId":"MSG-1"}`)), + }, nil + })} + server := Server{APIBaseURL: "http://api.test/api", HTTPClient: background, SubmitHTTPClient: submit} + result, err := server.submit(&net.TCPAddr{IP: net.ParseIP("127.0.0.1"), Port: 12000}, submitRequest{Account: "test"}) + if err != nil || !result.Accepted { + t.Fatalf("expected dedicated submit success, result=%+v err=%v", result, err) + } +} + +type roundTripFunc func(*http.Request) (*http.Response, error) + +func (fn roundTripFunc) RoundTrip(request *http.Request) (*http.Response, error) { + return fn(request) +} diff --git a/gateway/internal/inbound/server.go b/gateway/internal/inbound/server.go index 48b7ea4..78b872c 100644 --- a/gateway/internal/inbound/server.go +++ b/gateway/internal/inbound/server.go @@ -15,6 +15,7 @@ type Server struct { APIBaseURL string SecurityEventToken string HTTPClient *http.Client + SubmitHTTPClient *http.Client LogWriter io.Writer PendingFlushInterval time.Duration PresenceStore PresenceStore diff --git a/gateway/internal/inbound/submit.go b/gateway/internal/inbound/submit.go index 13d13bc..a7c09ff 100644 --- a/gateway/internal/inbound/submit.go +++ b/gateway/internal/inbound/submit.go @@ -290,7 +290,13 @@ func setInboundSubmitResponse(packet any, messageID uint64, result uint32) { func (s Server) submit(remote net.Addr, payload submitRequest) (submitResponse, error) { payload.RemoteIP = remoteIP(remote) var result submitResponse - err := s.post(context.Background(), "/gateway/events/inbound/submit", payload, &result) + client := s.SubmitHTTPClient + if client == nil { + client = s.HTTPClient + } + // Submit has a dedicated transport so protocol logs, receipt recovery and + // presence traffic cannot occupy the connections needed for SubmitResp. + err := s.postWithClient(context.Background(), client, "/gateway/events/inbound/submit", payload, &result) return result, err } diff --git a/gateway/internal/inbound/transport.go b/gateway/internal/inbound/transport.go index befcb75..38cbabb 100644 --- a/gateway/internal/inbound/transport.go +++ b/gateway/internal/inbound/transport.go @@ -15,7 +15,10 @@ import ( const maxAPIResponseBodyBytes int64 = 4 * 1024 * 1024 func (s Server) post(ctx context.Context, path string, payload any, result any) error { - client := s.HTTPClient + return s.postWithClient(ctx, s.HTTPClient, path, payload, result) +} + +func (s Server) postWithClient(ctx context.Context, client *http.Client, path string, payload any, result any) error { if client == nil { client = &http.Client{Timeout: defaultHTTPTimeout} }