perf(cmpp): isolate inbound transport capacity

This commit is contained in:
hectorzhao
2026-08-20 18:09:29 +08:00
parent 53073461e9
commit 6708f1f7c5
11 changed files with 95 additions and 8 deletions
+15 -5
View File
@@ -6,14 +6,24 @@ import { requestContext } from '../common/request-context';
@Injectable()
export class PrismaService extends PrismaClient implements OnModuleDestroy {
constructor() {
const databaseUrl = process.env.CMPP_PROCESS_ROLE === 'worker'
const workerRole = process.env.CMPP_PROCESS_ROLE === 'worker';
const databaseUrl = workerRole
? process.env.API_WORKER_DATABASE_URL || process.env.DATABASE_URL
: process.env.DATABASE_URL;
const configuredPoolMax = Number(workerRole
? process.env.API_WORKER_DB_POOL_MAX ?? 8
: process.env.API_DB_POOL_MAX ?? 32);
const poolMax = Number.isInteger(configuredPoolMax) && configuredPoolMax > 0
? configuredPoolMax
: workerRole ? 8 : 32;
super({
adapter: new PrismaPg(
databaseUrl ??
'postgresql://cmpp:cmpp_password@localhost:5432/cmpp_platform?schema=public',
),
adapter: new PrismaPg({
connectionString: databaseUrl
?? 'postgresql://cmpp:cmpp_password@localhost:5432/cmpp_platform?schema=public',
// API capacity must be reserved independently from the heavier Worker
// transactions; explicit bounds also protect PostgreSQL max_connections.
max: poolMax,
}),
});
const operationLog = this.operationLog;
Object.defineProperty(this, 'operationLog', {
+1
View File
@@ -1157,5 +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。连接复用、池上限和超时只属于传输/基础设施边界,不得渗入风控、计费、路由或消息状态机。
- `api/src/infrastructure-monitoring/`是运营端监控聚合与固定阈值应用边界:只消费代码白名单 PromQL,并通过版本化 PostgreSQL 单例、promtool 校验和原子规则热加载管理数值阈值;Exporter安装、端口隔离、固定规则模板和权限仍归`tools/monitoring/`治理。
- 活动告警已读也归该边界:Prometheus保留告警事实,Prisma仅持久化逐管理员、逐触发周期的阅读状态;全局布局只消费轻量未读汇总,不复制指纹、activeAt或用户隔离逻辑。
@@ -2108,6 +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`并为运维连接留余量。
- 系统监控标题说明必须明确标注数据来自 Prometheus;“服务关键指标”位于趋势/核心服务区域之后、活动告警之前,并提供统一的“告警阈值设置”入口。安全检测页不重复渲染大号标题,说明文字必须明确标注使用 Fail2ban。
- 告警阈值仅开放固定指标的警告/严重数值,不允许前端提交 PromQL、标签、文件路径或持续时间;必须满足警告值小于严重值。配置以 PostgreSQL 保存版本、期望值、生效值和应用状态,经 `promtool check rules` 校验、同目录原子替换和 Prometheus 热加载成功后才标记生效,失败保留上一生效规则并展示原因。
- 右上角预警中心增加“系统监控告警”,通过独立轻量接口统计 Prometheus 当前 firing/pending 告警及严重数,跳转系统监控活动告警区;任一预警域失败不得清空其他域。
+4 -2
View File
@@ -33,6 +33,8 @@ API_PORT=3000
API_HOST=127.0.0.1
API_METRICS_HOST=127.0.0.1
API_METRICS_PORT=9464
API_DB_POOL_MAX=32
API_WORKER_DB_POOL_MAX=8
HTTP_API_MASTER_KEY=<至少32位随机值,用于AES-256-GCM加密HTTP访问凭据和Webhook密钥>
HTTP_API_PUBLIC_ORIGIN=https://api.lisglo.com
API_ENABLE_SEND_WORKER=true
@@ -152,7 +154,7 @@ curl http://127.0.0.1:12026/
redis-cli -h 127.0.0.1 -p 6379 ping
pg_isready -d "$(grep '^DATABASE_URL=' /etc/cmpp-platform/cmpp-platform.env | cut -d= -f2-)"
grep -E '^(API_ENABLE_SEND_WORKER|API_SEND_WORKER_CONCURRENCY|GATEWAY_SUBMIT_WORKER_CONCURRENCY|GATEWAY_SUBMIT_RESULT_WORKER_CONCURRENCY|GATEWAY_CMPP_INBOUND_MAX_CONCURRENCY)=' /etc/cmpp-platform/cmpp-platform.env
grep -E '^(CMPP_INBOUND_FAST_PATH_ENABLED|CMPP_INBOUND_WORKFLOW_WORKER_ENABLED|API_INBOUND_WORKFLOW_CONCURRENCY|API_INBOUND_WORKFLOW_POLL_INTERVAL_MS|API_INBOUND_WORKFLOW_STALE_SECONDS|API_WORKER_METRICS_PORT)=' /etc/cmpp-platform/cmpp-platform.env
grep -E '^(CMPP_INBOUND_FAST_PATH_ENABLED|CMPP_INBOUND_WORKFLOW_WORKER_ENABLED|API_INBOUND_WORKFLOW_CONCURRENCY|API_INBOUND_WORKFLOW_POLL_INTERVAL_MS|API_INBOUND_WORKFLOW_STALE_SECONDS|API_DB_POOL_MAX|API_WORKER_DB_POOL_MAX|API_WORKER_METRICS_PORT)=' /etc/cmpp-platform/cmpp-platform.env
systemctl is-active cmpp-api cmpp-send-worker cmpp-gateway
curl -fsS http://127.0.0.1:9465/metrics | grep '^cmpp_worker_inbound_workflow_'
redis-cli --scan --pattern 'rate:gateway:channel:*'
@@ -181,6 +183,6 @@ bash tools/deploy/production-deploy.sh
## CMPP耐久Inbox与独立Worker发布门禁(2026-08-20
- 发布前恢复资产除PostgreSQL、运行源码和环境文件外,必须包含`cmpp-api.service``cmpp-send-worker.service`及其drop-in;逐项校验`pg_restore --list`、tar可读性和SHA-256。数据库回滚与运行代码必须成套执行,禁止只回退代码后让旧Prisma Client访问新状态机。
- 环境必须显式启用`CMPP_INBOUND_FAST_PATH_ENABLED=true``CMPP_INBOUND_WORKFLOW_WORKER_ENABLED=true`,并给出正整数`API_INBOUND_WORKFLOW_CONCURRENCY`;推荐初始值32、轮询100ms、租约300秒。API systemd角色必须是`api`Worker角色必须是`worker`Worker可通过`API_WORKER_DATABASE_URL`使用独立连接上限,未配置时仍使用同一数据库地址但保持独立进程连接池。
- 环境必须显式启用`CMPP_INBOUND_FAST_PATH_ENABLED=true``CMPP_INBOUND_WORKFLOW_WORKER_ENABLED=true`,并给出正整数`API_INBOUND_WORKFLOW_CONCURRENCY``API_DB_POOL_MAX``API_WORKER_DB_POOL_MAX`;推荐初始值分别为32槽、API池32、Worker池8、轮询100ms、租约300秒。API systemd角色必须是`api`Worker角色必须是`worker`Worker可通过`API_WORKER_DATABASE_URL`使用独立地址,未配置时仍使用同一数据库地址但保持独立进程和有界连接池。发布前必须核对两池之和、其他服务连接与运维余量不超过PostgreSQL`max_connections`
- Worker日志目录归`cmpp-api:cmpp-security`且仅服务可写,9465只监听回环并加入Prometheus `cmpp-send-worker` target。发布后必须验证两进程均为非root、API/Gateway health、Worker metrics、PostgreSQL/Redis,以及Inbox pending/processing/最老等待可观测。
- 回滚前先停止Gateway、API和Worker,保留故障现场Inbox及日志;如恢复旧数据库备份,必须同时恢复对应源码和环境/systemd资产。不得在回滚时删除pending Inbox或重投真实短信。
+2
View File
@@ -4741,3 +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-012 | API/Worker数据库池隔离 | API与Worker分别配置32/8连接并在Worker积压时持续提交 | 两进程使用各自有界连接池,API受理连接不被Worker抢占;总连接数不超过PostgreSQL上限且压后无`idle in transaction`泄漏 |
+7
View File
@@ -3791,3 +3791,10 @@ git diff --check
- Prisma新增第92条migration `20260820170000_add_cmpp_inbound_submission_inbox`,包含Inbox、日配额预留和号码频控预留三张真实PostgreSQL表及领取/租约/应用索引。发布脚本增加快路径/Worker强制门禁、systemd分进程、安全drop-in、Worker日志目录、健康检查和Prometheus target/告警;需求、系统用例、模块路线图和生产发布文档同步更新。
- 提交前验证:Prisma generate/validate、API TypeScript正式构建、前端TypeScript与Vite生产构建通过;API全量42套489项通过,Gateway全量`go test ./... -count=1``go vet ./...`通过,5份Stream契约、依赖/安全/部署门禁、R0/R6/R7/R10及`git diff --check`通过。R6契约只新增/更新本轮`submit.go`的稳定请求键声明。R9仍先被HEAD既有`dispatchDueScheduledTasks`哈希漂移阻断,与本轮文件和既有V5记录一致,未为通过本任务错误吸收该并行历史。Jest仍用`--forceExit`收尾仓库既有开放句柄,Vite仅保留既有大chunk告警。
- 当前尚未提交、部署或压测;下一步在排除受保护文件后提交,随后仅对`100.93.204.60`测试环境建立并校验PostgreSQL、运行源码、环境/systemd恢复资产,应用migration和独立Worker,再用隔离供应商执行同口径阶梯压测。预生产不发布、不回退、不压测;不发送、补发或重投真实短信,不修改真实通道账号、密码、启停状态、企业余额或客户连接。
## 2026-08-20 第一阶段测试环境首轮部署、UTC修复与连接池容量补丁
- 第一阶段提交`0b63bcd74e8f8be8b85aaa5f58a5e8ea7fdc636c`已部署到`100.93.204.60`测试环境。发布前恢复资产位于`/opt/cmpp-platform-backups/p1-20260820T173800-before-durable-inbox`,包含可由`pg_restore --list`读取的PostgreSQL custom dump、运行源码、环境/systemd、发布包与`SHA256SUMS`,tar和摘要均通过校验。第92条migration仅在测试数据库应用,API、Gateway、独立Worker和监控均健康;预生产未修改。
- 首轮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条/秒重新阶梯验证。
+3 -1
View File
@@ -50,13 +50,15 @@ func main() {
go func() {
log.Printf("cmpp gateway inbound server listening on %s", cmppAddr)
inboundConcurrency := positiveEnvInt("GATEWAY_CMPP_INBOUND_MAX_CONCURRENCY", 64)
if err := (inbound.Server{
Addr: cmppAddr,
APIBaseURL: apiBaseURL,
HTTPClient: inbound.NewAPIHTTPClient(inboundConcurrency),
PresenceStore: presenceStore,
RecoveryStore: recoveryStore,
GatewayInstanceID: getenv("GATEWAY_INSTANCE_ID", hostname()),
MaxSubmitConcurrency: positiveEnvInt("GATEWAY_CMPP_INBOUND_MAX_CONCURRENCY", 64),
MaxSubmitConcurrency: inboundConcurrency,
SecurityEventToken: os.Getenv("SECURITY_EVENT_TOKEN"),
}).ListenAndServe(); err != nil {
log.Fatalf("gateway inbound server stopped: %v", err)
+26
View File
@@ -0,0 +1,26 @@
package inbound
import (
"net"
"net/http"
"time"
)
// NewAPIHTTPClient keeps enough loopback connections warm for the bounded CMPP
// Submit window. Go's default of two idle connections per host otherwise causes
// connection churn exactly when SubmitResp latency matters most.
func NewAPIHTTPClient(maxConcurrency int) *http.Client {
if maxConcurrency < 1 {
maxConcurrency = 64
}
transport := http.DefaultTransport.(*http.Transport).Clone()
transport.MaxIdleConns = maxConcurrency + 16
transport.MaxIdleConnsPerHost = maxConcurrency
transport.MaxConnsPerHost = maxConcurrency
transport.IdleConnTimeout = 90 * time.Second
transport.DialContext = (&net.Dialer{
Timeout: 3 * time.Second,
KeepAlive: 30 * time.Second,
}).DialContext
return &http.Client{Transport: transport, Timeout: defaultHTTPTimeout}
}
@@ -0,0 +1,28 @@
package inbound
import (
"net/http"
"testing"
)
func TestNewAPIHTTPClientMatchesBoundedSubmitConcurrency(t *testing.T) {
client := NewAPIHTTPClient(48)
transport, ok := client.Transport.(*http.Transport)
if !ok {
t.Fatalf("expected *http.Transport, got %T", client.Transport)
}
if transport.MaxConnsPerHost != 48 || transport.MaxIdleConnsPerHost != 48 {
t.Fatalf("unexpected host connection bounds: max=%d idle=%d", transport.MaxConnsPerHost, transport.MaxIdleConnsPerHost)
}
if client.Timeout != defaultHTTPTimeout {
t.Fatalf("unexpected client timeout: %s", client.Timeout)
}
}
func TestNewAPIHTTPClientUsesSafeDefault(t *testing.T) {
client := NewAPIHTTPClient(0)
transport := client.Transport.(*http.Transport)
if transport.MaxConnsPerHost != 64 {
t.Fatalf("expected default max connections 64, got %d", transport.MaxConnsPerHost)
}
}
+4
View File
@@ -18,6 +18,8 @@ CMPP_INBOUND_WORKFLOW_WORKER_ENABLED="${CMPP_INBOUND_WORKFLOW_WORKER_ENABLED:-tr
API_INBOUND_WORKFLOW_CONCURRENCY="${API_INBOUND_WORKFLOW_CONCURRENCY:-32}"
API_INBOUND_WORKFLOW_POLL_INTERVAL_MS="${API_INBOUND_WORKFLOW_POLL_INTERVAL_MS:-100}"
API_INBOUND_WORKFLOW_STALE_SECONDS="${API_INBOUND_WORKFLOW_STALE_SECONDS:-300}"
API_DB_POOL_MAX="${API_DB_POOL_MAX:-32}"
API_WORKER_DB_POOL_MAX="${API_WORKER_DB_POOL_MAX:-8}"
GATEWAY_CONTROL_ADDR="${GATEWAY_CONTROL_ADDR:-127.0.0.1:8090}"
GATEWAY_CMPP_ADDR="${GATEWAY_CMPP_ADDR:-0.0.0.0:17890}"
CMPP_PUBLIC_HOST="${CMPP_PUBLIC_HOST:-8.160.169.106}"
@@ -205,6 +207,8 @@ CMPP_INBOUND_WORKFLOW_WORKER_ENABLED=${CMPP_INBOUND_WORKFLOW_WORKER_ENABLED}
API_INBOUND_WORKFLOW_CONCURRENCY=${API_INBOUND_WORKFLOW_CONCURRENCY}
API_INBOUND_WORKFLOW_POLL_INTERVAL_MS=${API_INBOUND_WORKFLOW_POLL_INTERVAL_MS}
API_INBOUND_WORKFLOW_STALE_SECONDS=${API_INBOUND_WORKFLOW_STALE_SECONDS}
API_DB_POOL_MAX=${API_DB_POOL_MAX}
API_WORKER_DB_POOL_MAX=${API_WORKER_DB_POOL_MAX}
DATABASE_URL=postgresql://${DB_USER}:${DB_PASSWORD}@127.0.0.1:5432/${DB_NAME}?schema=public
REDIS_HOST=127.0.0.1
REDIS_PORT=6379
+4
View File
@@ -38,6 +38,10 @@ if [[ ! "${API_INBOUND_WORKFLOW_CONCURRENCY:-}" =~ ^[1-9][0-9]*$ ]]; then
echo "API_INBOUND_WORKFLOW_CONCURRENCY must be a positive integer in $ENV_FILE." >&2
exit 1
fi
if [[ ! "${API_DB_POOL_MAX:-}" =~ ^[1-9][0-9]*$ || ! "${API_WORKER_DB_POOL_MAX:-}" =~ ^[1-9][0-9]*$ ]]; then
echo "API_DB_POOL_MAX and API_WORKER_DB_POOL_MAX must be positive integers in $ENV_FILE." >&2
exit 1
fi
if [[ -z "${CMPP_PUBLIC_HOST:-}" || ! "${CMPP_PUBLIC_PORT:-}" =~ ^[1-9][0-9]*$ ]]; then
echo "CMPP_PUBLIC_HOST and a positive CMPP_PUBLIC_PORT are required in $ENV_FILE; these are the customer-facing CMPP endpoint." >&2