fix(grok): 限制导入探测队列并去重任务

This commit is contained in:
IanShaw027
2026-08-07 16:36:12 +08:00
parent a3aae134ee
commit 962da308cc
2 changed files with 79 additions and 4 deletions
@@ -13,6 +13,7 @@ import (
const (
grokImportProbeConcurrency = 3
grokImportProbeTimeout = 25 * time.Second
grokImportProbeQueueLimit = 64
)
type grokImportProber interface {
@@ -27,6 +28,8 @@ type grokImportProbeTask struct {
type grokImportProbeScheduler struct {
mu sync.Mutex
queue []grokImportProbeTask
pending map[int64]struct{}
inFlight map[int64]struct{}
concurrency int
workers int
maxWorkers int
@@ -48,6 +51,8 @@ func newGrokImportProbeScheduler(concurrency int, timeout time.Duration) *grokIm
return &grokImportProbeScheduler{
concurrency: concurrency,
timeout: timeout,
pending: make(map[int64]struct{}),
inFlight: make(map[int64]struct{}),
}
}
@@ -60,7 +65,21 @@ func (s *grokImportProbeScheduler) schedule(prober grokImportProber, account *se
}
s.mu.Lock()
if _, exists := s.pending[account.ID]; exists {
s.mu.Unlock()
return
}
if _, exists := s.inFlight[account.ID]; exists {
s.mu.Unlock()
return
}
if len(s.queue) >= grokImportProbeQueueLimit {
s.mu.Unlock()
slog.Debug("grok_import_active_probe_dropped", "account_id", account.ID, "reason", "queue_full")
return
}
s.queue = append(s.queue, grokImportProbeTask{prober: prober, accountID: account.ID})
s.pending[account.ID] = struct{}{}
if s.workers < s.concurrency {
s.workers++
if s.workers > s.maxWorkers {
@@ -78,6 +97,7 @@ func (s *grokImportProbeScheduler) worker() {
return
}
s.run(task.prober, task.accountID)
s.finish(task.accountID)
}
}
@@ -94,9 +114,17 @@ func (s *grokImportProbeScheduler) nextTask() (grokImportProbeTask, bool) {
if len(s.queue) == 0 {
s.queue = nil
}
delete(s.pending, task.accountID)
s.inFlight[task.accountID] = struct{}{}
return task, true
}
func (s *grokImportProbeScheduler) finish(accountID int64) {
s.mu.Lock()
delete(s.inFlight, accountID)
s.mu.Unlock()
}
func (s *grokImportProbeScheduler) run(prober grokImportProber, accountID int64) {
defer func() {
if recovered := recover(); recovered != nil {
@@ -108,8 +136,6 @@ func (s *grokImportProbeScheduler) run(prober grokImportProber, accountID int64)
}
}()
// Queue time is intentionally excluded: every imported account is probed,
// while this timeout only bounds the actual upstream probe execution.
ctx, cancel := context.WithTimeout(context.Background(), s.timeout)
defer cancel()
result, err := prober.QueryQuota(ctx, accountID)
@@ -142,7 +142,7 @@ func TestGrokImportProbeSchedulerProbesSingleAccountOnce(t *testing.T) {
}
func TestGrokImportProbeSchedulerQueuesBatchWithoutPerTaskGoroutines(t *testing.T) {
const taskCount = 100
const taskCount = 50
release := make(chan struct{})
scheduler := newGrokImportProbeScheduler(3, time.Second)
prober := newGrokImportProbeStub(taskCount)
@@ -156,7 +156,7 @@ func TestGrokImportProbeSchedulerQueuesBatchWithoutPerTaskGoroutines(t *testing.
awaitGrokProbeSignal(t, prober.started)
}
snapshot := snapshotGrokImportProbeScheduler(scheduler)
require.Equal(t, 97, snapshot.queued)
require.Equal(t, taskCount-3, snapshot.queued)
require.Equal(t, 3, snapshot.workers)
require.Equal(t, 3, snapshot.maxWorkers)
select {
@@ -182,6 +182,55 @@ func TestGrokImportProbeSchedulerQueuesBatchWithoutPerTaskGoroutines(t *testing.
require.Equal(t, 3, snapshot.maxWorkers)
}
func TestGrokImportProbeSchedulerDeduplicatesPendingAndInFlightAccounts(t *testing.T) {
scheduler := newGrokImportProbeScheduler(1, time.Second)
prober := newGrokImportProbeStub(2)
release := make(chan struct{})
prober.block = release
account := newGrokOAuthImportAccount(501)
queued := newGrokOAuthImportAccount(502)
scheduler.schedule(prober, account)
require.Equal(t, int64(501), awaitGrokProbeSignal(t, prober.started))
scheduler.schedule(prober, account)
scheduler.schedule(prober, queued)
scheduler.schedule(prober, queued)
scheduler.mu.Lock()
require.Len(t, scheduler.queue, 1)
require.Contains(t, scheduler.inFlight, int64(501))
require.Contains(t, scheduler.pending, int64(502))
scheduler.mu.Unlock()
close(release)
require.Equal(t, int64(501), awaitGrokProbeSignal(t, prober.done))
require.Equal(t, int64(502), awaitGrokProbeSignal(t, prober.done))
calls, _, _ := prober.snapshot()
require.Equal(t, 1, calls[501])
require.Equal(t, 1, calls[502])
}
func TestGrokImportProbeSchedulerBoundsPendingQueue(t *testing.T) {
scheduler := newGrokImportProbeScheduler(1, time.Second)
prober := newGrokImportProbeStub(grokImportProbeQueueLimit + 1)
release := make(chan struct{})
prober.block = release
scheduler.schedule(prober, newGrokOAuthImportAccount(600))
require.Equal(t, int64(600), awaitGrokProbeSignal(t, prober.started))
for id := int64(601); id < 601+grokImportProbeQueueLimit+10; id++ {
scheduler.schedule(prober, newGrokOAuthImportAccount(id))
}
scheduler.mu.Lock()
require.Len(t, scheduler.queue, grokImportProbeQueueLimit)
scheduler.mu.Unlock()
close(release)
for i := 0; i < grokImportProbeQueueLimit+1; i++ {
awaitGrokProbeSignal(t, prober.done)
}
}
func TestGrokImportProbeSchedulerTimeoutCancelsProbe(t *testing.T) {
neverRelease := make(chan struct{})
scheduler := newGrokImportProbeScheduler(1, 20*time.Millisecond)