diff --git a/backend/internal/config/config.go b/backend/internal/config/config.go index 18baa34881..b8b880eab0 100644 --- a/backend/internal/config/config.go +++ b/backend/internal/config/config.go @@ -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) diff --git a/backend/internal/config/config_test.go b/backend/internal/config/config_test.go index bf7a327563..804155d1a9 100644 --- a/backend/internal/config/config_test.go +++ b/backend/internal/config/config_test.go @@ -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) diff --git a/backend/internal/repository/usage_log_repo.go b/backend/internal/repository/usage_log_repo.go index 885e63f9fd..24c648b0a5 100644 --- a/backend/internal/repository/usage_log_repo.go +++ b/backend/internal/repository/usage_log_repo.go @@ -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 } + // 同 ensureCreateBatcher:nil 检查放在 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) + } }) } diff --git a/backend/internal/repository/usage_log_repo_integration_test.go b/backend/internal/repository/usage_log_repo_integration_test.go index ed3050d89c..b43c6d56ef 100644 --- a/backend/internal/repository/usage_log_repo_integration_test.go +++ b/backend/internal/repository/usage_log_repo_integration_test.go @@ -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 到期才标记 dropped(issue #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 persisted(issue #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) { diff --git a/backend/internal/service/gateway_record_usage_test.go b/backend/internal/service/gateway_record_usage_test.go index c819eeca6e..2769251820 100644 --- a/backend/internal/service/gateway_record_usage_test.go +++ b/backend/internal/service/gateway_record_usage_test.go @@ -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) { diff --git a/backend/internal/service/gateway_service.go b/backend/internal/service/gateway_service.go index 160a92a8e1..54035345d9 100644 --- a/backend/internal/service/gateway_service.go +++ b/backend/internal/service/gateway_service.go @@ -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) } } diff --git a/backend/internal/service/usage_record_worker_pool.go b/backend/internal/service/usage_record_worker_pool.go index 5da0b89023..bb5ae452c8 100644 --- a/backend/internal/service/usage_record_worker_pool.go +++ b/backend/internal/service/usage_record_worker_pool.go @@ -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