Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1aa8a97f95 | ||
|
|
cc69a912f4 | ||
|
|
f4a0544903 |
+54
-8
@@ -36,6 +36,11 @@ func buildUserRunner(cfg config.Config, st *store.Store, secretStore ports.Secre
|
||||
}, secretStore)
|
||||
|
||||
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,
|
||||
runner.WithBackoff(cfg.FetchBackoff),
|
||||
@@ -57,6 +62,7 @@ type userLister interface {
|
||||
// isolation). Returns the stats summed across users.
|
||||
func runDiscoveryPass(
|
||||
ctx context.Context,
|
||||
pass int,
|
||||
lister userLister,
|
||||
runUser func(context.Context, string) (runner.Stats, error),
|
||||
log *slog.Logger,
|
||||
@@ -67,15 +73,17 @@ func runDiscoveryPass(
|
||||
return runner.Stats{}
|
||||
}
|
||||
|
||||
log.Info("scheduler: starting discovery pass", "users", len(users))
|
||||
var total runner.Stats
|
||||
// Keep only users with a video connection. A pass for a connectionless user
|
||||
// (e.g. a stale Dex-era orphan identity) only tries to resolve a token that
|
||||
// was never minted, logging a spurious "ref not found" every tick. Filtering
|
||||
// here — BEFORE rotation — also keeps fairness honest: rotation is over the
|
||||
// users that actually consume the caption budget, so a dead identity can't eat
|
||||
// a rotation slot and skew the lead share.
|
||||
var connected []store.UserIdentity
|
||||
for _, u := range users {
|
||||
if ctx.Err() != nil {
|
||||
break // shutting down: stop enumerating
|
||||
return runner.Stats{} // shutting down
|
||||
}
|
||||
// Skip users with no video connection. A discovery pass for them only
|
||||
// attempts to resolve a token that was never minted, logging a spurious
|
||||
// "ref not found" every tick (e.g. stale Dex-era orphan identities).
|
||||
conns, err := lister.ConnectionsForUser(ctx, u.UserID)
|
||||
if err != nil {
|
||||
log.Warn("scheduler: list connections failed", "user", u.UserID, "err", err)
|
||||
@@ -85,6 +93,22 @@ func runDiscoveryPass(
|
||||
log.Debug("scheduler: skipping user with no video connections", "user", u.UserID)
|
||||
continue
|
||||
}
|
||||
connected = append(connected, u)
|
||||
}
|
||||
|
||||
// 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 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 over the connected set gives each real user the lead in turn.
|
||||
connected = rotateUsers(connected, pass)
|
||||
|
||||
log.Info("scheduler: starting discovery pass", "users", len(connected))
|
||||
var total runner.Stats
|
||||
for _, u := range connected {
|
||||
if ctx.Err() != nil {
|
||||
break // shutting down: stop enumerating
|
||||
}
|
||||
stats, err := runUser(ctx, u.UserID)
|
||||
total = sumStats(total, stats)
|
||||
if err != nil {
|
||||
@@ -116,7 +140,8 @@ func runScheduler(
|
||||
return // disabled
|
||||
}
|
||||
|
||||
runDiscoveryPass(ctx, lister, runUser, log)
|
||||
pass := 0
|
||||
runDiscoveryPass(ctx, pass, lister, runUser, log)
|
||||
|
||||
ticker := time.NewTicker(interval)
|
||||
defer ticker.Stop()
|
||||
@@ -125,11 +150,32 @@ func runScheduler(
|
||||
case <-ctx.Done():
|
||||
return
|
||||
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
|
||||
// per-tick aggregate across all users.
|
||||
func sumStats(a, b runner.Stats) runner.Stats {
|
||||
|
||||
@@ -45,6 +45,7 @@ func (f fakeLister) ConnectionsForUser(_ context.Context, userID string) ([]stor
|
||||
type countingRunUser struct {
|
||||
mu sync.Mutex
|
||||
calls map[string]int
|
||||
order []string // userIDs in the order they were run, across all passes
|
||||
failFor map[string]bool
|
||||
}
|
||||
|
||||
@@ -60,12 +61,19 @@ func (c *countingRunUser) run(_ context.Context, userID string) (runner.Stats, e
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
c.calls[userID]++
|
||||
c.order = append(c.order, userID)
|
||||
if c.failFor[userID] {
|
||||
return runner.Stats{Errors: 1}, errors.New("boom")
|
||||
}
|
||||
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 {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
@@ -94,7 +102,7 @@ func TestDiscoveryPassRunsEveryUserOnce(t *testing.T) {
|
||||
lister := fakeLister{users: usersN("a", "b", "c")}
|
||||
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("b"))
|
||||
@@ -102,13 +110,46 @@ func TestDiscoveryPassRunsEveryUserOnce(t *testing.T) {
|
||||
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"))
|
||||
}
|
||||
|
||||
// A connectionless orphan must not consume a rotation slot: rotation is over the
|
||||
// connected users only, so two real users alternate the lead 50/50 even with a
|
||||
// dead identity listed between them.
|
||||
func TestDiscoveryPassRotationIgnoresConnectionlessUsers(t *testing.T) {
|
||||
lister := fakeLister{users: usersN("a", "orphan", "c"), noConn: map[string]bool{"orphan": true}}
|
||||
rc := newCountingRunUser()
|
||||
|
||||
runDiscoveryPass(context.Background(), 0, lister, rc.run, quietLog())
|
||||
runDiscoveryPass(context.Background(), 1, lister, rc.run, quietLog())
|
||||
|
||||
require.Equal(t, []string{"a", "c", "c", "a"}, rc.runOrder(),
|
||||
"only connected users rotate; the orphan never runs and never holds a slot")
|
||||
require.Equal(t, 0, rc.count("orphan"))
|
||||
}
|
||||
|
||||
func TestDiscoveryPassSkipsUsersWithoutConnections(t *testing.T) {
|
||||
// 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.
|
||||
lister := fakeLister{users: usersN("a", "b", "c"), noConn: map[string]bool{"b": true}}
|
||||
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, 0, rc.count("b"), "a user with no connection must be skipped, not run")
|
||||
@@ -121,7 +162,7 @@ func TestDiscoveryPassOneUserFailureDoesNotStopOthers(t *testing.T) {
|
||||
lister := fakeLister{users: usersN("a", "b", "c")}
|
||||
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("b"))
|
||||
@@ -134,7 +175,7 @@ func TestDiscoveryPassListerErrorIsContained(t *testing.T) {
|
||||
lister := fakeLister{err: errors.New("db down")}
|
||||
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, runner.Stats{}, stats)
|
||||
|
||||
Reference in New Issue
Block a user