diff --git a/cmd/tapir/main.go b/cmd/tapir/main.go index 3a5ff65..2ae7fd0 100644 --- a/cmd/tapir/main.go +++ b/cmd/tapir/main.go @@ -262,23 +262,36 @@ func cmdServe(ctx context.Context, log *slog.Logger) error { // One lock shared by the scheduler and connect-triggered passes (#6) so // they never fetch concurrently — the single-fetcher invariant (ADR-018). runUser := serialize(&sync.Mutex{}, rawRunUser) - // Onboarding burst (Feature 1): after the connect-triggered discovery pass, - // summarize up to OnboardSummarizeCount of the user's NEWEST unsummarized - // videos so a fresh account gets real summaries in its first session. Hard - // cap; explicit, so it bypasses the recency window — but every fetch still - // goes through globalFetchGate via the Processor. No-op when disabled - // (count 0) or queue-only (no Processor). + // Onboarding burst (Feature 1, refined by ADR-028): after the connect-triggered + // discovery pass, summarize up to OnboardSummarizeCount of the user's newest + // LIKELY-GOOD unsummarized videos so a fresh account gets a strong first + // session. Selection avoids known-junk (Shorts/over-long/livestream VODs via + // the persisted duration); the burst leads its chain with the stronger onboard + // model. Hard cap; explicit, so it bypasses the recency window — but every + // fetch still goes through globalFetchGate. No-op when disabled (count 0) or + // queue-only (no processor). + // + // burstProcessor leads with the stronger model (ADR-028); it collapses onto the + // shared Processor when the onboard model is empty/equal-to-primary or the + // engine config is incomplete. + burstProcessor := app.Processor + if burstEngine, berr := buildBurstProcessor(cfg, st); berr != nil { + return berr + } else if burstEngine != nil { + burstProcessor = &engineProcessor{engine: burstEngine, store: st} + log.Info("onboarding burst uses a stronger model", "onboard_model", cfg.OnboardSummarizerModel) + } onboard := func(ctx context.Context, userID string) { - if cfg.OnboardSummarizeCount <= 0 || app.Processor == nil { + if cfg.OnboardSummarizeCount <= 0 || burstProcessor == nil { return } - ids, err := st.NewestUnsummarizedVideoIDs(ctx, userID, cfg.OnboardSummarizeCount) + ids, err := st.OnboardBurstVideoIDs(ctx, userID, cfg.OnboardSummarizeCount, cfg.MinVideoSeconds, cfg.OnboardMaxVideoSeconds) if err != nil { - log.Warn("onboarding: list newest unsummarized", "user", userID, "err", err) + log.Warn("onboarding: list burst candidates", "user", userID, "err", err) return } for _, id := range ids { - if err := app.Processor.ProcessVideo(ctx, userID, id); err != nil { + if err := burstProcessor.ProcessVideo(ctx, userID, id); err != nil { log.Warn("onboarding: summarize", "user", userID, "video", id, "err", err) } } diff --git a/cmd/tapir/processor.go b/cmd/tapir/processor.go index 2c82ddf..bb15676 100644 --- a/cmd/tapir/processor.go +++ b/cmd/tapir/processor.go @@ -52,13 +52,7 @@ func (f videoFetcher) FetchVideo(ctx context.Context, userID, videoID string) (d // alias, not a second client config. Empty model entries are skipped, so a // client deployment can set the cloud fallback empty to keep content local. func buildSummarizer(cfg config.Config) *summarizer.Summarizer { - mk := func(model string) summarizer.Endpoint { - return summarizer.Endpoint{ - Client: llm.New(cfg.GatewayURL, cfg.GatewayKey, model, cfg.SummarizerTimeout, llm.WithMaxTokens(cfg.SummaryMaxTokens)), - Provider: providerOf(model), - Model: model, - } - } + mk := summarizerEndpoint(cfg) eps := []summarizer.Endpoint{mk(cfg.SummarizerModel)} if cfg.FallbackModel != "" && cfg.FallbackModel != cfg.SummarizerModel { eps = append(eps, mk(cfg.FallbackModel)) @@ -69,6 +63,55 @@ func buildSummarizer(cfg config.Config) *summarizer.Summarizer { return summarizer.NewChain(eps, cfg.MaxTranscriptChars) } +// summarizerEndpoint returns a constructor for a chain endpoint over the one +// LiteLLM gateway, varying only the model alias (the gateway fronts both +// llama-swap and berget). Shared by the standard and burst chains. +func summarizerEndpoint(cfg config.Config) func(model string) summarizer.Endpoint { + return func(model string) summarizer.Endpoint { + return summarizer.Endpoint{ + Client: llm.New(cfg.GatewayURL, cfg.GatewayKey, model, cfg.SummarizerTimeout, llm.WithMaxTokens(cfg.SummaryMaxTokens)), + Provider: providerOf(model), + Model: model, + } + } +} + +// burstChainModels is the ordered, deduped model list for the onboarding burst +// (ADR-028): the stronger onboard model leads, then the standard ADR-022 chain +// (primary → local fallback → cloud) follows as resilience. Empty entries are +// dropped and duplicates collapsed, so the NDA lever (empty cloud fallback) keeps +// the burst chain fully local exactly as the standard chain does. +func burstChainModels(cfg config.Config) []string { + var models []string + add := func(m string) { + if m == "" { + return + } + for _, e := range models { + if e == m { + return + } + } + models = append(models, m) + } + add(cfg.OnboardSummarizerModel) + add(cfg.SummarizerModel) + add(cfg.FallbackModel) + add(cfg.CloudFallbackModel) + return models +} + +// buildBurstSummarizer builds the onboarding-burst summarizer chain (ADR-028): +// the onboard model first, then the standard chain as fallback, deduped. +func buildBurstSummarizer(cfg config.Config) *summarizer.Summarizer { + mk := summarizerEndpoint(cfg) + var eps []summarizer.Endpoint + for _, m := range burstChainModels(cfg) { + eps = append(eps, mk(m)) + } + return summarizer.NewChain(eps, cfg.MaxTranscriptChars) +} + // chatModels is the ordered, local-first set of models offered in the chat // switcher (ADR-027), reusing the ADR-022 chain: primary → local fallback → // cloud. Empty entries are dropped and duplicates collapsed, so a client/NDA @@ -126,15 +169,7 @@ func buildProcessor(cfg config.Config, st *store.Store) (*usecase.Engine, error) return nil, nil } - secretStore := secrets.NewFileStore(cfg.SecretsFile) - src := youtube.New(youtube.Config{ - ClientID: cfg.YTClientID, - ClientSecret: cfg.YTClientSecret, - TokenSecretRef: cfg.YTTokenRef, - PreferredLanguages: []string{"en"}, - MinVideoSeconds: cfg.MinVideoSeconds, - }, secretStore) - + src := newYouTubeSource(cfg, secrets.NewFileStore(cfg.SecretsFile)) sum := buildSummarizer(cfg) // The store is both the summary sink and the shared transcript cache (ADR-021): @@ -145,6 +180,38 @@ func buildProcessor(cfg config.Config, st *store.Store) (*usecase.Engine, error) return eng, nil } +// newYouTubeSource builds the captions-first VideoSource shared by the standard +// and burst processors — same per-process YouTube credentials and ADR-023 Shorts +// filter; only the summarizer chain differs between them. +func newYouTubeSource(cfg config.Config, secretStore ports.SecretStore) ports.VideoSource { + return youtube.New(youtube.Config{ + ClientID: cfg.YTClientID, + ClientSecret: cfg.YTClientSecret, + TokenSecretRef: cfg.YTTokenRef, + PreferredLanguages: []string{"en"}, + MinVideoSeconds: cfg.MinVideoSeconds, + }, secretStore) +} + +// buildBurstProcessor wires a processor whose summarizer leads with the stronger +// onboard model (ADR-028), used only by the connect-time burst over the SAME +// store / transcript cache / sink — a wiring choice; the engine and ports are +// unchanged. Returns (nil, nil) — the collapse lever — when the onboard model is +// empty or equal to the primary (the burst then reuses the shared processor), or +// when the engine config is incomplete (queue-only, same as buildProcessor). +func buildBurstProcessor(cfg config.Config, st *store.Store) (*usecase.Engine, error) { + if cfg.OnboardSummarizerModel == "" || cfg.OnboardSummarizerModel == cfg.SummarizerModel { + return nil, nil + } + if cfg.GatewayURL == "" || cfg.YTClientID == "" || cfg.YTClientSecret == "" || cfg.SecretsFile == "" { + return nil, nil + } + src := newYouTubeSource(cfg, secrets.NewFileStore(cfg.SecretsFile)) + eng := usecase.NewEngine(src, buildBurstSummarizer(cfg), st) + eng.Transcripts = st + return eng, nil +} + // engineProcessor adapts the engine (which works in terms of a domain.Video) to // the web.Processor port (which works in terms of a stored video id): it loads the // video row, runs the engine, and — on a produced summary — clears the manual diff --git a/cmd/tapir/processor_test.go b/cmd/tapir/processor_test.go index ccb2c6b..3f93a14 100644 --- a/cmd/tapir/processor_test.go +++ b/cmd/tapir/processor_test.go @@ -45,3 +45,65 @@ func TestBuildProcessorNilOnIncompleteConfig(t *testing.T) { }) } } + +// TestBurstChainModelsLeadsWithOnboardModel: the onboarding burst chain (ADR-028) +// leads with the stronger onboard model, then falls back through the standard +// ADR-022 chain (primary -> local fallback -> cloud), deduped. +func TestBurstChainModelsLeadsWithOnboardModel(t *testing.T) { + got := burstChainModels(config.Config{ + OnboardSummarizerModel: "iguana/gemma4-26b", + SummarizerModel: "koala/phi4-mini", + FallbackModel: "iguana/gemma4-26b", // also the onboard model -> dedup + CloudFallbackModel: "berget/mistral-small", + }) + want := []string{"iguana/gemma4-26b", "koala/phi4-mini", "berget/mistral-small"} + if len(got) != len(want) { + t.Fatalf("burstChainModels = %v, want %v", got, want) + } + for i := range want { + if got[i] != want[i] { + t.Fatalf("burstChainModels = %v, want %v", got, want) + } + } +} + +// TestBurstChainModelsCloudAbsentWhenDisabled: the NDA lever holds for the burst +// too — empty cloud fallback keeps the burst chain fully local. +func TestBurstChainModelsCloudAbsentWhenDisabled(t *testing.T) { + got := burstChainModels(config.Config{ + OnboardSummarizerModel: "iguana/gemma4-26b", + SummarizerModel: "koala/phi4-mini", + CloudFallbackModel: "", + }) + for _, m := range got { + if m == "" || m == "berget/mistral-small" { + t.Fatalf("cloud model leaked into burst chain: %v", got) + } + } +} + +// TestBuildBurstProcessorNilWhenCollapsed: an empty or primary-equal onboard model +// collapses the burst onto the shared processor (buildBurstProcessor returns nil). +func TestBuildBurstProcessorNilWhenCollapsed(t *testing.T) { + base := config.Config{ + GatewayURL: "http://gw/v1", + YTClientID: "id", + YTClientSecret: "secret", + SecretsFile: "/tmp/secrets.json", + SummarizerModel: "koala/phi4-mini", + } + t.Run("empty onboard model", func(t *testing.T) { + base.OnboardSummarizerModel = "" + eng, err := buildBurstProcessor(base, nil) + if err != nil || eng != nil { + t.Fatalf("buildBurstProcessor = (%v, %v), want (nil, nil)", eng, err) + } + }) + t.Run("onboard model equals primary", func(t *testing.T) { + base.OnboardSummarizerModel = "koala/phi4-mini" + eng, err := buildBurstProcessor(base, nil) + if err != nil || eng != nil { + t.Fatalf("buildBurstProcessor = (%v, %v), want (nil, nil)", eng, err) + } + }) +} diff --git a/internal/adapters/store/videos.go b/internal/adapters/store/videos.go index fa73768..7db0fe6 100644 --- a/internal/adapters/store/videos.go +++ b/internal/adapters/store/videos.go @@ -86,44 +86,6 @@ func nullDuration(seconds int) *int { return &seconds } -// NewestUnsummarizedVideoIDs returns up to limit of the user's videos that have -// no summary yet, newest first (published_at DESC, NULLS LAST). It caps the -// connect-time onboarding burst (Feature 1) at a fixed count: the caller marks -// these for summarization through the shared rate gate. RLS-scoped via withUser, -// so it only ever sees the requesting user's rows. limit <= 0 returns nil. -func (s *Store) NewestUnsummarizedVideoIDs(ctx context.Context, userID string, limit int) ([]string, error) { - if limit <= 0 { - return nil, nil - } - var ids []string - if err := s.withUser(ctx, userID, func(tx pgx.Tx) error { - rows, err := tx.Query(ctx, - `SELECT v.id - FROM videos v - WHERE v.user_id = $1 - AND NOT EXISTS ( - SELECT 1 FROM summaries su - WHERE su.user_id = v.user_id AND su.video_id = v.id) - ORDER BY v.published_at DESC NULLS LAST, v.seen_at DESC - LIMIT $2`, userID, limit) - if err != nil { - return fmt.Errorf("store: newest unsummarized: %w", err) - } - defer rows.Close() - for rows.Next() { - var id string - if err := rows.Scan(&id); err != nil { - return fmt.Errorf("store: scan newest unsummarized: %w", err) - } - ids = append(ids, id) - } - return rows.Err() - }); err != nil { - return nil, err - } - return ids, nil -} - // OnboardBurstVideoIDs returns up to limit of the user's unsummarized videos for // the connect-time onboarding burst (ADR-028), newest-first but quality-aware: a // video is excluded when its duration is KNOWN and outside [minSeconds, maxSeconds] diff --git a/internal/adapters/store/videos_test.go b/internal/adapters/store/videos_test.go index bb396be..f0c453b 100644 --- a/internal/adapters/store/videos_test.go +++ b/internal/adapters/store/videos_test.go @@ -112,38 +112,6 @@ func TestUpsertVideo_PerUserIsolation(t *testing.T) { require.NotEqual(t, idA, idB, "same provider video for two users must be two distinct rows") } -func TestNewestUnsummarizedVideoIDs(t *testing.T) { - ctx := context.Background() - s := newStore(t) - resetDB(t, rawPool(t)) - - mk := func(user, pid string, day int) string { - v := ytVideo(user, pid, pid) - v.PublishedAt = time.Date(2026, 6, day, 12, 0, 0, 0, time.UTC) - id, err := s.UpsertVideo(ctx, v) - require.NoError(t, err) - return id - } - - _ = mk(userA, "a1vid000001", 1) - id2 := mk(userA, "a2vid000002", 2) - id3 := mk(userA, "a3vid000003", 3) - id4 := mk(userA, "a4vid000004", 4) - mk(userB, "b1vid000009", 9) // userB's newest — must never leak via RLS - - // The newest (v4) is summarized, so it's excluded from "unsummarized". - require.NoError(t, s.Deliver(ctx, summary(userA, id4, "done"))) - - // Cap 2, newest-first unsummarized: v3 then v2 (v4 excluded; userB excluded). - got, err := s.NewestUnsummarizedVideoIDs(ctx, userA, 2) - require.NoError(t, err) - require.Equal(t, []string{id3, id2}, got) - - none, err := s.NewestUnsummarizedVideoIDs(ctx, userA, 0) - require.NoError(t, err) - require.Empty(t, none, "limit 0 returns nothing") -} - func TestOnboardBurstVideoIDs(t *testing.T) { ctx := context.Background() s := newStore(t) @@ -158,14 +126,16 @@ func TestOnboardBurstVideoIDs(t *testing.T) { return id } - good1 := mk(userA, "good0000001", 5, 600) // 10m, newest known-good - tooLong := mk(userA, "toolong0001", 4, 20000) // > maxSeconds -> dropped - _ = tooLong - tooShort := mk(userA, "tooshort001", 3, 30) // < minSeconds -> dropped - _ = tooShort - unknown := mk(userA, "unknown0001", 2, 0) // NULL duration -> kept, ranked last - good2 := mk(userA, "good0000002", 1, 800) // known-good but oldest - mk(userB, "bvid0000009", 9, 600) // userB -> must not leak via RLS + summarized := mk(userA, "summ0000001", 6, 600) // newest known-good, but already summarized + good1 := mk(userA, "good0000001", 5, 600) // 10m, newest UNsummarized known-good + tooLong := mk(userA, "toolong0001", 4, 20000) // > maxSeconds -> dropped + tooShort := mk(userA, "tooshort001", 3, 30) // < minSeconds -> dropped + unknown := mk(userA, "unknown0001", 2, 0) // NULL duration -> kept, ranked last + good2 := mk(userA, "good0000002", 1, 800) // known-good but oldest + mk(userB, "bvid0000009", 9, 600) // userB -> must not leak via RLS + + // The newest video is summarized, so it is excluded from the burst. + require.NoError(t, s.Deliver(ctx, summary(userA, summarized, "done"))) const minSec, maxSec = 60, 14400