fix(billing-probe): govern the automatic account rate write-back

- 自动写回值域治理与留痕:上游声明值必须 > 0 且 <= 100 才写回。0 会让
  accountCost 恒为 0(账号总额/日/周配额与成本告警全部静默失效),极大值可
  一次打爆配额并污染成本报表;越界时保持原倍率、记 WARN,探测快照照常记 ok。
  写回成功时记结构化 slog(account_id / 旧值 / 新值 / source),并在快照中
  新增 synced_rate_multiplier 记录本次写回值——后台任务裸 SQL 不产生
  audit_logs,这两处是唯一可追溯来源。管理员手工设 0 不受影响。
- 写回改用 resolved_rate_multiplier(不含高峰的基准倍率):effective 含探测
  那一刻的高峰系数,写回会把一个探测周期的峰值/谷值冻结进静态列,而展示与
  调度用 upstreamBillingRateAt 按当前时间重算高峰,两者会持续不一致。
- 倍率解析失败不再污染公共探测路径:账号级值域/精度只在该账号已开启同步、
  真要写回时才有影响;未开同步的账号照常记 ok 快照,不累计 failure_count、
  不进入指数退避。
- 单账号编辑补 service 层守卫:同步开启时拒绝手工倍率(新增
  UPSTREAM_BILLING_RATE_SYNC_CONFLICT,与批量路径同族),此前只有前端
  disabled 挡人,直接 PUT /admin/accounts/{id} 可写入并活到下次成功探测。
  判断的是本次请求生效后的状态,"关同步 + 改倍率"同请求仍然放行。
- 修正开关反推方向:不再由 rate_sync=true 推出 probe=true,否则一条"同步开、
  探测键缺失"的僵尸记录会在任意一次无关编辑时静默打开周期性外呼;改为探测
  关闭/缺失一律把同步归零。
- 删除死代码 UpdateWithUpstreamBillingProbeEnabled(PR 删接口后生产已无调用
  方),其回滚测试改为直接覆盖生产路径 UpdateWithAccountBillingSettings。
  upstreamBillingRateSyncEnabled 不再是只服务测试的假门控,现为写回前置过滤,
  SQL CAS 仍是权威门控,两侧均加注释说明分工。
