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)
|
}, 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,15 +73,17 @@ func runDiscoveryPass(
|
|||||||
return runner.Stats{}
|
return runner.Stats{}
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Info("scheduler: starting discovery pass", "users", len(users))
|
// Keep only users with a video connection. A pass for a connectionless user
|
||||||
var total runner.Stats
|
// (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 {
|
for _, u := range users {
|
||||||
if ctx.Err() != nil {
|
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)
|
conns, err := lister.ConnectionsForUser(ctx, u.UserID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Warn("scheduler: list connections failed", "user", u.UserID, "err", err)
|
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)
|
log.Debug("scheduler: skipping user with no video connections", "user", u.UserID)
|
||||||
continue
|
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)
|
stats, err := runUser(ctx, u.UserID)
|
||||||
total = sumStats(total, stats)
|
total = sumStats(total, stats)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -116,7 +140,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 +150,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 {
|
||||||
|
|||||||
@@ -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,46 @@ 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"))
|
||||||
|
}
|
||||||
|
|
||||||
|
// 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) {
|
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 +162,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 +175,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)
|
||||||
|
|||||||
Reference in New Issue
Block a user