diff --git a/cmd/tapir/scheduler.go b/cmd/tapir/scheduler.go index f7cd560..f8ed8ad 100644 --- a/cmd/tapir/scheduler.go +++ b/cmd/tapir/scheduler.go @@ -49,10 +49,12 @@ func buildUserRunner(cfg config.Config, st *store.Store, secretStore ports.Secre runner.WithAutoWindow(cfg.AutoSummarizeWindow)), nil } -// userLister enumerates every registered user. *store.Store satisfies it via -// ListAllUsers. A small local interface keeps the scheduler testable with a fake. +// userLister enumerates every registered user and reports a user's video +// connections. *store.Store satisfies it via ListAllUsers + ConnectionsForUser. +// A small local interface keeps the scheduler testable with a fake. type userLister interface { ListAllUsers(ctx context.Context) ([]store.UserIdentity, error) + ConnectionsForUser(ctx context.Context, userID string) ([]store.Connection, error) } // runDiscoveryPass runs one discovery pass for every user. runUser performs a @@ -78,6 +80,18 @@ func runDiscoveryPass( if ctx.Err() != nil { break // shutting down: stop enumerating } + // 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) + continue + } + if len(conns) == 0 { + log.Debug("scheduler: skipping user with no video connections", "user", u.UserID) + continue + } stats, err := runUser(ctx, u.UserID) total = sumStats(total, stats) if err != nil { diff --git a/cmd/tapir/scheduler_test.go b/cmd/tapir/scheduler_test.go index 97fd536..a56b88d 100644 --- a/cmd/tapir/scheduler_test.go +++ b/cmd/tapir/scheduler_test.go @@ -21,14 +21,25 @@ func quietLog() *slog.Logger { // fakeLister returns a fixed user set (or an error) for the scheduler under test. type fakeLister struct { - users []store.UserIdentity - err error + users []store.UserIdentity + err error + noConn map[string]bool // users that have NOT connected a video source } func (f fakeLister) ListAllUsers(context.Context) ([]store.UserIdentity, error) { return f.users, f.err } +// ConnectionsForUser reports a single youtube connection for every user except +// those in noConn, which return zero — the connection-less case the scheduler +// must skip instead of running (and failing to resolve a token for). +func (f fakeLister) ConnectionsForUser(_ context.Context, userID string) ([]store.Connection, error) { + if f.noConn[userID] { + return nil, nil + } + return []store.Connection{{Provider: "youtube"}}, nil +} + // countingRunUser records how many passes each user got, optionally failing for // specific users, under a mutex so it is safe across the scheduler goroutine. type countingRunUser struct { @@ -91,6 +102,21 @@ func TestDiscoveryPassRunsEveryUserOnce(t *testing.T) { require.Equal(t, 3, stats.Summarized, "stats are summed across users") } +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()) + + 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, 1, rc.count("c")) + require.Equal(t, 2, stats.Summarized, "only connected users contribute") + require.Equal(t, 0, stats.Errors, "skipping is silent — no spurious error stat") +} + func TestDiscoveryPassOneUserFailureDoesNotStopOthers(t *testing.T) { lister := fakeLister{users: usersN("a", "b", "c")} rc := newCountingRunUser("b") // user b's pass errors