Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions desktop/backend/services.go
Original file line number Diff line number Diff line change
Expand Up @@ -293,8 +293,8 @@ func NewServices(application *app.Application, info app.Info, env Environment, s
UsageSyncIntervalSeconds: value.UsageSyncIntervalSeconds,
}, err
},
func(ctx context.Context) (usage.UsageSyncResult, error) {
return application.Usage().SyncProviderBackground(ctx, codexconfig.ProviderID)
func(ctx context.Context, onWorkDetected func()) (usage.BackgroundSyncOutcome, error) {
return application.Usage().SyncProviderBackground(ctx, codexconfig.ProviderID, onWorkDetected)
},
)
grokBuildUsageSync := newUsageAutoSyncRuntime(
Expand All @@ -305,8 +305,8 @@ func NewServices(application *app.Application, info app.Info, env Environment, s
UsageSyncIntervalSeconds: value.UsageSyncIntervalSeconds,
}, err
},
func(ctx context.Context) (usage.UsageSyncResult, error) {
return application.Usage().SyncProviderBackground(ctx, grokconfig.ProviderID)
func(ctx context.Context, onWorkDetected func()) (usage.BackgroundSyncOutcome, error) {
return application.Usage().SyncProviderBackground(ctx, grokconfig.ProviderID, onWorkDetected)
},
)
quota := newCodexQuotaRuntime(application.Codex().ListAutomationTargets, application.Codex().RunCredentialJob)
Expand Down
8 changes: 4 additions & 4 deletions desktop/backend/services_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -443,9 +443,9 @@ func TestCodexSettingsServiceKeepsConcurrentUsageIntervalUpdatesConsistent(t *te
}

