Skip to content

Commit 7399de6

Browse files
authored
Merge pull request Wei-Shaw#938 from xvhuan/fix/account-extra-scheduler-pressure-20260311
精准收紧 accounts.extra 观测字段触发的调度重建
2 parents c0110cb + 5c13ec3 commit 7399de6

2 files changed

Lines changed: 80 additions & 96 deletions

File tree

backend/internal/repository/account_repo.go

Lines changed: 24 additions & 81 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,18 @@ type accountRepository struct {
5151
schedulerCache service.SchedulerCache
5252
}
5353

54+
var schedulerNeutralExtraKeyPrefixes = []string{
55+
"codex_primary_",
56+
"codex_secondary_",
57+
"codex_5h_",
58+
"codex_7d_",
59+
}
60+
61+
var schedulerNeutralExtraKeys = map[string]struct{}{
62+
"codex_usage_updated_at": {},
63+
"session_window_utilization": {},
64+
}
65+
5466
// NewAccountRepository 创建账户仓储实例。
5567
// 这是对外暴露的构造函数,返回接口类型以便于依赖注入。
5668
func NewAccountRepository(client *dbent.Client, sqlDB *sql.DB, schedulerCache service.SchedulerCache) service.AccountRepository {
@@ -1190,8 +1202,10 @@ func (r *accountRepository) UpdateExtra(ctx context.Context, id int64, updates m
11901202
if err := enqueueSchedulerOutbox(ctx, r.sql, service.SchedulerOutboxEventAccountChanged, &id, nil, nil); err != nil {
11911203
logger.LegacyPrintf("repository.account", "[SchedulerOutbox] enqueue extra update failed: account=%d err=%v", id, err)
11921204
}
1193-
} else if shouldSyncSchedulerSnapshotForExtraUpdates(updates) {
1194-
// codex 限流快照仍需要让调度缓存尽快看见,避免 DB 抖动时丢失自愈链路。
1205+
} else {
1206+
// 观测型 extra 字段不需要触发 bucket 重建,但仍同步单账号快照,
1207+
// 让 sticky session / GetAccount 命中缓存时也能读到最新数据,
1208+
// 同时避免缓存局部 patch 覆盖掉并发写入的其它账号字段。
11951209
r.syncSchedulerAccountSnapshot(ctx, id)
11961210
}
11971211
return nil
@@ -1202,99 +1216,28 @@ func shouldEnqueueSchedulerOutboxForExtraUpdates(updates map[string]any) bool {
12021216
return false
12031217
}
12041218
for key := range updates {
1205-
if isSchedulerNeutralAccountExtraKey(key) {
1219+
if isSchedulerNeutralExtraKey(key) {
12061220
continue
12071221
}
12081222
return true
12091223
}
12101224
return false
12111225
}
12121226

1213-
func shouldSyncSchedulerSnapshotForExtraUpdates(updates map[string]any) bool {
1214-
return codexExtraIndicatesRateLimit(updates, "7d") || codexExtraIndicatesRateLimit(updates, "5h")
1215-
}
1216-
1217-
func isSchedulerNeutralAccountExtraKey(key string) bool {
1227+
func isSchedulerNeutralExtraKey(key string) bool {
12181228
key = strings.TrimSpace(key)
12191229
if key == "" {
12201230
return false
12211231
}
1222-
if key == "session_window_utilization" {
1232+
if _, ok := schedulerNeutralExtraKeys[key]; ok {
12231233
return true
12241234
}
1225-
return strings.HasPrefix(key, "codex_")
1226-
}
1227-
1228-
func codexExtraIndicatesRateLimit(updates map[string]any, window string) bool {
1229-
if len(updates) == 0 {
1230-
return false
1231-
}
1232-
usedValue, ok := updates["codex_"+window+"_used_percent"]
1233-
if !ok || !extraValueIndicatesExhausted(usedValue) {
1234-
return false
1235-
}
1236-
return extraValueHasResetMarker(updates["codex_"+window+"_reset_at"]) ||
1237-
extraValueHasPositiveNumber(updates["codex_"+window+"_reset_after_seconds"])
1238-
}
1239-
1240-
func extraValueIndicatesExhausted(value any) bool {
1241-
number, ok := extraValueToFloat64(value)
1242-
return ok && number >= 100-1e-9
1243-
}
1244-
1245-
func extraValueHasPositiveNumber(value any) bool {
1246-
number, ok := extraValueToFloat64(value)
1247-
return ok && number > 0
1248-
}
1249-
1250-
func extraValueHasResetMarker(value any) bool {
1251-
switch v := value.(type) {
1252-
case string:
1253-
return strings.TrimSpace(v) != ""
1254-
case time.Time:
1255-
return !v.IsZero()
1256-
case *time.Time:
1257-
return v != nil && !v.IsZero()
1258-
default:
1259-
return false
1260-
}
1261-
}
1262-
1263-
func extraValueToFloat64(value any) (float64, bool) {
1264-
switch v := value.(type) {
1265-
case float64:
1266-
return v, true
1267-
case float32:
1268-
return float64(v), true
1269-
case int:
1270-
return float64(v), true
1271-
case int8:
1272-
return float64(v), true
1273-
case int16:
1274-
return float64(v), true
1275-
case int32:
1276-
return float64(v), true
1277-
case int64:
1278-
return float64(v), true
1279-
case uint:
1280-
return float64(v), true
1281-
case uint8:
1282-
return float64(v), true
1283-
case uint16:
1284-
return float64(v), true
1285-
case uint32:
1286-
return float64(v), true
1287-
case uint64:
1288-
return float64(v), true
1289-
case json.Number:
1290-
parsed, err := v.Float64()
1291-
return parsed, err == nil
1292-
case string:
1293-
parsed, err := strconv.ParseFloat(strings.TrimSpace(v), 64)
1294-
return parsed, err == nil
1295-
default:
1296-
return 0, false
1235+
for _, prefix := range schedulerNeutralExtraKeyPrefixes {
1236+
if strings.HasPrefix(key, prefix) {
1237+
return true
1238+
}
12971239
}
1240+
return false
12981241
}
12991242

13001243
func (r *accountRepository) BulkUpdate(ctx context.Context, ids []int64, updates service.AccountBulkUpdate) (int64, error) {

backend/internal/repository/account_repo_integration_test.go

Lines changed: 56 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ type AccountRepoSuite struct {
2323

2424
type schedulerCacheRecorder struct {
2525
setAccounts []*service.Account
26+
accounts map[int64]*service.Account
2627
}
2728

2829
func (s *schedulerCacheRecorder) GetSnapshot(ctx context.Context, bucket service.SchedulerBucket) ([]*service.Account, bool, error) {
@@ -34,11 +35,20 @@ func (s *schedulerCacheRecorder) SetSnapshot(ctx context.Context, bucket service
3435
}
3536

3637
func (s *schedulerCacheRecorder) GetAccount(ctx context.Context, accountID int64) (*service.Account, error) {
37-
return nil, nil
38+
if s.accounts == nil {
39+
return nil, nil
40+
}
41+
return s.accounts[accountID], nil
3842
}
3943

4044
func (s *schedulerCacheRecorder) SetAccount(ctx context.Context, account *service.Account) error {
4145
s.setAccounts = append(s.setAccounts, account)
46+
if s.accounts == nil {
47+
s.accounts = make(map[int64]*service.Account)
48+
}
49+
if account != nil {
50+
s.accounts[account.ID] = account
51+
}
4252
return nil
4353
}
4454

@@ -623,21 +633,46 @@ func (s *AccountRepoSuite) TestUpdateExtra_NilExtra() {
623633
s.Require().Equal("val", got.Extra["key"])
624634
}
625635

626-
func (s *AccountRepoSuite) TestUpdateExtra_SchedulerNeutralKeysSkipOutbox() {
627-
account := mustCreateAccount(s.T(), s.client, &service.Account{Name: "acc-extra-neutral", Extra: map[string]any{}})
628-
_, err := s.repo.sql.ExecContext(s.ctx, "TRUNCATE scheduler_outbox")
629-
s.Require().NoError(err)
636+
func (s *AccountRepoSuite) TestUpdateExtra_SchedulerNeutralSkipsOutboxAndSyncsFreshSnapshot() {
637+
account := mustCreateAccount(s.T(), s.client, &service.Account{
638+
Name: "acc-extra-neutral",
639+
Platform: service.PlatformOpenAI,
640+
Extra: map[string]any{"codex_usage_updated_at": "old"},
641+
})
642+
cacheRecorder := &schedulerCacheRecorder{
643+
accounts: map[int64]*service.Account{
644+
account.ID: {
645+
ID: account.ID,
646+
Platform: account.Platform,
647+
Status: service.StatusDisabled,
648+
Extra: map[string]any{
649+
"codex_usage_updated_at": "old",
650+
},
651+
},
652+
},
653+
}
654+
s.repo.schedulerCache = cacheRecorder
630655

631-
s.Require().NoError(s.repo.UpdateExtra(s.ctx, account.ID, map[string]any{
632-
"codex_usage_updated_at": "2026-03-11T13:00:00Z",
633-
"codex_5h_used_percent": 12.5,
656+
updates := map[string]any{
657+
"codex_usage_updated_at": "2026-03-11T10:00:00Z",
658+
"codex_5h_used_percent": 88.5,
634659
"session_window_utilization": 0.42,
635-
}))
660+
}
661+
s.Require().NoError(s.repo.UpdateExtra(s.ctx, account.ID, updates))
636662

637-
var count int
638-
err = scanSingleRow(s.ctx, s.repo.sql, "SELECT COUNT(*) FROM scheduler_outbox", nil, &count)
663+
got, err := s.repo.GetByID(s.ctx, account.ID)
639664
s.Require().NoError(err)
640-
s.Require().Equal(0, count)
665+
s.Require().Equal("2026-03-11T10:00:00Z", got.Extra["codex_usage_updated_at"])
666+
s.Require().Equal(88.5, got.Extra["codex_5h_used_percent"])
667+
s.Require().Equal(0.42, got.Extra["session_window_utilization"])
668+
669+
var outboxCount int
670+
s.Require().NoError(scanSingleRow(s.ctx, s.repo.sql, "SELECT COUNT(*) FROM scheduler_outbox", nil, &outboxCount))
671+
s.Require().Zero(outboxCount)
672+
s.Require().Len(cacheRecorder.setAccounts, 1)
673+
s.Require().NotNil(cacheRecorder.accounts[account.ID])
674+
s.Require().Equal(service.StatusActive, cacheRecorder.accounts[account.ID].Status)
675+
s.Require().Equal("2026-03-11T10:00:00Z", cacheRecorder.accounts[account.ID].Extra["codex_usage_updated_at"])
641676
}
642677

643678
func (s *AccountRepoSuite) TestUpdateExtra_ExhaustedCodexSnapshotSyncsSchedulerCache() {
@@ -664,16 +699,22 @@ func (s *AccountRepoSuite) TestUpdateExtra_ExhaustedCodexSnapshotSyncsSchedulerC
664699
s.Require().Equal(0, count)
665700
s.Require().Len(cacheRecorder.setAccounts, 1)
666701
s.Require().Equal(account.ID, cacheRecorder.setAccounts[0].ID)
702+
s.Require().Equal(service.StatusActive, cacheRecorder.setAccounts[0].Status)
667703
s.Require().Equal(100.0, cacheRecorder.setAccounts[0].Extra["codex_7d_used_percent"])
668704
}
669705

670-
func (s *AccountRepoSuite) TestUpdateExtra_CustomKeysStillEnqueueOutbox() {
671-
account := mustCreateAccount(s.T(), s.client, &service.Account{Name: "acc-extra-custom", Extra: map[string]any{}})
706+
func (s *AccountRepoSuite) TestUpdateExtra_SchedulerRelevantStillEnqueuesOutbox() {
707+
account := mustCreateAccount(s.T(), s.client, &service.Account{
708+
Name: "acc-extra-mixed",
709+
Platform: service.PlatformAntigravity,
710+
Extra: map[string]any{},
711+
})
672712
_, err := s.repo.sql.ExecContext(s.ctx, "TRUNCATE scheduler_outbox")
673713
s.Require().NoError(err)
674714

675715
s.Require().NoError(s.repo.UpdateExtra(s.ctx, account.ID, map[string]any{
676-
"custom_scheduler_sensitive_key": true,
716+
"mixed_scheduling": true,
717+
"codex_usage_updated_at": "2026-03-11T10:00:00Z",
677718
}))
678719

679720
var count int

0 commit comments

Comments
 (0)