- en/zh 文案补充:同步的是不含高峰的基准倍率;开启同步会连带打开自动探测。
This commit is contained in:
shaw
2026-08-01 22:11:10 +08:00
parent b0f5007f04
commit 0b6b4ea956
9 changed files with 421 additions and 45 deletions
+4 -14
View File
@@ -401,17 +401,6 @@ func (r *accountRepository) Update(ctx context.Context, account *service.Account
return r.updateAccount(ctx, account, nil, nil, account.RateMultiplier)
}
// UpdateWithUpstreamBillingProbeEnabled applies an explicit probe switch in the
// same row-lock transaction as the rest of an admin account edit.
func (r *accountRepository) UpdateWithUpstreamBillingProbeEnabled(ctx context.Context, account *service.Account, enabled bool) error {
var rateSyncEnabled *bool
if !enabled {
disabled := false
rateSyncEnabled = &disabled
}
return r.updateAccount(ctx, account, &enabled, rateSyncEnabled, nil)
}
// UpdateWithAccountBillingSettings applies an admin account edit while
// preserving a concurrently probe-synchronized rate unless the request
// explicitly includes a manual rate.
@@ -704,10 +693,11 @@ func lockAndMergeAccountProbeExtra(
if explicitProbeEnabled != nil && !*explicitProbeEnabled {
rateSyncEnabled = false
rateSyncEnabledPresent = true
} else if rateSyncEnabled {
probeEnabled = true
probeEnabledPresent = true
}
// 同步依赖探测,方向是单向的:探测关闭(或探测键缺失)一律把同步归零。
// 不做反向推导——由 rate_sync=true 推出 probe=true 会让一条"同步开、探测键
// 缺失"的僵尸记录在任意一次无关编辑时静默打开周期性外呼。需要同时打开两个
// 开关的调用方(管理端编辑)自己显式传 explicitProbeEnabled=true。
if !probeEnabled {
rateSyncEnabled = false
}
@@ -104,6 +104,102 @@ func TestLockAndMergeAccountProbeExtraUsesCurrentDatabaseSnapshot(t *testing.T)
}
}
func probeBoolPtr(value bool) *bool {
return &value
}
// The probe switch drives the rate-sync switch and never the other way round:
// syncing depends on probing, so a row where sync is on but the probe key is
// missing must lose the sync flag instead of silently gaining periodic
// outbound calls on the next unrelated edit.
func TestLockAndMergeAccountProbeExtraNeverInfersProbeFromRateSync(t *testing.T) {
tests := []struct {
name string
databaseEnabled any
databaseRateSync any
explicitProbeEnabled *bool
explicitRateSync *bool
wantEnabled any
wantRateSync any
}{
{
name: "sync on with missing probe key zeroes sync and keeps probing off",
databaseEnabled: nil,
databaseRateSync: []byte(`true`),
wantEnabled: nil,
wantRateSync: false,
},
{
name: "sync on with probe off zeroes sync",
databaseEnabled: []byte(`false`),
databaseRateSync: []byte(`true`),
wantEnabled: false,
wantRateSync: false,
},
{
name: "both on in database stay on",
databaseEnabled: []byte(`true`),
databaseRateSync: []byte(`true`),
wantEnabled: true,
wantRateSync: true,
},
{
name: "admin enabling both explicitly still turns probing on",
databaseEnabled: nil,
databaseRateSync: nil,
explicitProbeEnabled: probeBoolPtr(true),
explicitRateSync: probeBoolPtr(true),
wantEnabled: true,
wantRateSync: true,
},
{
name: "explicit probe disable clears sync",
databaseEnabled: []byte(`true`),
databaseRateSync: []byte(`true`),
explicitProbeEnabled: probeBoolPtr(false),
wantEnabled: false,
wantRateSync: false,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
db, mock, err := sqlmock.New()
require.NoError(t, err)
t.Cleanup(func() { _ = db.Close() })
client := dbent.NewClient(dbent.Driver(entsql.OpenDB(dialect.Postgres, db)))
t.Cleanup(func() { _ = client.Close() })
mock.ExpectQuery(`(?s)`+regexp.QuoteMeta("SELECT")+`.*`+regexp.QuoteMeta("FOR NO KEY UPDATE")).
WithArgs(int64(31), service.PlatformOpenAI, service.AccountTypeAPIKey, `{"api_key":"sk-test"}`, nil).
WillReturnRows(sqlmock.NewRows([]string{"identity_unchanged", "ollama_group_unchanged", "ollama_proxy_unchanged", "enabled", "rate_sync_enabled", "snapshot", "ollama_session", "ollama_auto", "ollama_snapshot"}).
AddRow(true, false, true, tt.databaseEnabled, tt.databaseRateSync, nil, nil, nil, nil))
account := &service.Account{
ID: 31,
Platform: service.PlatformOpenAI,
Type: service.AccountTypeAPIKey,
Credentials: map[string]any{"api_key": "sk-test"},
}
got, err := lockAndMergeAccountProbeExtra(
context.Background(), client, account, tt.explicitProbeEnabled, tt.explicitRateSync,
)
require.NoError(t, err)
if tt.wantEnabled == nil {
require.NotContains(t, got, service.UpstreamBillingProbeEnabledExtraKey)
} else {
require.Equal(t, tt.wantEnabled, got[service.UpstreamBillingProbeEnabledExtraKey])
}
if tt.wantRateSync == nil {
require.NotContains(t, got, service.UpstreamBillingRateSyncEnabledExtraKey)
} else {
require.Equal(t, tt.wantRateSync, got[service.UpstreamBillingRateSyncEnabledExtraKey])
}
require.NoError(t, mock.ExpectationsWereMet())
})
}
}
func TestLockAndMergeAccountProbeExtraProtectsOllamaManagedFields(t *testing.T) {
for _, identityUnchanged := range []bool{true, false} {
t.Run(map[bool]string{true: "same identity keeps snapshot", false: "changed identity clears snapshot"}[identityUnchanged], func(t *testing.T) {
@@ -265,7 +361,7 @@ func TestUpdateCredentialsAtomicallyClearsProbeForOpenAIAPIKeyIdentityChange(t *
require.NoError(t, mock.ExpectationsWereMet())
}
func TestUpdateWithUpstreamBillingProbeEnabledRollsBackWhenOutboxFails(t *testing.T) {
func TestUpdateWithAccountBillingSettingsRollsBackWhenOutboxFails(t *testing.T) {
db, mock, err := sqlmock.New()
require.NoError(t, err)
t.Cleanup(func() { _ = db.Close() })
@@ -301,7 +397,8 @@ func TestUpdateWithUpstreamBillingProbeEnabledRollsBackWhenOutboxFails(t *testin
Schedulable: true,
}
err = repo.UpdateWithUpstreamBillingProbeEnabled(context.Background(), account, false)
probeDisabled := false
err = repo.UpdateWithAccountBillingSettings(context.Background(), account, &probeDisabled, nil, nil)
require.EqualError(t, err, "outbox failed")
require.Equal(t, false, account.Extra[service.UpstreamBillingProbeEnabledExtraKey])
@@ -792,6 +792,13 @@ func (s *adminServiceImpl) UpdateAccount(ctx context.Context, id int64, input *U
if *input.RateMultiplier < 0 {
return nil, errors.New("rate_multiplier must be >= 0")
}
// 同步开启时倍率归上游所有,手工值活不过下一次成功探测(表现为"改了又自己
// 变回去"),与批量路径一样直接拒绝。判断的是本次请求生效后的状态:上面
// 已把请求携带的两个开关落进 account.Extra,所以"同一请求关闭同步 + 改倍率"
// (用户显式收回所有权)会走到这里时读到 false,正常放行。
if upstreamBillingRateSyncEnabled(account) {
return nil, ErrUpstreamBillingRateSyncConflict
}
account.RateMultiplier = input.RateMultiplier
}
if input.LoadFactor != nil {
@@ -97,8 +97,14 @@ func TestUpdateAccountRoutesRateIntentThroughAtomicBillingUpdater(t *testing.T)
require.Nil(t, repo.lastExplicitRate)
require.Equal(t, concurrentRate, *updated.RateMultiplier)
// 手工倍率只有在同步不再开启时才被接受,所以同一请求先关闭同步再设值
// (同步仍开启时的手工倍率由 TestUpdateAccountRejectsManualRateWhileRateSyncEnabled 覆盖)。
zero := 0.0
updated, err = svc.UpdateAccount(context.Background(), accountID, &UpdateAccountInput{RateMultiplier: &zero})
syncDisabled := false
updated, err = svc.UpdateAccount(context.Background(), accountID, &UpdateAccountInput{
RateSyncEnabled: &syncDisabled,
RateMultiplier: &zero,
})
require.NoError(t, err)
require.Equal(t, 2, repo.updateCalls)
require.NotNil(t, repo.lastExplicitRate)
@@ -426,6 +432,85 @@ func TestUpdateAccountRateSyncControlsProbeAndManualMode(t *testing.T) {
require.Equal(t, false, updated.Extra[UpstreamBillingRateSyncEnabledExtraKey])
}
// 单账号编辑必须和批量路径语义一致:同步开启时倍率归上游所有,手工值会在下一次
// 成功探测时被覆盖,因此直接拒绝而不是静默接受。
func TestUpdateAccountRejectsManualRateWhileRateSyncEnabled(t *testing.T) {
newRepo := func(accountID int64, extra map[string]any) *upstreamBillingProbeAccountRepo {
initialRate := 0.25
return &upstreamBillingProbeAccountRepo{accounts: map[int64]*Account{
accountID: {
ID: accountID,
Platform: PlatformOpenAI,
Type: AccountTypeAPIKey,
Status: StatusActive,
RateMultiplier: &initialRate,
Extra: extra,
},
}}
}
manualRate := 3.5
syncEnabled := map[string]any{
UpstreamBillingProbeEnabledExtraKey: true,
UpstreamBillingRateSyncEnabledExtraKey: true,
}
t.Run("sync enabled rejects manual rate", func(t *testing.T) {
accountID := int64(153)
repo := newRepo(accountID, mergeMap(nil, syncEnabled))
_, err := (&adminServiceImpl{accountRepo: repo}).UpdateAccount(context.Background(), accountID, &UpdateAccountInput{
RateMultiplier: &manualRate,
})
require.ErrorIs(t, err, ErrUpstreamBillingRateSyncConflict)
require.Equal(t, 0.25, *repo.accounts[accountID].RateMultiplier)
})
t.Run("enabling sync in the same request rejects manual rate", func(t *testing.T) {
accountID := int64(154)
repo := newRepo(accountID, map[string]any{})
enable := true
_, err := (&adminServiceImpl{accountRepo: repo}).UpdateAccount(context.Background(), accountID, &UpdateAccountInput{
RateSyncEnabled: &enable,
RateMultiplier: &manualRate,
})
require.ErrorIs(t, err, ErrUpstreamBillingRateSyncConflict)
require.Equal(t, 0.25, *repo.accounts[accountID].RateMultiplier)
})
// 用户显式收回所有权:同一请求关闭同步并改倍率必须放行。
t.Run("disabling sync in the same request allows manual rate", func(t *testing.T) {
accountID := int64(155)
repo := newRepo(accountID, mergeMap(nil, syncEnabled))
disable := false
updated, err := (&adminServiceImpl{accountRepo: repo}).UpdateAccount(context.Background(), accountID, &UpdateAccountInput{
RateSyncEnabled: &disable,
RateMultiplier: &manualRate,
})
require.NoError(t, err)
require.Equal(t, false, updated.Extra[UpstreamBillingRateSyncEnabledExtraKey])
require.NotNil(t, updated.RateMultiplier)
require.Equal(t, manualRate, *updated.RateMultiplier)
})
t.Run("sync disabled allows manual rate", func(t *testing.T) {
accountID := int64(156)
repo := newRepo(accountID, map[string]any{UpstreamBillingProbeEnabledExtraKey: true})
updated, err := (&adminServiceImpl{accountRepo: repo}).UpdateAccount(context.Background(), accountID, &UpdateAccountInput{
RateMultiplier: &manualRate,
})
require.NoError(t, err)
require.NotNil(t, updated.RateMultiplier)
require.Equal(t, manualRate, *updated.RateMultiplier)
})
}
func TestUpdateAccountRejectsSyncWithExplicitlyDisabledProbe(t *testing.T) {
accountID := int64(152)
repo := &upstreamBillingProbeAccountRepo{accounts: map[int64]*Account{
@@ -8,6 +8,7 @@ import (
"errors"
"fmt"
"io"
"log/slog"
"math"
"math/rand/v2"
"net/http"
@@ -46,7 +47,6 @@ const (
// upstreamBillingProbeMaxPerCycle 个名额。
upstreamBillingProbeUnsupportedDelayFactor = 8
upstreamBillingProbeAccountRateScale = 10000.0
upstreamBillingProbeAccountRateMax = 999999.9999
upstreamBillingProbeLeaderLockKey = "upstream:billing:probe:leader"
upstreamBillingProbeLeaderLockTTL = 2 * time.Minute
)
@@ -54,6 +54,21 @@ const (
// UpstreamBillingProbeMaxBatchSize limits one manual batch and one runner cycle.
const UpstreamBillingProbeMaxBatchSize = upstreamBillingProbeMaxPerCycle
// upstreamBillingRateSyncMaxMultiplier bounds the value the automatic
// write-back may push into accounts.rate_multiplier.
//
// No other code path bounds that column from above — admins may type any
// non-negative number and the only ceiling is the DECIMAL(10,4) column itself
// (999999.9999). That ceiling is meaningless as a guard: rate_multiplier
// scales the per-request account cost that feeds quota_used, so a single
// declared 999999 would exhaust any account quota on the first request and
// poison cost reporting. 100 is picked as a deliberately generous bound: it is
// two orders of magnitude above the 1.0 default and far above any plausible
// upstream resale markup, so no legitimate declaration is rejected while an
// absurd or hostile one cannot reach the quota control plane unattended.
// It only constrains the automatic path; manual edits keep their old range.
const upstreamBillingRateSyncMaxMultiplier = 100.0
var (
ErrUpstreamBillingProbeUnavailable = infraerrors.ServiceUnavailable(
"UPSTREAM_BILLING_PROBE_UNAVAILABLE", "upstream billing probe is unavailable",
@@ -68,6 +83,10 @@ var (
"UPSTREAM_BILLING_RATE_SYNC_BULK_CONFLICT",
"account rate multiplier cannot be changed in bulk while upstream billing rate sync is enabled",
)
ErrUpstreamBillingRateSyncConflict = infraerrors.Conflict(
"UPSTREAM_BILLING_RATE_SYNC_CONFLICT",
"account rate multiplier cannot be changed while upstream billing rate sync is enabled",
)
)
const (
@@ -94,6 +113,12 @@ type UpstreamBillingProbeSnapshot struct {
FailureCount int `json:"failure_count,omitempty"`
HTTPStatus int `json:"http_status,omitempty"`
LastError string `json:"last_error,omitempty"`
// SyncedRateMultiplier records the value this probe wrote into
// accounts.rate_multiplier. It is only set when the account opted into rate
// sync and the declared value passed the write-back range check, so the
// stored snapshot always answers "did this probe move the account rate, and
// to what" without a separate history table.
SyncedRateMultiplier *float64 `json:"synced_rate_multiplier,omitempty"`
}
// UpstreamBillingProbeResult is returned by manual probe endpoints.
@@ -644,10 +669,6 @@ func (s *UpstreamBillingProbeService) probeLoadedAccount(ctx context.Context, ac
if err != nil {
return s.persistProbeFailure(ctx, account, intervalMinutes, now, resp.StatusCode, "invalid_response", retryAfter(resp.Header, now))
}
rateMultiplier, ok := upstreamBillingProbeAccountRate(data)
if !ok {
return s.persistProbeFailure(ctx, account, intervalMinutes, now, resp.StatusCode, "invalid_response", retryAfter(resp.Header, now))
}
snapshot := &UpstreamBillingProbeSnapshot{
Status: UpstreamBillingProbeStatusOK,
Data: data,
@@ -657,9 +678,39 @@ func (s *UpstreamBillingProbeService) probeLoadedAccount(ctx context.Context, ac
NextProbeAt: now.Add(nextProbeDelay(intervalMinutes, 0)),
HTTPStatus: resp.StatusCode,
}
if err := s.updateSnapshot(ctx, account, snapshot, &rateMultiplier); err != nil {
// 账号级值域与精度只在真要写回时才有影响:只观察上游声明、未开启同步的
// 账号不因声明值不适配 accounts.rate_multiplier 而被记成探测失败并进入
// 指数退避——探测本身成功了,原始声明照常存进快照供展示。
var syncRate *float64
previousRate := account.BillingRateMultiplier()
if upstreamBillingRateSyncEnabled(account) {
if value, valid := upstreamBillingProbeSyncRate(data); valid {
syncRate = &value
snapshot.SyncedRateMultiplier = &value
} else {
declared, _ := resolveAccountExtraNumber(data, "resolved_rate_multiplier")
slog.Warn("upstream_billing_rate_sync_rejected",
"source", "upstream_billing_probe",
"account_id", account.ID,
"declared_resolved_rate_multiplier", declared,
"max_rate_multiplier", upstreamBillingRateSyncMaxMultiplier,
"current_rate_multiplier", previousRate,
)
}
}
if err := s.updateSnapshot(ctx, account, snapshot, syncRate); err != nil {
return nil, err
}
if syncRate != nil {
// 写回是后台任务的裸 SQL,不经过管理端路由,因此不会产生 audit_logs 行。
// old_rate_multiplier 是本次探测开始时读到的值(写回的 CAS 不比对该列)。
slog.Info("upstream_billing_rate_sync_applied",
"source", "upstream_billing_probe",
"account_id", account.ID,
"old_rate_multiplier", previousRate,
"new_rate_multiplier", *syncRate,
)
}
return snapshot, nil
}
@@ -817,16 +868,32 @@ func upstreamBillingRateAt(data map[string]any, now time.Time) (float64, bool) {
return base, true
}
// upstreamBillingProbeAccountRate converts the declared effective multiplier
// to the precision supported by accounts.rate_multiplier (DECIMAL(10,4)).
func upstreamBillingProbeAccountRate(data map[string]any) (float64, bool) {
value, ok := resolveAccountExtraNumber(data, "effective_rate_multiplier")
if !ok || value < 0 || value > upstreamBillingProbeAccountRateMax ||
math.IsNaN(value) || math.IsInf(value, 0) {
// upstreamBillingProbeSyncRate converts the declared multiplier into the value
// the automatic write-back may store in accounts.rate_multiplier, at the
// precision that column supports (DECIMAL(10,4)).
//
// It reads resolved_rate_multiplier, not effective_rate_multiplier: the
// effective value folds in the peak coefficient that happened to apply at the
// instant of the probe, so writing it would freeze one probe cycle's peak (or
// off-peak) factor into a static column, while display and scheduling
// recompute the peak factor for the current time through upstreamBillingRateAt.
//
// The accepted range is deliberately narrower than the column:
// - 0 is rejected. accountCost multiplies the request cost by this value, so
// an upstream-declared 0 would stop quota_used from ever growing and every
// admin-configured account quota and cost alert would silently stop
// working. Admins may still set 0 by hand; only the automatic path refuses.
// - anything above upstreamBillingRateSyncMaxMultiplier is rejected.
//
// A rejected declaration leaves the current multiplier untouched; the probe
// still records an OK snapshot carrying the raw declaration for display.
func upstreamBillingProbeSyncRate(data map[string]any) (float64, bool) {
value, ok := resolveAccountExtraNumber(data, "resolved_rate_multiplier")
if !ok || math.IsNaN(value) || math.IsInf(value, 0) {
return 0, false
}
rounded := math.Round(value*upstreamBillingProbeAccountRateScale) / upstreamBillingProbeAccountRateScale
if value > 0 && rounded == 0 {
if rounded <= 0 || rounded > upstreamBillingRateSyncMaxMultiplier {
return 0, false
}
return rounded, true
@@ -973,6 +1040,10 @@ func upstreamBillingProbeEnabled(account *Account) bool {
return ok && enabled
}
// upstreamBillingRateSyncEnabled is the probe-side pre-filter deciding whether
// a rate is even proposed for write-back. It is a necessary condition, not the
// authority: the repository CAS re-checks both switches against the row it
// updates, so a switch flipped between load and write can never sneak a rate in.
func upstreamBillingRateSyncEnabled(account *Account) bool {
if account == nil || account.Extra == nil {
return false
@@ -324,8 +324,12 @@ func TestUpstreamBillingProbeSuccessPersistsSanitizedSnapshot(t *testing.T) {
require.Equal(t, fixedNow.Add(time.Hour), *snapshot.FreshUntil)
require.False(t, snapshot.NextProbeAt.Before(fixedNow.Add(24*time.Minute)))
require.False(t, snapshot.NextProbeAt.After(fixedNow.Add(36*time.Minute)))
// 写回的是不含高峰因子的 resolved 倍率(0.6),不是探测那一刻含高峰的
// effective 倍率(0.9)——否则一个探测周期的峰值会被冻结进静态列。
require.NotNil(t, account.RateMultiplier)
require.Equal(t, 0.9, *account.RateMultiplier)
require.Equal(t, 0.6, *account.RateMultiplier)
require.NotNil(t, snapshot.SyncedRateMultiplier)
require.Equal(t, 0.6, *snapshot.SyncedRateMultiplier)
require.Equal(t, "https://upstream.example/v1/sub2api/billing", upstream.lastReq.URL.String())
require.Equal(t, http.MethodGet, upstream.lastReq.Method)
require.Equal(t, "Bearer sk-sensitive", upstream.lastReq.Header.Get("Authorization"))
@@ -336,7 +340,7 @@ func TestUpstreamBillingProbeSuccessPersistsSanitizedSnapshot(t *testing.T) {
require.Equal(t, snapshot.Status, persisted.Status)
}
func TestUpstreamBillingProbeSyncsEffectiveRateForAllAPIKeyPlatforms(t *testing.T) {
func TestUpstreamBillingProbeSyncsResolvedRateForAllAPIKeyPlatforms(t *testing.T) {
for _, platform := range []string{
PlatformOpenAI,
PlatformAnthropic,
@@ -402,24 +406,27 @@ func TestUpstreamBillingProbeOnlyDoesNotChangeAccountRate(t *testing.T) {
require.Contains(t, account.Extra, UpstreamBillingProbeExtraKey)
}
func TestUpstreamBillingProbeAccountRateMatchesDatabasePrecision(t *testing.T) {
func TestUpstreamBillingProbeSyncRateRangeAndPrecision(t *testing.T) {
tests := []struct {
name string
value float64
want float64
ok bool
}{
{name: "zero", value: 0, want: 0, ok: true},
{name: "round to four decimals", value: 0.07654, want: 0.0765, ok: true},
{name: "maximum", value: upstreamBillingProbeAccountRateMax, want: upstreamBillingProbeAccountRateMax, ok: true},
{name: "positive below database precision", value: 0.00001, ok: false},
{name: "overflow before rounding", value: 999999.99991, ok: false},
{name: "overflow after rounding", value: 999999.99996, ok: false},
{name: "maximum", value: upstreamBillingRateSyncMaxMultiplier, want: upstreamBillingRateSyncMaxMultiplier, ok: true},
// 0 会让 accountCost 恒为 0,账号配额与成本告警全部静默失效,
// 自动写回一律拒绝(管理员手工设 0 仍然允许)。
{name: "zero is rejected", value: 0, ok: false},
{name: "positive below database precision rounds to zero", value: 0.00001, ok: false},
{name: "just above the write-back ceiling", value: 100.0001, ok: false},
{name: "column ceiling is far above the write-back ceiling", value: 999999.9999, ok: false},
{name: "negative", value: -1, ok: false},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got, ok := upstreamBillingProbeAccountRate(map[string]any{"effective_rate_multiplier": tt.value})
got, ok := upstreamBillingProbeSyncRate(map[string]any{"resolved_rate_multiplier": tt.value})
require.Equal(t, tt.ok, ok)
if tt.ok {
require.Equal(t, tt.want, got)
@@ -428,6 +435,122 @@ func TestUpstreamBillingProbeAccountRateMatchesDatabasePrecision(t *testing.T) {
}
}
// 只读取 resolved(时间无关的基准倍率):effective 含探测那一刻的高峰系数,
// 写回它会把一个探测周期的峰值/谷值冻结进静态列。
func TestUpstreamBillingProbeSyncRateIgnoresEffectiveRate(t *testing.T) {
got, ok := upstreamBillingProbeSyncRate(map[string]any{
"resolved_rate_multiplier": 0.6,
"effective_rate_multiplier": 0.9,
})
require.True(t, ok)
require.Equal(t, 0.6, got)
_, ok = upstreamBillingProbeSyncRate(map[string]any{"effective_rate_multiplier": 0.9})
require.False(t, ok)
}
// 上游声明超出自动写回值域时保持原倍率,但探测本身是成功的:
// 快照照常记 ok,不累计 failure_count、不进入退避。
func TestUpstreamBillingProbeKeepsRateWhenDeclarationOutOfSyncRange(t *testing.T) {
for _, tt := range []struct {
name string
declared string
}{
{name: "zero", declared: "0"},
{name: "above ceiling", declared: "1000"},
} {
t.Run(tt.name, func(t *testing.T) {
initialRate := 0.25
account := &Account{
ID: 21,
Platform: PlatformOpenAI,
Type: AccountTypeAPIKey,
Status: StatusActive,
Concurrency: 1,
RateMultiplier: &initialRate,
Credentials: map[string]any{
"api_key": "sk-sensitive",
"base_url": "https://upstream.example",
},
Extra: map[string]any{
UpstreamBillingProbeEnabledExtraKey: true,
UpstreamBillingRateSyncEnabledExtraKey: true,
},
}
repo := &upstreamBillingProbeAccountRepo{accounts: map[int64]*Account{account.ID: account}}
upstream := &httpUpstreamRecorder{resp: &http.Response{
StatusCode: http.StatusOK,
Header: http.Header{"Content-Type": []string{"application/json"}},
Body: io.NopCloser(strings.NewReader(fmt.Sprintf(`{
"object":"sub2api.key_billing",
"schema_version":1,
"billing_scope":"token",
"group_rate_multiplier":%[1]s,
"resolved_rate_multiplier":%[1]s,
"peak_rate_enabled":false,
"effective_rate_multiplier":%[1]s,
"observed_at":"2026-07-13T01:00:00Z"
}`, tt.declared))),
}}
svc := newUpstreamBillingProbeTestService(repo, upstream, &upstreamBillingProbeSettingRepo{})
snapshot, err := svc.ProbeAccount(context.Background(), account.ID)
require.NoError(t, err)
require.Equal(t, UpstreamBillingProbeStatusOK, snapshot.Status)
require.Zero(t, snapshot.FailureCount)
require.Nil(t, snapshot.SyncedRateMultiplier)
require.NotNil(t, account.RateMultiplier)
require.Equal(t, initialRate, *account.RateMultiplier)
// 原始声明仍进快照供展示。
require.Equal(t, snapshot.Data["resolved_rate_multiplier"], snapshot.Data["effective_rate_multiplier"])
})
}
}
// 未开启同步的账号只观察上游声明:声明值不适配 accounts.rate_multiplier
// 不得被记成探测失败(否则会累计 failure_count 并进入指数退避)。
func TestUpstreamBillingProbeWithoutSyncIgnoresUnusableDeclaredRate(t *testing.T) {
account := &Account{
ID: 22,
Platform: PlatformOpenAI,
Type: AccountTypeAPIKey,
Status: StatusActive,
Concurrency: 1,
Credentials: map[string]any{
"api_key": "sk-sensitive",
"base_url": "https://upstream.example",
},
Extra: map[string]any{UpstreamBillingProbeEnabledExtraKey: true},
}
repo := &upstreamBillingProbeAccountRepo{accounts: map[int64]*Account{account.ID: account}}
upstream := &httpUpstreamRecorder{resp: &http.Response{
StatusCode: http.StatusOK,
Header: http.Header{"Content-Type": []string{"application/json"}},
Body: io.NopCloser(strings.NewReader(`{
"object":"sub2api.key_billing",
"schema_version":1,
"billing_scope":"token",
"group_rate_multiplier":0,
"resolved_rate_multiplier":0,
"peak_rate_enabled":false,
"effective_rate_multiplier":0,
"observed_at":"2026-07-13T01:00:00Z"
}`)),
}}
svc := newUpstreamBillingProbeTestService(repo, upstream, &upstreamBillingProbeSettingRepo{})
snapshot, err := svc.ProbeAccount(context.Background(), account.ID)
require.NoError(t, err)
require.Equal(t, UpstreamBillingProbeStatusOK, snapshot.Status)
require.Zero(t, snapshot.FailureCount)
require.Empty(t, snapshot.LastError)
require.Nil(t, snapshot.SyncedRateMultiplier)
require.Equal(t, float64(0), snapshot.Data["resolved_rate_multiplier"])
require.Nil(t, account.RateMultiplier)
}
func TestUpstreamBillingProbeRejectsMissingRequiredMultiplier(t *testing.T) {
_, err := parseUpstreamBillingProbeResponse([]byte(`{
"object":"sub2api.key_billing",
@@ -227,9 +227,9 @@ export default {
autoProbe: 'Automatically probe upstream declared rate',
autoProbeHint: 'Refresh the upstream declared rate on the global interval. This switch alone does not change the account rate.',
syncRate: 'Sync upstream declared rate',
syncRateHint: 'Update the account rate after each successful probe. Failed or invalid declarations leave it unchanged.',
syncRateManagedHint: 'The current rate is maintained automatically from the upstream declared rate.',
syncedRateTooltip: 'This account rate is synchronized from the upstream declared rate',
syncRateHint: 'Update the account rate after each successful probe, using the base rate excluding peak hours. Failed probes or declarations outside the allowed range leave it unchanged. Enabling this also turns on "Automatically probe upstream declared rate".',
syncRateManagedHint: 'The current rate is maintained automatically from the upstream declared base rate (excluding peak hours).',
syncedRateTooltip: 'This account rate is synchronized from the upstream declared base rate (excluding peak hours)',
manualProbe: 'Probe upstream rate now',
stale: 'Stale',
unsupported: 'Unsupported',
@@ -180,9 +180,9 @@ export default {
autoProbe: '自动探测上游声明倍率',
autoProbeHint: '启用后按全局周期刷新上游声明倍率;此开关本身不会修改账号倍率。',
syncRate: '同步上游声明倍率',
syncRateHint: '成功探测后自动更新账号倍率;探测失败或声明无效时保持不变。',
syncRateManagedHint: '当前倍率由上游声明倍率自动维护。',
syncedRateTooltip: '该账号倍率由上游声明倍率自动同步',
syncRateHint: '成功探测后自动更新账号倍率,同步的是不含高峰的基准倍率;探测失败或声明超出允许范围时保持不变。开启本项会同时打开“自动探测上游声明倍率”。',
syncRateManagedHint: '当前倍率由上游声明的基准倍率(不含高峰)自动维护。',
syncedRateTooltip: '该账号倍率由上游声明的基准倍率(不含高峰)自动同步',
manualProbe: '立即探测上游倍率',
stale: '已过期',
unsupported: '不支持',
+3
View File
@@ -970,6 +970,9 @@ export interface UpstreamBillingProbeSnapshot {
failure_count?: number
http_status?: number
last_error?: string
// Value this probe wrote into the account rate multiplier; absent when the
// probe did not sync a rate.
synced_rate_multiplier?: number
}
export interface UpstreamBillingProbeSettings {