start := make(chan struct{})
errorsByUpdate := make(chan error, 4)
errorsByUpdate := make(chan error, 5)
var wg sync.WaitGroup
for _, interval := range []int{5, 15, 30, 60} {
for _, interval := range []int{15, 30, 60, 120, 300} {
interval := interval
wg.Add(1)
go func() {
Expand Down Expand Up @@ -493,9 +493,9 @@ func TestGrokBuildSettingsKeepProviderIntervalsIndependentUnderConcurrency(t *te
codexBefore := services.codexUsageSync.Status()

start := make(chan struct{})
errorsByUpdate := make(chan error, 4)
errorsByUpdate := make(chan error, 5)
var wg sync.WaitGroup
for _, interval := range []int{5, 15, 30, 60} {
for _, interval := range []int{15, 30, 60, 120, 300} {
interval := interval
wg.Add(1)
go func() {
Expand Down
102 changes: 85 additions & 17 deletions desktop/backend/usage_auto_sync.go
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,8 @@ type realUsageAutoSyncTicker struct {
ticker *time.Ticker
}

type backgroundUsageSync func(context.Context, func()) (usage.BackgroundSyncOutcome, error)

func (t realUsageAutoSyncTicker) C() <-chan time.Time {
return t.ticker.C
}
Expand All @@ -71,6 +73,7 @@ type usageAutoSyncRuntime struct {
emitter func(UsageAutoSyncStatus)
intervalRevision uint64
syncRevision uint64
syncRunning bool

lifecycleMu sync.Mutex
started bool
Expand All @@ -88,13 +91,13 @@ type usageAutoSyncRuntime struct {
startDelayFunc func() time.Duration
timeout time.Duration
loadSettings func(context.Context) (usage.ProviderSyncSettings, error)
syncProvider func(context.Context) (usage.UsageSyncResult, error)
syncProvider backgroundUsageSync
}

func newUsageAutoSyncRuntime(
providerID string,
loadSettings func(context.Context) (usage.ProviderSyncSettings, error),
syncProvider func(context.Context) (usage.UsageSyncResult, error),
syncProvider backgroundUsageSync,
) *usageAutoSyncRuntime {
providerID = strings.TrimSpace(providerID)
return &usageAutoSyncRuntime{
Expand Down Expand Up @@ -245,11 +248,12 @@ func (r *usageAutoSyncRuntime) SyncNow(ctx context.Context) UsageAutoSyncStatus
r.startSync(r.runCtx)
r.mu.RLock()
done := r.syncDone
running := r.syncRunning
status := cloneUsageAutoSyncStatus(r.status)
r.mu.RUnlock()
r.lifecycleMu.Unlock()

if done == nil || !status.Syncing {
if done == nil || !running {
return status
}
select {
Expand Down Expand Up @@ -342,20 +346,15 @@ func (r *usageAutoSyncRuntime) startSyncWithRevision(
requireRevision bool,
) bool {
r.mu.Lock()
if r.status.Syncing || parent.Err() != nil ||
if r.syncRunning || parent.Err() != nil ||
(requireRevision && r.syncRevision != expectedRevision) {
r.mu.Unlock()
return false
}
r.syncRevision++
r.status.Syncing = true
r.status.Outcome = UsageAutoSyncOutcomeSyncing
r.status.Error = nil
r.status.LastStartedAtUnixMS = r.now().UnixMilli()
r.status.Revision++
r.syncRunning = true
r.syncDone = make(chan struct{})
r.mu.Unlock()
r.emitStatus()

r.workerWG.Add(1)
go func() {
Expand All @@ -366,13 +365,15 @@ func (r *usageAutoSyncRuntime) startSyncWithRevision(
}
ctx, cancel := context.WithTimeout(usage.WithPhaseTimeout(parent, timeout), 2*timeout)
defer cancel()
result, err := r.syncProvider(ctx)
outcome, err := r.runProviderSync(ctx)
if parent.Err() != nil {
r.mu.Lock()
r.status.Syncing = false
r.status.Outcome = UsageAutoSyncOutcomeIdle
r.status.Error = nil
r.status.Revision++
if r.status.Syncing {
r.status.Syncing = false
r.status.Outcome = UsageAutoSyncOutcomeIdle
r.status.Error = nil
r.status.Revision++
}
r.finishSyncLocked()
r.mu.Unlock()
return
Expand All @@ -381,18 +382,83 @@ func (r *usageAutoSyncRuntime) startSyncWithRevision(
r.completeWithError(err)
return
}
r.completeWithResult(result)
if !outcome.Performed {
r.completeWithoutWork(outcome.Result)
return
}
r.markWorkDetected()
r.completeWithResult(outcome.Result)
}()
return true
}

func (r *usageAutoSyncRuntime) runProviderSync(ctx context.Context) (usage.BackgroundSyncOutcome, error) {
if r.syncProvider == nil {
return usage.BackgroundSyncOutcome{}, errors.New("usage sync provider is unavailable")
}
return r.syncProvider(ctx, r.markWorkDetected)
}

func (r *usageAutoSyncRuntime) markWorkDetected() {
r.mu.Lock()
if !r.syncRunning || r.status.Syncing {
r.mu.Unlock()
return
}
r.status.Syncing = true
r.status.Outcome = UsageAutoSyncOutcomeSyncing
r.status.Error = nil
r.status.LastStartedAtUnixMS = r.now().UnixMilli()
r.status.Revision++
r.mu.Unlock()
r.emitStatus()
}

func (r *usageAutoSyncRuntime) completeWithoutWork(result usage.UsageSyncResult) {
r.mu.Lock()
if len(result.Errors) == 0 {
if r.status.Outcome == UsageAutoSyncOutcomeError ||
r.status.Outcome == UsageAutoSyncOutcomeWarning ||
r.status.Error != nil || r.status.ImportErrorCount != 0 {
r.status.Syncing = false
r.status.Outcome = UsageAutoSyncOutcomeIdle
r.status.ImportErrorCount = 0
r.status.Error = nil
r.status.Revision++
r.finishSyncLocked()
r.mu.Unlock()
r.emitStatus()
return
}
r.finishSyncLocked()
r.mu.Unlock()
return
}
count := int64(len(result.Errors))
if r.status.Outcome == UsageAutoSyncOutcomeWarning &&
r.status.ImportErrorCount == count && r.status.Error == nil {
r.finishSyncLocked()
r.mu.Unlock()
return
}
r.status.Syncing = false
r.status.Outcome = UsageAutoSyncOutcomeWarning
r.status.ImportErrorCount = count
r.status.Error = nil
r.status.Revision++
r.finishSyncLocked()
r.mu.Unlock()
r.emitStatus()
}

func (r *usageAutoSyncRuntime) completeWithResult(result usage.UsageSyncResult) {
completedAt := r.now().UnixMilli()
outcome := UsageAutoSyncOutcomeSuccess
if len(result.Errors) > 0 {
outcome = UsageAutoSyncOutcomeWarning
}
r.mu.Lock()
r.syncRunning = false
r.status.Syncing = false
r.status.Outcome = outcome
r.status.LastCompletedAtUnixMS = completedAt
Expand All @@ -408,6 +474,7 @@ func (r *usageAutoSyncRuntime) completeWithResult(result usage.UsageSyncResult)
func (r *usageAutoSyncRuntime) completeWithError(err error) {
completedAt := r.now().UnixMilli()
r.mu.Lock()
r.syncRunning = false
r.status.Syncing = false
r.status.Outcome = UsageAutoSyncOutcomeError
r.status.LastCompletedAtUnixMS = completedAt
Expand All @@ -424,7 +491,7 @@ func (r *usageAutoSyncRuntime) reportStartupError(err error, expectedSyncRevisio
r.mu.Lock()
// SyncNow can join the runtime while its startup settings read is still in
// flight. That read must not complete or supersede the active Provider sync.
if r.status.Syncing || r.syncRevision != expectedSyncRevision {
if r.syncRunning || r.syncRevision != expectedSyncRevision {
r.mu.Unlock()
return
}
Expand All @@ -438,6 +505,7 @@ func (r *usageAutoSyncRuntime) reportStartupError(err error, expectedSyncRevisio
}

func (r *usageAutoSyncRuntime) finishSyncLocked() {
r.syncRunning = false
if r.syncDone == nil {
return
}
Expand Down
Loading