Merge pull request #3684 from Wei-Shaw/fix/usage-log-queue-silent-drop

fix: prevent silent usage_logs drops under queue overflow
This commit is contained in:
Wesley Liddick
2026-07-03 22:17:24 +08:00
committed by GitHub
7 changed files with 92 additions and 31 deletions
+4 -1
View File
@@ -1945,7 +1945,10 @@ func setDefaults() {
viper.SetDefault("gateway.usage_record.worker_count", 128)
viper.SetDefault("gateway.usage_record.queue_size", 16384)
viper.SetDefault("gateway.usage_record.task_timeout_seconds", 5)
viper.SetDefault("gateway.usage_record.overflow_policy", UsageRecordOverflowPolicySample)
// 默认 sync:队列满时由提交方内联执行(提交点在响应写出之后,不阻塞客户端)。
// sample/drop 会在溢出时静默丢弃计费任务,造成扣费与 usage_logs 对账缺口(issue #3656),
// 仅供显式配置的运维场景使用。
viper.SetDefault("gateway.usage_record.overflow_policy", UsageRecordOverflowPolicySync)
viper.SetDefault("gateway.usage_record.overflow_sample_percent", 10)
viper.SetDefault("gateway.usage_record.auto_scale_enabled", true)
viper.SetDefault("gateway.usage_record.auto_scale_min_workers", 128)
+2 -2
View File
@@ -1903,8 +1903,8 @@ func TestLoad_DefaultGatewayUsageRecordConfig(t *testing.T) {
if cfg.Gateway.UsageRecord.TaskTimeoutSeconds != 5 {
t.Fatalf("task_timeout_seconds = %d, want 5", cfg.Gateway.UsageRecord.TaskTimeoutSeconds)
}
if cfg.Gateway.UsageRecord.OverflowPolicy != UsageRecordOverflowPolicySample {
t.Fatalf("overflow_policy = %s, want %s", cfg.Gateway.UsageRecord.OverflowPolicy, UsageRecordOverflowPolicySample)
if cfg.Gateway.UsageRecord.OverflowPolicy != UsageRecordOverflowPolicySync {
t.Fatalf("overflow_policy = %s, want %s", cfg.Gateway.UsageRecord.OverflowPolicy, UsageRecordOverflowPolicySync)
}
if cfg.Gateway.UsageRecord.OverflowSamplePercent != 10 {
t.Fatalf("overflow_sample_percent = %d, want 10", cfg.Gateway.UsageRecord.OverflowSamplePercent)
+17 -10
View File
@@ -372,12 +372,13 @@ func (r *usageLogRepository) CreateBestEffort(ctx context.Context, log *service.
}
}
// 队列满时阻塞等待而非立即丢弃:批处理器持续排空队列,短暂等待即可入队。
// 立即丢弃会造成“已扣费但无 usage_log”的永久数据缺口(issue #3656);
// 阻塞上限由调用方 ctx 期限约束,超时后由上层同步兜底。
select {
case r.bestEffortBatchCh <- req:
case <-ctx.Done():
return service.MarkUsageLogCreateDropped(ctx.Err())
default:
return service.MarkUsageLogCreateDropped(errors.New("usage log best-effort queue full"))
}
select {
@@ -493,12 +494,12 @@ func (r *usageLogRepository) createBatched(ctx context.Context, log *service.Usa
resultCh: make(chan usageLogCreateResult, 1),
}
// 队列满时阻塞等待而非立即报错:本路径是 best-effort 丢弃后的最后兜底,
// 立即失败会让日志永久丢失;阻塞上限由调用方 ctx 期限约束。
select {
case r.createBatchCh <- req:
case <-ctx.Done():
return false, service.MarkUsageLogCreateNotPersisted(ctx.Err())
default:
return false, service.MarkUsageLogCreateNotPersisted(errors.New("usage log create batch queue full"))
}
select {
@@ -520,22 +521,28 @@ func (r *usageLogRepository) createBatched(ctx context.Context, log *service.Usa
}
func (r *usageLogRepository) ensureCreateBatcher() {
if r == nil || r.db == nil || r.createBatchCh != nil {
if r == nil || r.db == nil {
return
}
// nil 检查必须在 Once 内部:在外层做无同步快路径读会与 Once 内的写构成数据竞争。
r.createBatchOnce.Do(func() {
r.createBatchCh = make(chan usageLogCreateRequest, usageLogCreateBatchQueueCap)
go r.runCreateBatcher(r.db)
if r.createBatchCh == nil {
r.createBatchCh = make(chan usageLogCreateRequest, usageLogCreateBatchQueueCap)
go r.runCreateBatcher(r.db)
}
})
}
func (r *usageLogRepository) ensureBestEffortBatcher() {
if r == nil || r.db == nil || r.bestEffortBatchCh != nil {
if r == nil || r.db == nil {
return
}
// 同 ensureCreateBatchernil 检查放在 Once 内部以避免数据竞争。
r.bestEffortBatchOnce.Do(func() {
r.bestEffortBatchCh = make(chan usageLogBestEffortRequest, usageLogBestEffortBatchQueueCap)
go r.runBestEffortBatcher(r.db)
if r.bestEffortBatchCh == nil {
r.bestEffortBatchCh = make(chan usageLogBestEffortRequest, usageLogBestEffortBatchQueueCap)
go r.runBestEffortBatcher(r.db)
}
})
}
@@ -288,21 +288,21 @@ func TestUsageLogRepositoryCreateBestEffort_BatchPathDuplicateRequestID(t *testi
}, 3*time.Second, 20*time.Millisecond)
}
func TestUsageLogRepositoryCreateBestEffort_QueueFullReturnsDropped(t *testing.T) {
ctx := context.Background()
func TestUsageLogRepositoryCreateBestEffort_QueueFullBlocksUntilCtxDeadline(t *testing.T) {
// 队列满时不再立即丢弃:阻塞等待入队,直到调用方 ctx 到期才标记 droppedissue #3656)。
client := testEntClient(t)
repo := newUsageLogRepositoryWithSQL(client, integrationDB)
repo.bestEffortBatchCh = make(chan usageLogBestEffortRequest, 1)
repo.bestEffortBatchCh <- usageLogBestEffortRequest{}
user := mustCreateUser(t, client, &service.User{Email: fmt.Sprintf("usage-best-effort-full-%d@example.com", time.Now().UnixNano())})
apiKey := mustCreateApiKey(t, client, &service.APIKey{UserID: user.ID, Key: "sk-usage-best-effort-full-" + uuid.NewString(), Name: "k"})
account := mustCreateAccount(t, client, &service.Account{Name: "acc-usage-best-effort-full-" + uuid.NewString()})
ctx, cancel := context.WithTimeout(context.Background(), 200*time.Millisecond)
defer cancel()
start := time.Now()
err := repo.CreateBestEffort(ctx, &service.UsageLog{
UserID: user.ID,
APIKeyID: apiKey.ID,
AccountID: account.ID,
UserID: 1,
APIKeyID: 2,
AccountID: 3,
RequestID: uuid.NewString(),
Model: "claude-3",
InputTokens: 10,
@@ -314,6 +314,40 @@ func TestUsageLogRepositoryCreateBestEffort_QueueFullReturnsDropped(t *testing.T
require.Error(t, err)
require.True(t, service.IsUsageLogCreateDropped(err))
require.GreaterOrEqual(t, time.Since(start), 150*time.Millisecond)
}
func TestUsageLogRepositoryCreateBestEffort_QueueFullWaitsForDrain(t *testing.T) {
// 队列满但批处理器随后排空时,阻塞的入队应成功完成而非丢弃。
client := testEntClient(t)
repo := newUsageLogRepositoryWithSQL(client, integrationDB)
repo.bestEffortBatchCh = make(chan usageLogBestEffortRequest, 1)
repo.bestEffortBatchCh <- usageLogBestEffortRequest{}
go func() {
time.Sleep(100 * time.Millisecond)
<-repo.bestEffortBatchCh // 排空占位请求,为阻塞中的入队腾出空间
req := <-repo.bestEffortBatchCh
sendUsageLogBestEffortResult(req.resultCh, nil)
}()
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
err := repo.CreateBestEffort(ctx, &service.UsageLog{
UserID: 1,
APIKeyID: 2,
AccountID: 3,
RequestID: uuid.NewString(),
Model: "claude-3",
InputTokens: 10,
OutputTokens: 20,
TotalCost: 0.5,
ActualCost: 0.5,
CreatedAt: time.Now().UTC(),
})
require.NoError(t, err)
}
func TestUsageLogRepositoryCreate_BatchPathCanceledContextMarksNotPersisted(t *testing.T) {
@@ -346,7 +380,7 @@ func TestUsageLogRepositoryCreate_BatchPathCanceledContextMarksNotPersisted(t *t
}
func TestUsageLogRepositoryCreate_BatchPathQueueFullMarksNotPersisted(t *testing.T) {
ctx := context.Background()
// 队列满时阻塞等待入队,直到调用方 ctx 到期才标记 not persistedissue #3656)。
client := testEntClient(t)
repo := newUsageLogRepositoryWithSQL(client, integrationDB)
repo.createBatchCh = make(chan usageLogCreateRequest, 1)
@@ -356,6 +390,10 @@ func TestUsageLogRepositoryCreate_BatchPathQueueFullMarksNotPersisted(t *testing
apiKey := mustCreateApiKey(t, client, &service.APIKey{UserID: user.ID, Key: "sk-usage-create-full-" + uuid.NewString(), Name: "k"})
account := mustCreateAccount(t, client, &service.Account{Name: "acc-usage-create-full-" + uuid.NewString()})
ctx, cancel := context.WithTimeout(context.Background(), 200*time.Millisecond)
defer cancel()
start := time.Now()
inserted, err := repo.Create(ctx, &service.UsageLog{
UserID: user.ID,
APIKeyID: apiKey.ID,
@@ -372,6 +410,7 @@ func TestUsageLogRepositoryCreate_BatchPathQueueFullMarksNotPersisted(t *testing
require.False(t, inserted)
require.Error(t, err)
require.True(t, service.IsUsageLogCreateNotPersisted(err))
require.GreaterOrEqual(t, time.Since(start), 150*time.Millisecond)
}
func TestUsageLogRepositoryCreate_BatchPathCanceledAfterQueueMarksNotPersisted(t *testing.T) {
@@ -440,7 +440,9 @@ func TestGatewayServiceRecordUsage_GeneratesRequestIDWhenAllSourcesMissing(t *te
require.Equal(t, billingRepo.lastCmd.RequestID, usageRepo.lastLog.RequestID)
}
func TestGatewayServiceRecordUsage_DroppedUsageLogDoesNotSyncFallback(t *testing.T) {
func TestGatewayServiceRecordUsage_DroppedUsageLogFallsBackToSyncCreate(t *testing.T) {
// 计费成功后 best-effort 写入被丢弃(队列超时)时必须同步兜底,
// 否则出现“已扣费但无 usage_log”的对账缺口(issue #3656)。
usageRepo := &openAIRecordUsageBestEffortLogRepoStub{
bestEffortErr: MarkUsageLogCreateDropped(errors.New("usage log best-effort queue full")),
}
@@ -464,7 +466,9 @@ func TestGatewayServiceRecordUsage_DroppedUsageLogDoesNotSyncFallback(t *testing
require.NoError(t, err)
require.Equal(t, 1, usageRepo.bestEffortCalls)
require.Equal(t, 0, usageRepo.createCalls)
require.Equal(t, 1, usageRepo.createCalls)
// 兜底调用使用的 ctx 必须仍然存活,不能带着已死的 ctx 走过场。
require.NoError(t, usageRepo.lastCtxErr)
}
func TestGatewayServiceRecordUsage_BillingErrorSkipsUsageLogWrite(t *testing.T) {
+10 -3
View File
@@ -9473,10 +9473,17 @@ func writeUsageLogBestEffort(ctx context.Context, repo UsageLogRepository, usage
if writer, ok := repo.(usageLogBestEffortWriter); ok {
if err := writer.CreateBestEffort(usageCtx, usageLog); err != nil {
logger.LegacyPrintf(logKey, "Create usage log failed: %v", err)
if IsUsageLogCreateDropped(err) {
return
// 计费已在此前完成,日志必须落库:dropped(批处理队列超时)同样走同步兜底,
// 否则会出现“已扣费但无 usage_log”的对账缺口(issue #3656)。
// 重复写入由 usage_logs 的 ON CONFLICT (request_id, api_key_id) DO NOTHING 防护。
fallbackCtx := usageCtx
if usageCtx.Err() != nil {
// usageCtx 已耗尽(best-effort 入队阻塞到期限):换新的 detached 窗口,避免兜底必然失败。
var fallbackCancel context.CancelFunc
fallbackCtx, fallbackCancel = detachedBillingContext(context.Background())
defer fallbackCancel()
}
if _, syncErr := repo.Create(usageCtx, usageLog); syncErr != nil {
if _, syncErr := repo.Create(fallbackCtx, usageLog); syncErr != nil {
logger.LegacyPrintf(logKey, "Create usage log sync fallback failed: %v", syncErr)
}
}
@@ -15,10 +15,11 @@ import (
)
const (
defaultUsageRecordWorkerCount = 128
defaultUsageRecordQueueSize = 16384
defaultUsageRecordTaskTimeoutSeconds = 5
defaultUsageRecordOverflowPolicy = config.UsageRecordOverflowPolicySample
defaultUsageRecordWorkerCount = 128
defaultUsageRecordQueueSize = 16384
defaultUsageRecordTaskTimeoutSeconds = 5
// 默认 sync:溢出时提交方内联执行,保证计费任务不被静默丢弃(issue #3656)。
defaultUsageRecordOverflowPolicy = config.UsageRecordOverflowPolicySync
defaultUsageRecordOverflowSampleRatio = 10
defaultUsageRecordAutoScaleEnabled = true
defaultUsageRecordAutoScaleMinWorkers = 128