diff --git a/backend/internal/repository/batch_image_repo.go b/backend/internal/repository/batch_image_repo.go index 932633eb7b..1d2cf0c13f 100644 --- a/backend/internal/repository/batch_image_repo.go +++ b/backend/internal/repository/batch_image_repo.go @@ -301,18 +301,20 @@ WHERE batch_id = $1 return appendBatchImageEventWithSQL(ctx, sqlq, params.BatchID, "settlement_completed", params.EventPayload) } -func (r *batchImageRepository) SetBatchImageJobSettlementFailed(ctx context.Context, batchID, code, message string) error { - _, err := r.sql.ExecContext(ctx, ` +func (r *batchImageRepository) SetBatchImageJobSettlementFailed(ctx context.Context, batchID, code, message string) (int, error) { + var retryCount int + err := r.sql.QueryRowContext(ctx, ` UPDATE batch_image_jobs SET last_error_code = $2, last_error_message = $3, retry_count = retry_count + 1, updated_at = $4 -WHERE batch_id = $1`, batchID, code, message, time.Now()) +WHERE batch_id = $1 +RETURNING retry_count`, batchID, code, message, time.Now()).Scan(&retryCount) if err != nil { - return err + return 0, translatePersistenceError(err, service.ErrBatchImageJobNotFound, nil) } - return appendBatchImageEventWithSQL(ctx, r.sql, batchID, "settlement_failed", map[string]any{ + return retryCount, appendBatchImageEventWithSQL(ctx, r.sql, batchID, "settlement_failed", map[string]any{ "error_code": code, }) } diff --git a/backend/internal/repository/batch_image_repo_integration_test.go b/backend/internal/repository/batch_image_repo_integration_test.go index 5d43b98a22..973e3d8c2e 100644 --- a/backend/internal/repository/batch_image_repo_integration_test.go +++ b/backend/internal/repository/batch_image_repo_integration_test.go @@ -287,8 +287,9 @@ func TestBatchImageRepository_SetBatchImageJobSettlementFailed(t *testing.T) { }) require.NoError(t, err) - err = repo.SetBatchImageJobSettlementFailed(ctx, batchID, "SETTLEMENT_BILLING_FAILED", "temporary") + retryCount, err := repo.SetBatchImageJobSettlementFailed(ctx, batchID, "SETTLEMENT_BILLING_FAILED", "temporary") require.NoError(t, err) + require.Equal(t, 1, retryCount) job, err := repo.GetBatchImageJobByBatchID(ctx, batchID) require.NoError(t, err) diff --git a/backend/internal/service/batch_image.go b/backend/internal/service/batch_image.go index 992f567841..68dd217f22 100644 --- a/backend/internal/service/batch_image.go +++ b/backend/internal/service/batch_image.go @@ -306,7 +306,7 @@ type BatchImageRepository interface { UpdateBatchImageJobProviderSubmit(ctx context.Context, params UpdateBatchImageJobProviderSubmitParams) error RecordBatchImageJobSubmitFailure(ctx context.Context, batchID, code, message string, markFailed bool) error MarkBatchImageJobSettled(ctx context.Context, params MarkBatchImageJobSettledParams) error - SetBatchImageJobSettlementFailed(ctx context.Context, batchID, code, message string) error + SetBatchImageJobSettlementFailed(ctx context.Context, batchID, code, message string) (int, error) CreateBatchImageItem(ctx context.Context, params CreateBatchImageItemParams) (*BatchImageItem, error) BulkCreateBatchImageItems(ctx context.Context, params []CreateBatchImageItemParams) error ReplaceBatchImageItemsForJob(ctx context.Context, batchID string, items []CreateBatchImageItemParams, counts BatchImageCounts) error diff --git a/backend/internal/service/batch_image_processor_test.go b/backend/internal/service/batch_image_processor_test.go index f4c96f3e19..8a37126161 100644 --- a/backend/internal/service/batch_image_processor_test.go +++ b/backend/internal/service/batch_image_processor_test.go @@ -334,12 +334,13 @@ func (p *fakeProcessorProvider) Cleanup(context.Context, *BatchImageJob, *Accoun } type fakeBatchImageRepository struct { - jobs map[string]*BatchImageJob - items map[string][]CreateBatchImageItemParams - counts map[string]BatchImageCounts - transitions map[string][]string - events map[string][]string - replaceCalls int + jobs map[string]*BatchImageJob + items map[string][]CreateBatchImageItemParams + counts map[string]BatchImageCounts + transitions map[string][]string + events map[string][]string + transitionErr error + replaceCalls int } func newFakeBatchImageRepository() *fakeBatchImageRepository { @@ -473,6 +474,9 @@ func (r *fakeBatchImageRepository) TransitionBatchImageJobStatus(_ context.Conte if !CanTransitionBatchImageJob(job.Status, toStatus) { return ErrBatchImageInvalidTransition } + if r.transitionErr != nil { + return r.transitionErr + } job.Status = toStatus job.LastErrorCode = opts.ErrorCode job.LastErrorMessage = opts.ErrorMessage @@ -558,15 +562,16 @@ func (r *fakeBatchImageRepository) MarkBatchImageJobSettled(_ context.Context, p return nil } -func (r *fakeBatchImageRepository) SetBatchImageJobSettlementFailed(_ context.Context, batchID, code, message string) error { +func (r *fakeBatchImageRepository) SetBatchImageJobSettlementFailed(_ context.Context, batchID, code, message string) (int, error) { job, ok := r.jobs[batchID] if !ok { - return ErrBatchImageJobNotFound + return 0, ErrBatchImageJobNotFound } job.LastErrorCode = batchImageStringPtr(code) job.LastErrorMessage = batchImageOptionalStringPtr(message) + job.RetryCount++ r.events[batchID] = append(r.events[batchID], "settlement_failed") - return nil + return job.RetryCount, nil } func (r *fakeBatchImageRepository) CreateBatchImageItem(_ context.Context, params CreateBatchImageItemParams) (*BatchImageItem, error) { diff --git a/backend/internal/service/batch_image_settlement.go b/backend/internal/service/batch_image_settlement.go index 870883b360..cbd1ca7ae6 100644 --- a/backend/internal/service/batch_image_settlement.go +++ b/backend/internal/service/batch_image_settlement.go @@ -16,6 +16,7 @@ import ( const ( batchImageSettlementRequestPrefix = "batch_image_settlement:" batchImageSettlementRetryDelay = time.Minute + batchImageSettlementMaxRetries = 5 batchImageCostEpsilon = 0.00000001 ) @@ -109,6 +110,9 @@ func (s *BatchImageSettlementService) Settle(ctx context.Context, batchID string if job.AccountID == nil || *job.AccountID <= 0 { return nil, ErrBatchImageSettlementMissingAccountID } + if isBatchImageSettlementRetryExhausted(job) { + return nil, s.failExhaustedSettlement(ctx, job, manifestHash, "settlement billing retry limit reached") + } unitPrice, err := s.settlementUnitPrice(ctx, job) if err != nil { @@ -125,13 +129,17 @@ func (s *BatchImageSettlementService) Settle(ctx context.Context, batchID string } if actualCost-holdAmount > batchImageCostEpsilon { msg := fmt.Sprintf("actual cost %.10f exceeds held amount %.10f", actualCost, holdAmount) - _ = s.Repo.SetBatchImageJobSettlementFailed(ctx, job.BatchID, "SETTLEMENT_COST_EXCEEDS_HOLD", msg) + _, _ = s.Repo.SetBatchImageJobSettlementFailed(ctx, job.BatchID, "SETTLEMENT_COST_EXCEEDS_HOLD", msg) return nil, ErrBatchImageSettlementCostExceedsHold } if err := captureBatchImageBalanceHold(ctx, s.BillingRepo, job, actualCost, manifestHash); err != nil { msg := truncateBatchImageMessage(err.Error(), batchImageMaxErrorMessageLength) - _ = s.Repo.SetBatchImageJobSettlementFailed(ctx, job.BatchID, "SETTLEMENT_BILLING_FAILED", msg) + retryCount, recordErr := s.Repo.SetBatchImageJobSettlementFailed(ctx, job.BatchID, "SETTLEMENT_BILLING_FAILED", msg) + if recordErr == nil && retryCount >= batchImageSettlementMaxRetries { + job.RetryCount = retryCount + return nil, s.failExhaustedSettlement(ctx, job, manifestHash, msg) + } return nil, err } s.invalidateAuthCache(ctx, job.UserID) @@ -160,6 +168,41 @@ func (s *BatchImageSettlementService) Settle(ctx context.Context, batchID string return result, nil } +func isBatchImageSettlementRetryExhausted(job *BatchImageJob) bool { + return job != nil && + job.Status == BatchImageJobStatusSettling && + job.RetryCount >= batchImageSettlementMaxRetries && + batchImageDerefString(job.LastErrorCode) == "SETTLEMENT_BILLING_FAILED" +} + +func (s *BatchImageSettlementService) failExhaustedSettlement(ctx context.Context, job *BatchImageJob, manifestHash, message string) error { + if s == nil || s.Repo == nil { + return ErrBatchImageSettlementBillingFailed + } + if err := releaseBatchImageBalanceHold(ctx, s.BillingRepo, job, manifestHash); err != nil { + msg := truncateBatchImageMessage(err.Error(), batchImageMaxErrorMessageLength) + _, _ = s.Repo.SetBatchImageJobSettlementFailed(ctx, job.BatchID, "SETTLEMENT_RELEASE_FAILED", msg) + return ErrBatchImageSettlementBillingFailed.WithCause(err) + } + s.invalidateAuthCache(ctx, job.UserID) + msg := strings.TrimSpace(message) + if msg == "" { + msg = "settlement billing retry limit reached" + } + if err := s.Repo.TransitionBatchImageJobStatus(ctx, job.BatchID, BatchImageJobStatusFailed, BatchImageTransitionOptions{ + ErrorCode: batchImageStringPtr("SETTLEMENT_BILLING_RETRY_EXHAUSTED"), + ErrorMessage: batchImageStringPtr(msg), + EventType: "settlement_retry_exhausted", + EventPayload: map[string]any{ + "batch_id": job.BatchID, + "retry_count": job.RetryCount, + }, + }); err != nil { + return err + } + return ErrBatchImageSettlementBillingFailed +} + func (s *BatchImageSettlementService) recordUsageLog(ctx context.Context, job *BatchImageJob, actualCost float64, requestID string, createdAt time.Time) { if s == nil || s.UsageLogRepo == nil || job == nil || job.APIKeyID == nil || job.AccountID == nil { return @@ -263,6 +306,10 @@ func (p *BatchImagePipelineProcessor) Process(ctx context.Context, batchID strin _, err := p.SettlementService.Settle(ctx, batchID) if err != nil { if errors.Is(err, ErrBatchImageSettlementBillingFailed) { + updated, getErr := p.ProviderProcessor.Repo.GetBatchImageJobByBatchID(ctx, batchID) + if getErr == nil && IsTerminalBatchImageJobStatus(updated.Status) { + return BatchImageProcessResult{Terminal: true}, nil + } delay := p.RetryDelay if delay <= 0 { delay = batchImageSettlementRetryDelay diff --git a/backend/internal/service/batch_image_settlement_test.go b/backend/internal/service/batch_image_settlement_test.go index c09f0a6cf7..8837a85358 100644 --- a/backend/internal/service/batch_image_settlement_test.go +++ b/backend/internal/service/batch_image_settlement_test.go @@ -231,6 +231,53 @@ func TestBatchImagePipelineProcessor_RequeuesTransientSettlementFailure(t *testi require.Equal(t, BatchImageJobStatusSettling, repo.jobs[job.BatchID].Status) } +func TestBatchImagePipelineProcessor_FailsAndReleasesAfterSettlementRetryLimit(t *testing.T) { + repo := newFakeBatchImageRepository() + job := testSettlingBatchImageJob("imgbatch_pipeline_retry_exhausted") + job.RetryCount = batchImageSettlementMaxRetries - 1 + repo.jobs[job.BatchID] = job + billing := &fakeBatchImageBillingRepo{captureErr: errors.New("temporary billing timeout")} + settlement := &BatchImageSettlementService{Repo: repo, BillingRepo: billing, Pricing: &fakeBatchImagePricingResolver{unitPrice: 0.25}} + processor := &BatchImagePipelineProcessor{ + ProviderProcessor: &BatchImageProviderProcessor{Repo: repo, ProviderRegistry: NewBatchImageProviderRegistry(&fakeProcessorProvider{}), AccountResolver: &fakeBatchImageAccountResolver{account: &Account{}}}, + SettlementService: settlement, + } + + result, err := processor.Process(context.Background(), job.BatchID) + require.NoError(t, err) + require.True(t, result.Terminal) + require.Equal(t, BatchImageJobStatusFailed, repo.jobs[job.BatchID].Status) + require.Equal(t, "SETTLEMENT_BILLING_RETRY_EXHAUSTED", batchImageDerefString(repo.jobs[job.BatchID].LastErrorCode)) + require.Len(t, billing.captures, 1) + require.Len(t, billing.releases, 1) + require.Equal(t, BatchImageReleaseRequestID(job.BatchID), billing.releases[0].RequestID) +} + +func TestBatchImageSettlementRetryExhaustedReleaseIsIdempotentAfterTransitionFailure(t *testing.T) { + repo := newFakeBatchImageRepository() + job := testSettlingBatchImageJob("imgbatch_retry_exhausted_transition_fail") + job.RetryCount = batchImageSettlementMaxRetries + job.LastErrorCode = batchImageStringPtr("SETTLEMENT_BILLING_FAILED") + repo.jobs[job.BatchID] = job + repo.transitionErr = errors.New("temporary transition failure") + billing := &fakeBatchImageBillingRepo{} + svc := &BatchImageSettlementService{Repo: repo, BillingRepo: billing, Pricing: &fakeBatchImagePricingResolver{unitPrice: 0.25}} + + _, err := svc.Settle(context.Background(), job.BatchID) + require.ErrorContains(t, err, "temporary transition failure") + require.Equal(t, BatchImageJobStatusSettling, repo.jobs[job.BatchID].Status) + require.Len(t, billing.releases, 1) + require.Len(t, billing.seen, 1) + + repo.transitionErr = nil + _, err = svc.Settle(context.Background(), job.BatchID) + require.ErrorIs(t, err, ErrBatchImageSettlementBillingFailed) + require.Equal(t, BatchImageJobStatusFailed, repo.jobs[job.BatchID].Status) + require.Len(t, billing.releases, 2) + require.Equal(t, billing.releases[0].RequestID, billing.releases[1].RequestID) + require.Len(t, billing.seen, 1) +} + func TestBatchImageSettlementManifestHash(t *testing.T) { job := testSettlingBatchImageJob("imgbatch_hash") first := BuildBatchImageSettlementManifestHash(job) diff --git a/frontend/src/views/user/BatchImageGuideView.vue b/frontend/src/views/user/BatchImageGuideView.vue index 8f1eebdbe3..afea8d97c0 100644 --- a/frontend/src/views/user/BatchImageGuideView.vue +++ b/frontend/src/views/user/BatchImageGuideView.vue @@ -638,7 +638,7 @@