From 2447c44f85ae709766e7a53362dc883a8af0b61e Mon Sep 17 00:00:00 2001 From: Brisbanehuang Date: Sat, 25 Jul 2026 15:14:04 -0400 Subject: [PATCH] fix(repository): parse nanosecond next_probe_at in due probe scheduling Go persists upstream_billing_probe.next_probe_at via RFC3339Nano, but jsonpath datetime() parses at most 6 fractional digits, so every stored timestamp failed to parse and was treated as malformed: fail-open due, ordered into the invalid bucket by id ASC. With more enabled accounts than the per-cycle limit, the same lowest IDs monopolized every cycle and higher IDs were never probed again. Trim the fraction to microseconds before datetime(), mirroring ListDueOllamaCloudUsageAccounts, and pin the behavior with integration regressions: nanosecond parsing, due-time ordering beyond the limit, preserved fail-open for truly invalid dates. --- backend/internal/repository/account_repo.go | 12 +- ...ream_billing_probe_due_integration_test.go | 104 ++++++++++++++++++ ...nt_repo_upstream_billing_probe_due_test.go | 1 + 3 files changed, 116 insertions(+), 1 deletion(-) 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")