diff --git a/backend/internal/repository/account_repo.go b/backend/internal/repository/account_repo.go index 6ab4446ba4..212523a4fb 100644 --- a/backend/internal/repository/account_repo.go +++ b/backend/internal/repository/account_repo.go @@ -3358,6 +3358,12 @@ func (r *accountRepository) FindByExtraField(ctx context.Context, key string, va // ListDueUpstreamBillingProbeAccounts bounds result hydration and network work // to limit. PostgreSQL must still filter and order all enabled candidates; // MATERIALIZED avoids repeating the defensive timestamp parse expression. +// Go writes next_probe_at via RFC3339Nano (up to 9 fractional digits) while +// jsonpath datetime() parses at most microseconds, so fractions beyond 6 +// digits are trimmed first — mirroring ListDueOllamaCloudUsageAccounts. +// Without this, every nanosecond timestamp is treated as malformed and the +// fail-open ordering pins the cycle to the lowest account IDs, starving the +// rest of the pool. func (r *accountRepository) ListDueUpstreamBillingProbeAccounts(ctx context.Context, now time.Time, limit int) ([]service.Account, error) { if limit <= 0 { return []service.Account{}, nil @@ -3387,7 +3393,11 @@ func (r *accountRepository) ListDueUpstreamBillingProbeAccounts(ctx context.Cont jsonb_path_query_first_tz( jsonb_build_object( 'value', - replace(regexp_replace(next_probe_at, 'Z$', '+00:00'), 'T', ' ') + replace(regexp_replace(regexp_replace( + next_probe_at, + '(\.[0-9]{6})[0-9]+(Z|[+-][0-9]{2}:[0-9]{2})$', + '\1\2' + ), 'Z$', '+00:00'), 'T', ' ') ), '$.value.datetime()', '{}'::jsonb, diff --git a/backend/internal/repository/account_repo_upstream_billing_probe_due_integration_test.go b/backend/internal/repository/account_repo_upstream_billing_probe_due_integration_test.go index 64d29553fe..104771af03 100644 --- a/backend/internal/repository/account_repo_upstream_billing_probe_due_integration_test.go +++ b/backend/internal/repository/account_repo_upstream_billing_probe_due_integration_test.go @@ -49,3 +49,107 @@ func TestListDueUpstreamBillingProbeAccountsHandlesInvalidCalendarDate(t *testin require.Equal(t, invalidID, accounts[0].ID) require.Equal(t, dueID, accounts[1].ID) } + +func insertUpstreamBillingProbeAccount(ctx context.Context, t *testing.T, tx sqlQueryer, name, nextProbeAt string) int64 { + t.Helper() + var id int64 + extra := fmt.Sprintf(`{ + "upstream_billing_probe_enabled": true, + "upstream_billing_probe": {"status": "ok", "next_probe_at": %q} + }`, nextProbeAt) + err := scanSingleRow(ctx, tx, ` + INSERT INTO accounts (name, platform, type, status, extra) + VALUES ($1, 'openai', $2, 'active', $3::jsonb) + RETURNING id + `, []any{name, service.AccountTypeAPIKey, extra}, &id) + require.NoError(t, err) + return id +} + +// rfc3339WithFraction renders t in UTC with an explicit fractional-second +// suffix, e.g. "2026-07-25T17:29:00.123456789Z". The probes' Go writer uses +// RFC3339Nano, so stored fractions carry up to 9 digits. +func rfc3339WithFraction(t time.Time, fraction string) string { + return fmt.Sprintf("%s.%sZ", t.UTC().Format("2006-01-02T15:04:05"), fraction) +} + +// Regression for the scheduled-probe starvation bug: Go persists next_probe_at +// via RFC3339Nano (7-9 fractional digits), which jsonpath datetime() cannot +// parse. Before the trim fix every such row was treated as malformed and +// fail-open "due", so the cycle always returned the lowest account IDs and +// higher IDs were never probed. +func TestListDueUpstreamBillingProbeAccountsParsesNanosecondTimestamps(t *testing.T) { + ctx := context.Background() + tx := testEntTx(t) + repo := newAccountRepositoryWithSQL(tx.Client(), tx, nil) + now := time.Date(2026, time.July, 25, 17, 30, 0, 0, time.UTC) + _, err := tx.ExecContext(ctx, ` + UPDATE accounts + SET extra = extra - 'upstream_billing_probe_enabled' - 'upstream_billing_probe' + `) + require.NoError(t, err) + + // 22 low-ID accounts that are NOT due yet, stored with 9 fractional digits + // exactly as the legacy writer produced them (no migration). + notDue := now.Add(time.Hour) + for i := 0; i < 22; i++ { + insertUpstreamBillingProbeAccount(ctx, t, tx, + fmt.Sprintf("probe-nano-not-due-%02d", i), + rfc3339WithFraction(notDue.Add(time.Duration(i)*time.Second), "123456789")) + } + // 3 high-ID accounts that ARE due, with 9/8/7 fractional digits. Their + // parsed order must follow due time, not insertion/ID order. + dueThird := insertUpstreamBillingProbeAccount(ctx, t, tx, + "probe-nano-due-9digits", rfc3339WithFraction(now.Add(-time.Minute), "123456789")) + dueFirst := insertUpstreamBillingProbeAccount(ctx, t, tx, + "probe-nano-due-8digits", rfc3339WithFraction(now.Add(-3*time.Minute), "12345678")) + dueSecond := insertUpstreamBillingProbeAccount(ctx, t, tx, + "probe-nano-due-7digits", rfc3339WithFraction(now.Add(-2*time.Minute), "1234567")) + + accounts, err := repo.ListDueUpstreamBillingProbeAccounts(ctx, now, 20) + require.NoError(t, err) + require.Len(t, accounts, 3) + require.Equal(t, dueFirst, accounts[0].ID) + require.Equal(t, dueSecond, accounts[1].ID) + require.Equal(t, dueThird, accounts[2].ID) +} + +// With more due accounts than the cycle limit, selection must follow the +// parsed due time so accounts beyond the first IDs still get their +// turn; before the fix the same lowest IDs monopolized every cycle. +func TestListDueUpstreamBillingProbeAccountsSelectsEarliestDueAcrossIDs(t *testing.T) { + ctx := context.Background() + tx := testEntTx(t) + repo := newAccountRepositoryWithSQL(tx.Client(), tx, nil) + now := time.Date(2026, time.July, 25, 17, 30, 0, 0, time.UTC) + _, err := tx.ExecContext(ctx, ` + UPDATE accounts + SET extra = extra - 'upstream_billing_probe_enabled' - 'upstream_billing_probe' + `) + require.NoError(t, err) + + // 25 due accounts; the LOWEST IDs carry the LATEST due times, so a + // correct query must pick the 20 highest-ID rows here. + ids := make([]int64, 0, 25) + for i := 0; i < 25; i++ { + due := now.Add(-time.Duration(i+1) * time.Minute) + ids = append(ids, insertUpstreamBillingProbeAccount(ctx, t, tx, + fmt.Sprintf("probe-nano-order-%02d", i), + rfc3339WithFraction(due, "123456789"))) + } + + accounts, err := repo.ListDueUpstreamBillingProbeAccounts(ctx, now, 20) + require.NoError(t, err) + require.Len(t, accounts, 20) + got := make([]int64, 0, len(accounts)) + for _, account := range accounts { + got = append(got, account.ID) + } + // Earliest due first == highest index first; the five most recently due + // (lowest index / lowest IDs) fall outside the limit this cycle. + want := make([]int64, 0, 20) + for i := 24; i >= 5; i-- { + want = append(want, ids[i]) + } + require.Equal(t, want, got) +} diff --git a/backend/internal/repository/account_repo_upstream_billing_probe_due_test.go b/backend/internal/repository/account_repo_upstream_billing_probe_due_test.go index 350b95d4a2..aa1fcb4303 100644 --- a/backend/internal/repository/account_repo_upstream_billing_probe_due_test.go +++ b/backend/internal/repository/account_repo_upstream_billing_probe_due_test.go @@ -32,6 +32,7 @@ func TestAccountRepositoryListDueUpstreamBillingProbeAccountsBoundsQuery(t *test require.Contains(t, normalized, "type = 'apikey'") require.Contains(t, normalized, `extra @> '{"upstream_billing_probe_enabled": true}'::jsonb`) require.Contains(t, normalized, "jsonb_path_query_first_tz") + require.Contains(t, normalized, `'(\.[0-9]{6})[0-9]+(Z|[+-][0-9]{2}:[0-9]{2})$'`) require.Contains(t, normalized, "parsed AS MATERIALIZED") require.Contains(t, normalized, "parsed_next_probe_at::timestamptz <= $1") require.Contains(t, normalized, "LIMIT $2")