Compare commits

..
2 Commits
Author SHA1 Message Date
mathiasandClaude Opus 4.8 cc69a912f4 fix(scheduler): rotate lead user each pass so caption budget is shared
CI / Lint / Test / Vet (push) Successful in 12s
CI / Build & Import (push) Successful in 10s
Caption fetches share one per-egress-IP rate budget; whoever runs first each pass
spends the pre-throttle window before YouTube starts 429ing. ListAllUsers order
is unspecified and was stable, so the last-listed user was permanently starved —
a friendly-pilot user got 0 fetches in 12h (all rate_limited) while the
first-listed user got every successful fetch. rotateUsers left-rotates the user
order by pass index so each user leads 1/N passes and the lead slot is shared.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-10 18:07:01 +02:00
mathiasandClaude Opus 4.8 f4a0544903 fix(scheduler): cache transcripts on the scheduled path (ADR-021 regression)
buildUserRunner built the engine without engine.Transcripts = st, so the
scheduler — unlike the web "Summarize now" path — never read or wrote the shared
transcript cache. Every discovery pass re-fetched transcripts it had already
fetched, burning the scarce per-egress-IP caption budget (ADR-014) on redundant
work and starving other users' first-time fetches. The transcripts table was
empty despite summaries existing. Wire the cache on this path too.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-10 18:07:01 +02:00
2 changed files with 67 additions and 6 deletions
+37 -2
View File
@@ -36,6 +36,11 @@ func buildUserRunner(cfg config.Config, st *store.Store, secretStore ports.Secre
}, secretStore) }, secretStore)
engine := usecase.NewEngine(src, buildSummarizer(cfg), st) engine := usecase.NewEngine(src, buildSummarizer(cfg), st)
// Share the transcript cache (ADR-021) on the scheduler path too — without
// this every scheduled pass re-fetches transcripts it already had, burning the
// scarce per-IP caption budget (ADR-014) and starving other users. The web
// "Summarize now" path already sets this; the scheduler omitting it was a bug.
engine.Transcripts = st
return runner.New(src, st, engine, userID, log, return runner.New(src, st, engine, userID, log,
runner.WithBackoff(cfg.FetchBackoff), runner.WithBackoff(cfg.FetchBackoff),
@@ -57,6 +62,7 @@ type userLister interface {
// isolation). Returns the stats summed across users. // isolation). Returns the stats summed across users.
func runDiscoveryPass( func runDiscoveryPass(
ctx context.Context, ctx context.Context,
pass int,
lister userLister, lister userLister,
runUser func(context.Context, string) (runner.Stats, error), runUser func(context.Context, string) (runner.Stats, error),
log *slog.Logger, log *slog.Logger,
@@ -67,6 +73,13 @@ func runDiscoveryPass(
return runner.Stats{} return runner.Stats{}
} }
// Rotate who goes first each pass. Caption fetches share one per-egress-IP
// rate budget (ADR-014); whoever runs first each pass spends the pre-throttle
// window, so a FIXED user order permanently starves whoever is last (a new
// pilot user got 0 fetches for 12h while the first-listed user got all of
// them). Rotation gives every user the lead slot in turn.
users = rotateUsers(users, pass)
log.Info("scheduler: starting discovery pass", "users", len(users)) log.Info("scheduler: starting discovery pass", "users", len(users))
var total runner.Stats var total runner.Stats
for _, u := range users { for _, u := range users {
@@ -116,7 +129,8 @@ func runScheduler(
return // disabled return // disabled
} }
runDiscoveryPass(ctx, lister, runUser, log) pass := 0
runDiscoveryPass(ctx, pass, lister, runUser, log)
ticker := time.NewTicker(interval) ticker := time.NewTicker(interval)
defer ticker.Stop() defer ticker.Stop()
@@ -125,11 +139,32 @@ func runScheduler(
case <-ctx.Done(): case <-ctx.Done():
return return
case <-ticker.C: case <-ticker.C:
runDiscoveryPass(ctx, lister, runUser, log) pass++
runDiscoveryPass(ctx, pass, lister, runUser, log)
} }
} }
} }
// rotateUsers left-rotates users by pass positions so a different user leads each
// pass. With n users, user i leads on every pass where pass ≡ i (mod n). A pass
// offset that is negative or exceeds n is normalised. Order within the rotation
// is otherwise preserved, so the set of users run is unchanged — only who is
// first (and thus wins the scarce caption-fetch budget) rotates.
func rotateUsers(users []store.UserIdentity, pass int) []store.UserIdentity {
n := len(users)
if n <= 1 {
return users
}
off := ((pass % n) + n) % n
if off == 0 {
return users
}
out := make([]store.UserIdentity, 0, n)
out = append(out, users[off:]...)
out = append(out, users[:off]...)
return out
}
// sumStats adds two passes' stats field-wise, so runDiscoveryPass can report a // sumStats adds two passes' stats field-wise, so runDiscoveryPass can report a
// per-tick aggregate across all users. // per-tick aggregate across all users.
func sumStats(a, b runner.Stats) runner.Stats { func sumStats(a, b runner.Stats) runner.Stats {
+30 -4
View File
@@ -45,6 +45,7 @@ func (f fakeLister) ConnectionsForUser(_ context.Context, userID string) ([]stor
type countingRunUser struct { type countingRunUser struct {
mu sync.Mutex mu sync.Mutex
calls map[string]int calls map[string]int
order []string // userIDs in the order they were run, across all passes
failFor map[string]bool failFor map[string]bool
} }
@@ -60,12 +61,19 @@ func (c *countingRunUser) run(_ context.Context, userID string) (runner.Stats, e
c.mu.Lock() c.mu.Lock()
defer c.mu.Unlock() defer c.mu.Unlock()
c.calls[userID]++ c.calls[userID]++
c.order = append(c.order, userID)
if c.failFor[userID] { if c.failFor[userID] {
return runner.Stats{Errors: 1}, errors.New("boom") return runner.Stats{Errors: 1}, errors.New("boom")
} }
return runner.Stats{Summarized: 1}, nil return runner.Stats{Summarized: 1}, nil
} }
func (c *countingRunUser) runOrder() []string {
c.mu.Lock()
defer c.mu.Unlock()
return append([]string(nil), c.order...)
}
func (c *countingRunUser) count(userID string) int { func (c *countingRunUser) count(userID string) int {
c.mu.Lock() c.mu.Lock()
defer c.mu.Unlock() defer c.mu.Unlock()
@@ -94,7 +102,7 @@ func TestDiscoveryPassRunsEveryUserOnce(t *testing.T) {
lister := fakeLister{users: usersN("a", "b", "c")} lister := fakeLister{users: usersN("a", "b", "c")}
rc := newCountingRunUser() rc := newCountingRunUser()
stats := runDiscoveryPass(context.Background(), lister, rc.run, quietLog()) stats := runDiscoveryPass(context.Background(), 0, lister, rc.run, quietLog())
require.Equal(t, 1, rc.count("a")) require.Equal(t, 1, rc.count("a"))
require.Equal(t, 1, rc.count("b")) require.Equal(t, 1, rc.count("b"))
@@ -102,13 +110,31 @@ func TestDiscoveryPassRunsEveryUserOnce(t *testing.T) {
require.Equal(t, 3, stats.Summarized, "stats are summed across users") require.Equal(t, 3, stats.Summarized, "stats are summed across users")
} }
// Caption fetches share one per-IP budget; a fixed user order starves whoever is
// last. Each pass must rotate which user leads so the lead slot is shared.
func TestDiscoveryPassRotatesLeadUser(t *testing.T) {
lister := fakeLister{users: usersN("a", "b", "c")}
rc := newCountingRunUser()
runDiscoveryPass(context.Background(), 0, lister, rc.run, quietLog())
runDiscoveryPass(context.Background(), 1, lister, rc.run, quietLog())
runDiscoveryPass(context.Background(), 2, lister, rc.run, quietLog())
require.Equal(t, []string{"a", "b", "c", "b", "c", "a", "c", "a", "b"}, rc.runOrder(),
"each pass left-rotates the user order so every user leads in turn")
// Fairness: over a full rotation cycle every user ran the same number of times.
require.Equal(t, 3, rc.count("a"))
require.Equal(t, 3, rc.count("b"))
require.Equal(t, 3, rc.count("c"))
}
func TestDiscoveryPassSkipsUsersWithoutConnections(t *testing.T) { func TestDiscoveryPassSkipsUsersWithoutConnections(t *testing.T) {
// b never connected a video source (e.g. a stale Dex-era orphan identity). // b never connected a video source (e.g. a stale Dex-era orphan identity).
// It must be skipped silently — not run and logged as a token error every pass. // It must be skipped silently — not run and logged as a token error every pass.
lister := fakeLister{users: usersN("a", "b", "c"), noConn: map[string]bool{"b": true}} lister := fakeLister{users: usersN("a", "b", "c"), noConn: map[string]bool{"b": true}}
rc := newCountingRunUser() rc := newCountingRunUser()
stats := runDiscoveryPass(context.Background(), lister, rc.run, quietLog()) stats := runDiscoveryPass(context.Background(), 0, lister, rc.run, quietLog())
require.Equal(t, 1, rc.count("a")) require.Equal(t, 1, rc.count("a"))
require.Equal(t, 0, rc.count("b"), "a user with no connection must be skipped, not run") require.Equal(t, 0, rc.count("b"), "a user with no connection must be skipped, not run")
@@ -121,7 +147,7 @@ func TestDiscoveryPassOneUserFailureDoesNotStopOthers(t *testing.T) {
lister := fakeLister{users: usersN("a", "b", "c")} lister := fakeLister{users: usersN("a", "b", "c")}
rc := newCountingRunUser("b") // user b's pass errors rc := newCountingRunUser("b") // user b's pass errors
stats := runDiscoveryPass(context.Background(), lister, rc.run, quietLog()) stats := runDiscoveryPass(context.Background(), 0, lister, rc.run, quietLog())
require.Equal(t, 1, rc.count("a")) require.Equal(t, 1, rc.count("a"))
require.Equal(t, 1, rc.count("b")) require.Equal(t, 1, rc.count("b"))
@@ -134,7 +160,7 @@ func TestDiscoveryPassListerErrorIsContained(t *testing.T) {
lister := fakeLister{err: errors.New("db down")} lister := fakeLister{err: errors.New("db down")}
rc := newCountingRunUser() rc := newCountingRunUser()
stats := runDiscoveryPass(context.Background(), lister, rc.run, quietLog()) stats := runDiscoveryPass(context.Background(), 0, lister, rc.run, quietLog())
require.Equal(t, 0, rc.total(), "no users enumerated → no passes") require.Equal(t, 0, rc.total(), "no users enumerated → no passes")
require.Equal(t, runner.Stats{}, stats) require.Equal(t, runner.Stats{}, stats)