fix: default gateway callback batching off

This commit is contained in:
hectorzhao
2026-08-26 13:59:21 +08:00
parent 4a17df78b8
commit 8170f727a3
5 changed files with 38 additions and 1 deletions
+1
View File
@@ -57,5 +57,6 @@ GATEWAY_SUBMIT_RESULT_STREAM=gateway.submit.results
GATEWAY_SUBMIT_RESULT_GROUP=cmpp-api-callback GATEWAY_SUBMIT_RESULT_GROUP=cmpp-api-callback
GATEWAY_SUBMIT_RESULT_CONSUMER=gateway-1 GATEWAY_SUBMIT_RESULT_CONSUMER=gateway-1
GATEWAY_SUBMIT_RESULT_WORKER_CONCURRENCY=8 GATEWAY_SUBMIT_RESULT_WORKER_CONCURRENCY=8
GATEWAY_CALLBACK_BATCH_ENABLED=false
GATEWAY_CMPP_INBOUND_MAX_CONCURRENCY=64 GATEWAY_CMPP_INBOUND_MAX_CONCURRENCY=64
CMPP_DOWNSTREAM_ACK_TIMEOUT_SECONDS=30 CMPP_DOWNSTREAM_ACK_TIMEOUT_SECONDS=30
+1
View File
@@ -188,3 +188,4 @@ bash tools/deploy/production-deploy.sh
- 回滚前先停止Gateway、API和Worker,保留故障现场Inbox及日志;如恢复旧数据库备份,必须同时恢复对应源码和环境/systemd资产。不得在回滚时删除pending Inbox或重投真实短信。 - 回滚前先停止Gateway、API和Worker,保留故障现场Inbox及日志;如恢复旧数据库备份,必须同时恢复对应源码和环境/systemd资产。不得在回滚时删除pending Inbox或重投真实短信。
- 第三阶段要求环境显式设置`API_INBOUND_WORKFLOW_BATCH_ENABLED=true`和正整数`API_INBOUND_WORKFLOW_BATCH_SIZE`,初始建议64且不得大于Worker业务槽的可解释倍数。批次增大前必须核对PostgreSQL参数数量、单事务持续时间、Worker RSS和租约时长;付费短信仍走逐条账务锁,不能用零计费批量结果替代付费链路验收。 - 第三阶段要求环境显式设置`API_INBOUND_WORKFLOW_BATCH_ENABLED=true`和正整数`API_INBOUND_WORKFLOW_BATCH_SIZE`,初始建议64且不得大于Worker业务槽的可解释倍数。批次增大前必须核对PostgreSQL参数数量、单事务持续时间、Worker RSS和租约时长;付费短信仍走逐条账务锁,不能用零计费批量结果替代付费链路验收。
- 单企业微批阶段要求显式设置`API_INBOUND_WORKFLOW_BATCH_WAIT_MS=40``API_INBOUND_WORKFLOW_TARGET_BATCH_SIZE=32`,目标不得超过批次上限;发布脚本拒绝负等待、超过250ms或目标越界。API还必须设置`API_HTTP_KEEP_ALIVE_TIMEOUT_MS=120000`和更大的`API_HTTP_HEADERS_TIMEOUT_MS=125000`,确保服务端keep-alive长于Gateway 90秒空闲池;发布后用响应头回读实际timeout并检查Gateway日志无loopback reset。 - 单企业微批阶段要求显式设置`API_INBOUND_WORKFLOW_BATCH_WAIT_MS=40``API_INBOUND_WORKFLOW_TARGET_BATCH_SIZE=32`,目标不得超过批次上限;发布脚本拒绝负等待、超过250ms或目标越界。API还必须设置`API_HTTP_KEEP_ALIVE_TIMEOUT_MS=120000`和更大的`API_HTTP_HEADERS_TIMEOUT_MS=125000`,确保服务端keep-alive长于Gateway 90秒空闲池;发布后用响应头回读实际timeout并检查Gateway日志无loopback reset。
- `GATEWAY_CALLBACK_BATCH_ENABLED`默认关闭,只有独立Gateway callback进程已启用、`GATEWAY_CALLBACK_API_BASE_URL`指向其回环端口且批量路由健康检查通过时才允许显式设为`true`。主API回调模式必须保持`false`;发布后同时检查`gateway.submit.results``pending/lag`,任何持续增长均视为发布失败,禁止手工ACK掩盖回调未持久化。
+8
View File
@@ -1,5 +1,13 @@
# 第一版系统化测试进度 # 第一版系统化测试进度
## 2026-08-26 发送拦截修复与预生产发布
- P0发送拦截修复提交`4a17df78b85f842f724dd0c405febc86f1c34d43`已推送并发布到预生产;部署标记回读一致。发布前有效恢复资产为`/opt/cmpp-platform-backups/releases/20260826T053341Z-before-4a17df7`PostgreSQL custom dump、运行源码、环境/systemd/Nginx配置、RPM/Python包清单及Fail2ban安装文件清单均纳入`SHA256SUMS`并通过校验,`pg_restore --list`和tar可读性通过。
- Alibaba Cloud Linux 4系统仓库不提供兼容的Fail2ban包;预生产按Fail2ban官方固定`1.1.0`源码安装,未增加永久第三方软件源。Fail2ban配置测试、服务状态、SSH、Nginx、nftables空封禁集合和安全代理Socket权限均通过。
- 数据库由90条推进至94条migration,实际新增Submit结果幂等、入站Inbox、Submit Outbox和上行eventId四条迁移。API、独立Send Worker、Gateway、安全代理、Fail2ban、PostgreSQL、Redis、Nginx均activeAPI/Gateway/前端健康,Inbox最终无pending/processingSubmit命令Stream为`pending=0/lag=0`
- 发布后发现主API回调模式下未显式配置批量开关时,Gateway却默认启用批量回调;批量路由只存在于独立callback进程,导致`gateway.submit.results`一度累积到97条pending。预生产显式设置`GATEWAY_CALLBACK_BATCH_ENABLED=false`并重启Gateway后,PEL按幂等逐条回调从97持续降至0并连续六次保持`pending=0/lag=0`,未手工ACK或删除事件。
- 最小代码修复将批量回调改为仅在环境变量严格等于`true`时启用,并增加默认关闭单元测试、示例环境和发布门禁说明。全程未进行压力测试,未由测试人员主动发送、补发或重投短信,未修改通道账号、启停状态、企业余额、白名单、临时号段或测试账号;发布窗口存在预生产客户自然流量,均按真实链路处理。多Gateway P2未实施。
> 环境命名:`8.160.169.106:12026`Web/API)和 `8.160.169.106:17890`(CMPP 入站)实例统一定义为“预生产环境”;`100.93.204.60`统一定义为“虚拟机测试环境”或“测试机”。“虚拟机”不得再用于指代预生产。历史记录中涉及这两个实例的验证和部署按其明确IP归属理解;`production-deploy.sh``NODE_ENV=production`及正式生产安全/备份规范保留原有技术语义,不代表测试机或预生产为正式生产。 > 环境命名:`8.160.169.106:12026`Web/API)和 `8.160.169.106:17890`(CMPP 入站)实例统一定义为“预生产环境”;`100.93.204.60`统一定义为“虚拟机测试环境”或“测试机”。“虚拟机”不得再用于指代预生产。历史记录中涉及这两个实例的验证和部署按其明确IP归属理解;`production-deploy.sh``NODE_ENV=production`及正式生产安全/备份规范保留原有技术语义,不代表测试机或预生产为正式生产。
## 2026-08-12 企业签名弹窗、充值回执、通道列表与金额显示优化(已提交、已部署) ## 2026-08-12 企业签名弹窗、充值回执、通道列表与金额显示优化(已提交、已部署)
+5 -1
View File
@@ -95,7 +95,7 @@ func main() {
resultOutbox.Consumer = getenv("GATEWAY_SUBMIT_RESULT_CONSUMER", "gateway-1") resultOutbox.Consumer = getenv("GATEWAY_SUBMIT_RESULT_CONSUMER", "gateway-1")
resultOutbox.APIBaseURL = callbackBaseURL resultOutbox.APIBaseURL = callbackBaseURL
resultOutbox.Concurrency = positiveEnvInt("GATEWAY_SUBMIT_RESULT_WORKER_CONCURRENCY", 8) resultOutbox.Concurrency = positiveEnvInt("GATEWAY_SUBMIT_RESULT_WORKER_CONCURRENCY", 8)
resultOutbox.BatchEnabled = os.Getenv("GATEWAY_CALLBACK_BATCH_ENABLED") != "false" resultOutbox.BatchEnabled = enabledEnv("GATEWAY_CALLBACK_BATCH_ENABLED")
resultOutbox.BatchSize = positiveEnvInt("GATEWAY_CALLBACK_BATCH_SIZE", 50) resultOutbox.BatchSize = positiveEnvInt("GATEWAY_CALLBACK_BATCH_SIZE", 50)
resultOutbox.BatchWait = time.Duration(positiveEnvInt("GATEWAY_CALLBACK_BATCH_WAIT_MS", 10)) * time.Millisecond resultOutbox.BatchWait = time.Duration(positiveEnvInt("GATEWAY_CALLBACK_BATCH_WAIT_MS", 10)) * time.Millisecond
resultOutbox.GatewayInstanceID = gatewayInstanceID resultOutbox.GatewayInstanceID = gatewayInstanceID
@@ -234,6 +234,10 @@ func positiveEnvInt(key string, fallback int) int {
return value return value
} }
func enabledEnv(key string) bool {
return os.Getenv(key) == "true"
}
func hostname() string { func hostname() string {
name, err := os.Hostname() name, err := os.Hostname()
if err != nil || name == "" { if err != nil || name == "" {
+23
View File
@@ -0,0 +1,23 @@
package main
import "testing"
func TestEnabledEnvRequiresExplicitTrue(t *testing.T) {
for _, testCase := range []struct {
name string
value string
want bool
}{
{name: "unset", want: false},
{name: "false", value: "false", want: false},
{name: "other value", value: "TRUE", want: false},
{name: "true", value: "true", want: true},
} {
t.Run(testCase.name, func(t *testing.T) {
t.Setenv("GATEWAY_CALLBACK_BATCH_ENABLED", testCase.value)
if got := enabledEnv("GATEWAY_CALLBACK_BATCH_ENABLED"); got != testCase.want {
t.Fatalf("enabledEnv() = %v, want %v", got, testCase.want)
}
})
}
}