feat(onboard): burst picks likely-good videos, summarizes with stronger model (ADR-028)
Wire the onboarding burst to its quality-aware selection and a stronger model: - main.go onboard uses OnboardBurstVideoIDs (junk-avoiding) instead of pure newest-first, bounded by MinVideoSeconds / OnboardMaxVideoSeconds. - buildBurstProcessor builds a burst-only summarizer chain led by the onboard model (burstChainModels: onboard -> standard ADR-022 chain, deduped, NDA lever intact), over the same store/cache/sink. Collapses onto the shared Processor when the onboard model is empty/equal-to-primary or config is incomplete. - summarizerEndpoint + newYouTubeSource extracted so standard and burst wiring share one definition. - Remove now-superseded NewestUnsummarizedVideoIDs: OnboardBurstVideoIDs(.,0,0) is identical pure-newest behaviour and its test covers RLS + ordering.
This commit is contained in:
+23
-10
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
+83
-16
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -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,15 +126,17 @@ func TestOnboardBurstVideoIDs(t *testing.T) {
|
||||
return id
|
||||
}
|
||||
|
||||
good1 := mk(userA, "good0000001", 5, 600) // 10m, newest known-good
|
||||
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
|
||||
_ = 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
|
||||
|
||||
// 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
|
||||
|
||||
// Known-good ranked before unknown, each newest-first within its group; the
|
||||
|
||||
Reference in New Issue
Block a user