From e9b5a3f3e7a3952264b3925e4cacae34f801950d Mon Sep 17 00:00:00 2001 From: Mathias Date: Wed, 10 Jun 2026 08:36:16 +0200 Subject: [PATCH] =?UTF-8?q?feat(summarizer):=20resilient=20endpoint=20chai?= =?UTF-8?q?n=20with=20local=E2=86=92cloud=20fallback=20(ADR-022)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The first friendly-pilot live run produced zero summaries: koala/phi4-mini hit three silent failure modes — 8k context overflow on long transcripts (HTTP 400), intermittent malformed JSON (highlights as a bare string), and no fallback wired at all (summarizer.New(primary, nil)). Keep phi4-mini as the fast primary and add resilience around it: - Ordered endpoint chain (summarizer.NewChain): phi4-mini → koala/phi4-14b (local) → berget/mistral-small (worst-case external). All reached through the one LiteLLM gateway by alias. - A parse failure now advances the chain like a transport error — the old Primary→Fallback shape returned the parse error without trying anyone else. - Tolerant parse: highlights/takeaways coerce string→[]string, absorbing the common small-model quirk without spending a fallback round-trip. - Transcript truncation (TAPIR_MAX_TRANSCRIPT_CHARS=18000) prevents the overflow rather than recovering from it; validated to fit phi4-mini's 8k window. - Bounded completion budget (TAPIR_SUMMARY_MAX_TOKENS=1500) — the old 8192 budget itself contributed to the overflow. Local-first guarantee preserved by ordering: external endpoint is tried only after every local one fails. TAPIR_CLOUD_FALLBACK_MODEL="" disables it entirely for client/NDA deployments. Co-Authored-By: Claude Opus 4.8 (1M context) --- DECISIONS.md | 52 +++++ cmd/tapir/processor.go | 43 +++- cmd/tapir/scheduler.go | 9 +- docs/homelab-integration.md | 12 ++ docs/use-cases/ai_routing.feature | 14 +- internal/adapters/llm/client.go | 24 ++- internal/adapters/llm/client_test.go | 21 ++ internal/adapters/summarizer/summarizer.go | 184 ++++++++++++++---- .../adapters/summarizer/summarizer_test.go | 98 +++++++++- internal/config/config.go | 50 ++++- internal/config/config_test.go | 58 ++++++ test/acceptance/scenario_coverage_test.go | 9 +- 12 files changed, 509 insertions(+), 65 deletions(-) diff --git a/DECISIONS.md b/DECISIONS.md index 4eeda4b..731ad58 100644 --- a/DECISIONS.md +++ b/DECISIONS.md @@ -803,6 +803,58 @@ never exposed in the UI. --- +## ADR-022 — Summarizer is a resilient endpoint chain, not a single model + +**Status:** Accepted (2026-06-10). **Extends ADR-004** (the copied `llm` Primary→Fallback +routing). Triggered by the first friendly-pilot live run, where a connected user got **zero** +summaries after 12h. + +**Context.** Stage-0 ran a single summarizer model (`koala/phi4-mini`) with no fallback wired +(`summarizer.New(primary, nil)`). The live run exposed three independent failure modes, each of +which silently produced no summary: + +1. **Context overflow.** `phi4-mini` has an 8k context. Real transcripts (one was 11,602 tokens) + exceed it and the gateway returns HTTP 400 — and the request also sent `max_tokens=8192`, so + even a short transcript plus the completion budget could overflow the window. +2. **Malformed model output.** `phi4-mini` intermittently emits `highlights` as a bare string + instead of an array, producing `cannot unmarshal string into []string`. The old code returned + the parse error **without** trying any other model — a 200-with-bad-JSON short-circuited. +3. **No fallback existed at all** — any primary failure was terminal for that video. + +`phi4-mini` is kept as primary deliberately: it is fast and, on transcripts that fit, correct. +The fix is resilience around it, not replacing it. + +**Decision.** +1. **Ordered endpoint chain (`summarizer.NewChain`).** Endpoints are tried in order; the first to + return a *parseable* summary wins. Default chain: + `koala/phi4-mini` (primary, local) → `koala/phi4-14b` (fallback, local) → + `berget/mistral-small` (worst-case, external). All three are reached through the **one** LiteLLM + gateway by alias — the gateway already fronts both llama-swap and berget — so a fallback is a + different alias, not a second client config. +2. **A parse failure advances the chain, same as a transport error.** "Reliably summarized" means + *parseable summary returned*, not *HTTP 200*. This is the behaviour the old Primary→Fallback + shape missed. +3. **Tolerant parse.** `highlights`/`takeaways` coerce from a bare string (or a mixed scalar + array) to `[]string`, so the most common small-model quirk is absorbed **without** spending a + fallback round-trip — keeping the fast path fast. +4. **Transcript truncation (`TAPIR_MAX_TRANSCRIPT_CHARS`, default 18000).** Input is bounded + up-front to fit a small-context primary, so overflow is prevented rather than recovered-from. +5. **Bounded completion budget (`TAPIR_SUMMARY_MAX_TOKENS`, default 1500).** A summary needs few + hundred tokens; the old 8192 budget itself contributed to 8k-window overflow. + +**Local-first guarantee preserved.** The chain ordering *is* the guarantee: locals are tried +first, so content reaches the external endpoint only after every local endpoint has failed. +`TAPIR_CLOUD_FALLBACK_MODEL=""` removes the external endpoint entirely — the lever a +**client/NDA deployment** pulls so content never leaves the local stack. With no external endpoint +configured the `ai_routing.feature` "content only local" scenarios hold unchanged. + +**Reversibility.** Pure wiring + config. Setting `TAPIR_FALLBACK_MODEL` and +`TAPIR_CLOUD_FALLBACK_MODEL` empty collapses the chain back to single-primary behaviour; the +tolerant parse and truncation are strict supersets of the old behaviour (a previously-parseable +reply still parses; a transcript within budget is unchanged). + +--- + ## Rejected alternatives Approaches considered during the 2026-06-02 planning + grill session and **deliberately not diff --git a/cmd/tapir/processor.go b/cmd/tapir/processor.go index 5a2ba5a..84a41f7 100644 --- a/cmd/tapir/processor.go +++ b/cmd/tapir/processor.go @@ -3,6 +3,7 @@ package main import ( "context" "fmt" + "strings" "gitea.d-ma.be/mathias/tapir/internal/adapters/llm" "gitea.d-ma.be/mathias/tapir/internal/adapters/secrets" @@ -42,6 +43,40 @@ func (f videoFetcher) FetchVideo(ctx context.Context, userID, videoID string) (d // queue-only fallback: the web UI keeps working (the button just queues) and // `tapir run` reports the gap via its own ValidateForRun. Missing engine config // is never an error here. +// buildSummarizer wires the summarization endpoint chain (ADR-022) shared by the +// web "Summarize now" path and the scheduler's per-user runners. The chain is: +// primary (local, fast) → local fallback → cloud fallback (worst case). Each +// endpoint reaches the same LiteLLM gateway with a different model alias — the +// gateway fronts both llama-swap and berget — so a fallback is just a different +// 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, + } + } + eps := []summarizer.Endpoint{mk(cfg.SummarizerModel)} + if cfg.FallbackModel != "" && cfg.FallbackModel != cfg.SummarizerModel { + eps = append(eps, mk(cfg.FallbackModel)) + } + if cfg.CloudFallbackModel != "" && cfg.CloudFallbackModel != cfg.SummarizerModel { + eps = append(eps, mk(cfg.CloudFallbackModel)) + } + return summarizer.NewChain(eps, cfg.MaxTranscriptChars) +} + +// providerOf maps a model alias to the domain AIProvider recorded on summaries. +// A "berget/" alias is an external provider; everything else is the local stack. +func providerOf(model string) string { + if strings.HasPrefix(model, "berget/") { + return "berget" + } + return "local" +} + func buildProcessor(cfg config.Config, st *store.Store) (*usecase.Engine, error) { if cfg.GatewayURL == "" || cfg.YTClientID == "" || cfg.YTClientSecret == "" || cfg.SecretsFile == "" { return nil, nil @@ -55,13 +90,7 @@ func buildProcessor(cfg config.Config, st *store.Store) (*usecase.Engine, error) PreferredLanguages: []string{"en"}, }, secretStore) - // Local Primary only; no BYO fallback for the demo (fallback nil). - primary := summarizer.Endpoint{ - Client: llm.New(cfg.GatewayURL, cfg.GatewayKey, cfg.SummarizerModel, cfg.SummarizerTimeout), - Provider: "local", - Model: cfg.SummarizerModel, - } - sum := summarizer.New(primary, nil) + sum := buildSummarizer(cfg) // The store is both the summary sink and the shared transcript cache (ADR-021): // the engine reads stored transcripts before any caption fetch and writes diff --git a/cmd/tapir/scheduler.go b/cmd/tapir/scheduler.go index f8ed8ad..8b02c25 100644 --- a/cmd/tapir/scheduler.go +++ b/cmd/tapir/scheduler.go @@ -6,9 +6,7 @@ import ( "log/slog" "time" - "gitea.d-ma.be/mathias/tapir/internal/adapters/llm" "gitea.d-ma.be/mathias/tapir/internal/adapters/store" - "gitea.d-ma.be/mathias/tapir/internal/adapters/summarizer" "gitea.d-ma.be/mathias/tapir/internal/adapters/youtube" "gitea.d-ma.be/mathias/tapir/internal/config" "gitea.d-ma.be/mathias/tapir/internal/ports" @@ -37,12 +35,7 @@ func buildUserRunner(cfg config.Config, st *store.Store, secretStore ports.Secre PreferredLanguages: []string{"en"}, }, secretStore) - primary := summarizer.Endpoint{ - Client: llm.New(cfg.GatewayURL, cfg.GatewayKey, cfg.SummarizerModel, cfg.SummarizerTimeout), - Provider: "local", - Model: cfg.SummarizerModel, - } - engine := usecase.NewEngine(src, summarizer.New(primary, nil), st) + engine := usecase.NewEngine(src, buildSummarizer(cfg), st) return runner.New(src, st, engine, userID, log, runner.WithBackoff(cfg.FetchBackoff), diff --git a/docs/homelab-integration.md b/docs/homelab-integration.md index 4901f20..e2f8b0d 100644 --- a/docs/homelab-integration.md +++ b/docs/homelab-integration.md @@ -27,6 +27,18 @@ it** — endpoints and aliases drift, and this file is a snapshot (2026-06-06), `iguana/deepseek-r1-14b`) is preferred for summary quality if its latency/output is acceptable. The `max_tokens` fix below means thinking models no longer return empty content, so they are now viable choices, not blocked ones. Do not assume a coder alias is right for prose. +- **Summarizer fallback chain (ADR-022).** The primary alias is the *first* of an ordered chain; + on failure or unparseable output the summarizer advances to the next model. All reached through + the same gateway by alias. + - `TAPIR_FALLBACK_MODEL` — local fallback. **Default `koala/phi4-14b`.** Empty disables it. + - `TAPIR_CLOUD_FALLBACK_MODEL` — worst-case EXTERNAL fallback. **Default `berget/mistral-small`.** + **Set this empty (`""`) for any client/NDA deployment** so content never leaves the local + stack — the chain then contains only local endpoints. + - `TAPIR_SUMMARY_MAX_TOKENS` — per-summary completion budget. **Default `1500`.** Small on + purpose: with the old 8192 budget, prompt + completion overflowed `phi4-mini`'s 8k window. + - `TAPIR_MAX_TRANSCRIPT_CHARS` — transcript truncation budget sent to the model. **Default + `18000`** (~fits an 8k-context model). `0` disables truncation. Prevents the context-overflow + HTTP 400 a long transcript caused on `phi4-mini`. - **Thinking models need an explicit `max_tokens`.** qwen3 / deepseek-r1 spend the budget on reasoning and return **empty content** if `max_tokens` is too low (or unset). The summarizer's parser treats an empty summary as an error for exactly this reason. **Done (2026-06-02, Worker F):** diff --git a/docs/use-cases/ai_routing.feature b/docs/use-cases/ai_routing.feature index 7ad1360..8848f24 100644 --- a/docs/use-cases/ai_routing.feature +++ b/docs/use-cases/ai_routing.feature @@ -31,5 +31,17 @@ Feature: Local-first AI with optional BYO fallback When any transcript is summarized Then my content is only ever sent to the local AI stack - # "Reliably" is operationalized as: Primary returned without error within timeout. + Scenario: A model returns unparseable output and the next endpoint succeeds + Given the local AI stack is available + But the primary model returns output that cannot be parsed into a summary + And a fallback model is configured + When a transcript is summarized + Then Tapir falls back to the next model in the chain + And the summary records fallback_used as true + + # "Reliably" is operationalized as: an endpoint returned a PARSEABLE summary + # within timeout. A 200 with malformed JSON (or highlights emitted as a bare + # string) counts as a failure and advances the chain (ADR-022). Endpoints are + # tried in order, locals first, so the external worst-case model only ever sees + # content after every local endpoint has failed. # Quality scoring may be added later without changing these scenarios. diff --git a/internal/adapters/llm/client.go b/internal/adapters/llm/client.go index 8e40af7..430876e 100644 --- a/internal/adapters/llm/client.go +++ b/internal/adapters/llm/client.go @@ -34,15 +34,35 @@ type Client struct { httpClient *http.Client } +// Option configures a Client at construction. Variadic so the existing 4-arg +// call sites stay valid as new knobs are added. +type Option func(*Client) + +// WithMaxTokens overrides the per-request completion budget. The summarizer uses +// this to cap completion for small-context models (e.g. koala/phi4-mini, 8k): +// with the default 8192 budget, prompt + max_tokens overflows an 8k context and +// the gateway returns HTTP 400. A non-positive n is ignored (keeps the default). +func WithMaxTokens(n int) Option { + return func(c *Client) { + if n > 0 { + c.maxTokens = n + } + } +} + // New constructs a Client. -func New(baseURL, apiKey, model string, timeout time.Duration) *Client { - return &Client{ +func New(baseURL, apiKey, model string, timeout time.Duration, opts ...Option) *Client { + c := &Client{ baseURL: strings.TrimRight(baseURL, "/"), apiKey: apiKey, model: model, maxTokens: defaultMaxTokens, httpClient: &http.Client{Timeout: timeout}, } + for _, opt := range opts { + opt(c) + } + return c } type chatRequest struct { diff --git a/internal/adapters/llm/client_test.go b/internal/adapters/llm/client_test.go index 1f0dc91..1ec7b74 100644 --- a/internal/adapters/llm/client_test.go +++ b/internal/adapters/llm/client_test.go @@ -64,6 +64,27 @@ func TestClient_SendsMaxTokens(t *testing.T) { } } +// TestClient_WithMaxTokens overrides the completion budget — the summarizer caps +// it small so prompt + max_tokens fits a small-context model's window (8k). +func TestClient_WithMaxTokens(t *testing.T) { + var body chatRequest + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + _ = json.NewDecoder(r.Body).Decode(&body) + _ = json.NewEncoder(w).Encode(map[string]any{ + "choices": []map[string]any{{"message": map[string]any{"content": "ok"}}}, + }) + })) + defer srv.Close() + + c := New(srv.URL, "", "test-model", 10*time.Second, WithMaxTokens(1500)) + if _, err := c.Complete(context.Background(), "sys", "user"); err != nil { + t.Fatalf("Complete: %v", err) + } + if body.MaxTokens != 1500 { + t.Errorf("max_tokens = %d, want 1500", body.MaxTokens) + } +} + func TestClient_ReturnsErrorOnNon200(t *testing.T) { srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { http.Error(w, "overloaded", http.StatusServiceUnavailable) diff --git a/internal/adapters/summarizer/summarizer.go b/internal/adapters/summarizer/summarizer.go index c387370..0312ee7 100644 --- a/internal/adapters/summarizer/summarizer.go +++ b/internal/adapters/summarizer/summarizer.go @@ -1,17 +1,21 @@ // Package summarizer implements ports.Summarizer backed by the copied llm -// package's local-Primary -> BYO-Fallback routing (ADR-004). It is the only -// place content ever leaves the engine toward an AI model, so it is also the -// enforcement point for the local-first guarantee in -// docs/use-cases/ai_routing.feature: a user with no BYO provider configured has -// their content sent to the local stack and nowhere else. +// package's routing (ADR-004, extended by ADR-022). It is the only place content +// ever leaves the engine toward an AI model, so it is also the enforcement point +// for the local-first guarantee in docs/use-cases/ai_routing.feature: endpoints +// are tried in order, locals first, so content only reaches an external model +// after every local endpoint has failed — and never at all when no external +// endpoint is configured. package summarizer import ( + "bytes" "context" "encoding/json" + "errors" "fmt" "strings" "time" + "unicode/utf8" "gitea.d-ma.be/mathias/tapir/internal/domain" ) @@ -29,19 +33,41 @@ type Endpoint struct { Model string // resolved alias, e.g. "iguana/deepseek-r1-14b" } -// Summarizer routes a transcript through the local endpoint first, then the -// optional BYO endpoint. It owns its routing (rather than delegating to -// llm.Router) so it can record which provider answered and whether the fallback -// was used — information llm.Router collapses away. +// Summarizer routes a transcript through an ordered chain of endpoints, trying +// each in turn until one returns a parseable summary. It owns its routing +// (rather than delegating to llm.Router) so it can record which provider answered +// and whether a fallback was used — information llm.Router collapses away. The +// chain ordering is the local-first guarantee: callers place local endpoints +// first and any external endpoint last, so content only reaches an external model +// after every local endpoint has failed. type Summarizer struct { - primary Endpoint - fallback *Endpoint // nil => no BYO; primary errors are returned, never sent externally - now func() time.Time + endpoints []Endpoint + maxInputChars int // transcript truncation budget; 0 = no limit + now func() time.Time } -// New constructs a Summarizer. fallback may be nil (no BYO provider configured). +// New constructs a Summarizer from a primary endpoint and an optional fallback +// (the historical local-Primary -> BYO-Fallback shape, ADR-004). A nil fallback +// means a single-endpoint chain: errors are returned, content never leaves it. func New(primary Endpoint, fallback *Endpoint) *Summarizer { - return &Summarizer{primary: primary, fallback: fallback, now: time.Now} + eps := []Endpoint{primary} + if fallback != nil { + eps = append(eps, *fallback) + } + return &Summarizer{endpoints: eps, now: time.Now} +} + +// NewChain constructs a Summarizer over an ordered endpoint chain (ADR-022). +// endpoints are tried in order; the first to return a parseable summary wins, and +// FallbackUsed is recorded true for any endpoint past the first. maxInputChars +// bounds the transcript text sent to every endpoint (0 = unbounded), so a long +// transcript does not overflow a small-context primary model's window. It panics +// on an empty chain — a wiring bug, not a runtime condition. +func NewChain(endpoints []Endpoint, maxInputChars int) *Summarizer { + if len(endpoints) == 0 { + panic("summarizer: NewChain requires at least one endpoint") + } + return &Summarizer{endpoints: endpoints, maxInputChars: maxInputChars, now: time.Now} } const systemPrompt = `You are Tapir, a video-summarization assistant. @@ -53,31 +79,35 @@ Respond with ONLY a JSON object, no prose and no code fences: - "takeaways": the actionable conclusions a viewer should leave with. Output the JSON object and nothing else.` -// Summarize implements ports.Summarizer. +// Summarize implements ports.Summarizer. It walks the endpoint chain in order: +// the first endpoint whose reply parses into a non-empty summary wins. An +// endpoint is considered failed — and the next one tried — when the model call +// errors OR when its reply cannot be parsed (a 200 with malformed JSON or a +// highlights field the model emitted as a bare string). Truncation is applied +// once, up front, so every endpoint sees the same bounded prompt. When the whole +// chain fails, the joined error is returned so the engine queues the work for +// retry and delivers no summary. func (s *Summarizer) Summarize(ctx context.Context, v domain.Video, t domain.Transcript) (domain.Summary, error) { if !t.HasText() { return domain.Summary{}, fmt.Errorf("summarize: transcript for video %s has no text", v.ID) } - user := buildUserPrompt(v, t) + user := buildUserPrompt(v, t, s.maxInputChars) - // Primary = local stack. Only on its failure is anything sent externally, - // and only when a BYO fallback is configured. - out, err := s.primary.Client.Complete(ctx, systemPrompt, user) - if err == nil { - return s.build(v, s.primary, false, out) + var errs []error + for i, ep := range s.endpoints { + out, err := ep.Client.Complete(ctx, systemPrompt, user) + if err != nil { + errs = append(errs, fmt.Errorf("%s/%s call: %w", ep.Provider, ep.Model, err)) + continue + } + sum, perr := s.build(v, ep, i > 0, out) + if perr != nil { + errs = append(errs, fmt.Errorf("%s/%s output: %w", ep.Provider, ep.Model, perr)) + continue + } + return sum, nil } - - if s.fallback == nil { - // No BYO: content was sent to the local stack only. Surface the error so - // the engine can queue the work for retry; deliver no summary. - return domain.Summary{}, fmt.Errorf("summarize: local AI failed and no BYO provider configured: %w", err) - } - - out, ferr := s.fallback.Client.Complete(ctx, systemPrompt, user) - if ferr != nil { - return domain.Summary{}, fmt.Errorf("summarize: local AI failed: %w; BYO %s failed: %v", err, s.fallback.Provider, ferr) - } - return s.build(v, *s.fallback, true, out) + return domain.Summary{}, fmt.Errorf("summarize: all %d endpoint(s) failed: %w", len(s.endpoints), errors.Join(errs...)) } func (s *Summarizer) build(v domain.Video, ep Endpoint, fallbackUsed bool, raw string) (domain.Summary, error) { @@ -89,8 +119,8 @@ func (s *Summarizer) build(v domain.Video, ep Endpoint, fallbackUsed bool, raw s UserID: v.UserID, VideoID: v.ID, Summary: parsed.Summary, - Highlights: parsed.Highlights, - Takeaways: parsed.Takeaways, + Highlights: []string(parsed.Highlights), + Takeaways: []string(parsed.Takeaways), AIProvider: ep.Provider, AIModel: ep.Model, FallbackUsed: fallbackUsed, @@ -98,20 +128,94 @@ func (s *Summarizer) build(v domain.Video, ep Endpoint, fallbackUsed bool, raw s }, nil } -func buildUserPrompt(v domain.Video, t domain.Transcript) string { +func buildUserPrompt(v domain.Video, t domain.Transcript, maxInputChars int) string { var b strings.Builder fmt.Fprintf(&b, "Title: %s\n", v.Title) if v.URL != "" { fmt.Fprintf(&b, "URL: %s\n", v.URL) } - fmt.Fprintf(&b, "\nTranscript:\n%s", t.Content) + fmt.Fprintf(&b, "\nTranscript:\n%s", truncate(t.Content, maxInputChars)) return b.String() } +// truncate caps content to max bytes on a UTF-8 rune boundary, appending a +// marker so the model knows the transcript was cut. A non-positive max (or a +// content already within budget) returns content unchanged. Bounding the input +// keeps a long transcript from overflowing a small-context model's window — the +// production failure mode where koala/phi4-mini's 8k context returned HTTP 400 on +// a 11.6k-token transcript. +func truncate(content string, max int) string { + if max <= 0 || len(content) <= max { + return content + } + cut := max + for cut > 0 && !utf8.RuneStart(content[cut]) { + cut-- + } + return content[:cut] + "\n…[transcript truncated to fit the model context]" +} + +// flexStrings is a []string that also unmarshals from a single JSON string or a +// JSON array of scalars. Small local models (koala/phi4-mini) sometimes emit +// "highlights": "one point" instead of an array, or mix in a number; rather than +// fail the whole summary on that quirk, coerce to []string. Empty/whitespace +// elements are dropped. +type flexStrings []string + +func (f *flexStrings) UnmarshalJSON(b []byte) error { + b = bytes.TrimSpace(b) + if len(b) == 0 || string(b) == "null" { + *f = nil + return nil + } + if b[0] == '[' { + var raw []json.RawMessage + if err := json.Unmarshal(b, &raw); err != nil { + return err + } + out := make([]string, 0, len(raw)) + for _, r := range raw { + s, err := rawToString(r) + if err != nil { + return err + } + if strings.TrimSpace(s) != "" { + out = append(out, s) + } + } + *f = out + return nil + } + s, err := rawToString(b) + if err != nil { + return err + } + if strings.TrimSpace(s) == "" { + *f = nil + } else { + *f = flexStrings{s} + } + return nil +} + +// rawToString renders a JSON scalar as text: a quoted string is unquoted; any +// other scalar (number, bool) is kept as its literal source so no content is lost. +func rawToString(r json.RawMessage) (string, error) { + r = bytes.TrimSpace(r) + if len(r) > 0 && r[0] == '"' { + var s string + if err := json.Unmarshal(r, &s); err != nil { + return "", err + } + return s, nil + } + return string(r), nil +} + type parsedSummary struct { - Summary string `json:"summary"` - Highlights []string `json:"highlights"` - Takeaways []string `json:"takeaways"` + Summary string `json:"summary"` + Highlights flexStrings `json:"highlights"` + Takeaways flexStrings `json:"takeaways"` } // parse extracts the JSON object from a model reply. Thinking models (qwen3, diff --git a/internal/adapters/summarizer/summarizer_test.go b/internal/adapters/summarizer/summarizer_test.go index 07e54f5..e74dec5 100644 --- a/internal/adapters/summarizer/summarizer_test.go +++ b/internal/adapters/summarizer/summarizer_test.go @@ -129,8 +129,8 @@ func TestSummarize_NoBYO_ContentOnlyLocal(t *testing.T) { local := &fakeClient{reply: goodReply} s := New(Endpoint{Client: local, Provider: "local", Model: "iguana/deepseek-r1-14b"}, nil) - if s.fallback != nil { - t.Fatal("no BYO configured but fallback endpoint is non-nil") + if len(s.endpoints) != 1 { + t.Fatalf("no BYO configured but chain has %d endpoints, want 1", len(s.endpoints)) } for i := 0; i < 3; i++ { sum, err := s.Summarize(context.Background(), testVideo(), testTranscript()) @@ -176,3 +176,97 @@ func TestParse_EmptySummaryRejected(t *testing.T) { t.Fatal("want error for empty summary (thinking model returned no content)") } } + +// parse tolerates a small model emitting "highlights" as a bare string instead +// of an array — the production koala/phi4-mini quirk that errored with +// "cannot unmarshal string into Go struct field ... highlights of type []string". +func TestParse_ToleratesStringHighlights(t *testing.T) { + p, err := parse(`{"summary":"s","highlights":"one big point","takeaways":["a","b"]}`) + if err != nil { + t.Fatalf("parse: %v", err) + } + if len(p.Highlights) != 1 || p.Highlights[0] != "one big point" { + t.Errorf("highlights = %v, want [\"one big point\"]", p.Highlights) + } + if len(p.Takeaways) != 2 { + t.Errorf("takeaways = %v, want 2", p.Takeaways) + } +} + +// Chain: an endpoint that returns a 200 with unparseable output is treated as a +// failure, and the next endpoint in the chain is tried. This is the case the old +// primary->fallback shape missed — a parse error short-circuited instead of +// falling back. +func TestSummarize_FallsBackOnMalformedOutput(t *testing.T) { + bad := &fakeClient{reply: `{"summary": not json`} + good := &fakeClient{reply: goodReply} + s := NewChain([]Endpoint{ + {Client: bad, Provider: "local", Model: "koala/phi4-mini"}, + {Client: good, Provider: "local", Model: "koala/phi4-14b"}, + }, 0) + + sum, err := s.Summarize(context.Background(), testVideo(), testTranscript()) + if err != nil { + t.Fatalf("Summarize: %v", err) + } + if sum.AIModel != "koala/phi4-14b" { + t.Errorf("AIModel = %q, want koala/phi4-14b (fell back past malformed primary)", sum.AIModel) + } + if !sum.FallbackUsed { + t.Error("FallbackUsed = false, want true") + } + if bad.calls != 1 || good.calls != 1 { + t.Errorf("calls: bad=%d good=%d, want 1 and 1", bad.calls, good.calls) + } +} + +// Chain: when every endpoint fails, no summary is produced and the joined error +// names each failure so the engine queues the work for retry. +func TestSummarize_ChainAllEndpointsFail(t *testing.T) { + a := &fakeClient{err: errors.New("context overflow")} + b := &fakeClient{reply: "not even json"} + s := NewChain([]Endpoint{ + {Client: a, Provider: "local", Model: "m1"}, + {Client: b, Provider: "berget", Model: "m2"}, + }, 0) + + if _, err := s.Summarize(context.Background(), testVideo(), testTranscript()); err == nil { + t.Fatal("want error when all endpoints fail") + } + if a.calls != 1 || b.calls != 1 { + t.Errorf("calls: a=%d b=%d, want 1 and 1", a.calls, b.calls) + } +} + +// A transcript longer than the chain's input budget is truncated before it +// reaches any model, so a small-context primary does not overflow its window. +func TestSummarize_TruncatesLongTranscript(t *testing.T) { + local := &fakeClient{reply: goodReply} + const budget = 100 + s := NewChain([]Endpoint{{Client: local, Provider: "local", Model: "m"}}, budget) + + long := domain.Transcript{ + VideoID: "vid-1", UserID: "user-1", Source: domain.SourceCaptions, + Content: strings.Repeat("word ", 1000), // 5000 bytes, well over budget + } + if _, err := s.Summarize(context.Background(), testVideo(), long); err != nil { + t.Fatalf("Summarize: %v", err) + } + // The prompt carries title/URL framing plus the truncation marker, so allow + // headroom over the raw transcript budget — but it must be far below 5000. + if len(local.lastUser) > budget+300 { + t.Errorf("prompt length = %d, want <= %d (transcript not truncated)", len(local.lastUser), budget+300) + } + if !strings.Contains(local.lastUser, "truncated") { + t.Error("truncation marker missing from prompt") + } +} + +func TestNewChain_PanicsOnEmptyChain(t *testing.T) { + defer func() { + if recover() == nil { + t.Fatal("want panic on empty endpoint chain") + } + }() + NewChain(nil, 0) +} diff --git a/internal/config/config.go b/internal/config/config.go index 24877f7..b8cc86d 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -28,8 +28,24 @@ type Config struct { GatewayURL string // GatewayKey authorizes the gateway. Read from env, never committed. GatewayKey string - // SummarizerModel is the alias in host/name form, e.g. "koala/phi4-mini". + // SummarizerModel is the primary summarizer alias in host/name form, tried + // first on every video, e.g. "koala/phi4-mini". SummarizerModel string + // FallbackModel is the LOCAL fallback alias tried when the primary fails or + // returns unparseable output (ADR-022). Kept local so content stays on the + // homelab stack. Empty disables it. Default a bigger-context local model. + FallbackModel string + // CloudFallbackModel is the worst-case EXTERNAL fallback alias, tried only + // after every local endpoint has failed (ADR-022). For client deployments set + // this empty so content never leaves the local stack. Default a berget alias. + CloudFallbackModel string + // SummaryMaxTokens caps the completion budget per summary call. Small-context + // models (koala/phi4-mini, 8k) overflow when prompt + max_tokens exceeds the + // window; a summary needs only a few hundred tokens, so the default is small. + SummaryMaxTokens int + // MaxTranscriptChars bounds the transcript text sent to the model so a long + // transcript does not overflow a small-context primary. 0 disables truncation. + MaxTranscriptChars int // SummarizerTimeout bounds a single completion call. Thinking models are // slow, so the default is generous. SummarizerTimeout time.Duration @@ -118,6 +134,10 @@ func (c Config) DexConfigured() bool { return strings.TrimSpace(c.OIDCIssuer) != const ( defaultGatewayURL = "http://koala:30401/v1" defaultSummarizerModel = "koala/phi4-mini" + defaultFallbackModel = "koala/phi4-14b" + defaultCloudFallbackModel = "berget/mistral-small" + defaultSummaryMaxTokens = 1500 + defaultMaxTranscriptChars = 18000 defaultSummarizerTimeout = 5 * time.Minute defaultYTTokenRef = "youtube/refresh_token" defaultYTConnectRedirectURL = "https://tapir.d-ma.be/oauth/youtube/callback" @@ -141,6 +161,8 @@ func Load() (Config, error) { GatewayURL: envOr("TAPIR_GATEWAY_URL", defaultGatewayURL), GatewayKey: os.Getenv("TAPIR_GATEWAY_KEY"), SummarizerModel: envOr("TAPIR_SUMMARIZER_MODEL", defaultSummarizerModel), + FallbackModel: lookupOr("TAPIR_FALLBACK_MODEL", defaultFallbackModel), + CloudFallbackModel: lookupOr("TAPIR_CLOUD_FALLBACK_MODEL", defaultCloudFallbackModel), DBDSN: os.Getenv("TAPIR_DB_DSN"), YTClientID: os.Getenv("TAPIR_YT_CLIENT_ID"), YTClientSecret: os.Getenv("TAPIR_YT_CLIENT_SECRET"), @@ -193,6 +215,21 @@ func Load() (Config, error) { } c.AutoSummarizeWindow = autoWindow + summaryTokens, err := intOr("TAPIR_SUMMARY_MAX_TOKENS", defaultSummaryMaxTokens) + if err != nil { + return Config{}, err + } + c.SummaryMaxTokens = summaryTokens + + maxChars, err := intOr("TAPIR_MAX_TRANSCRIPT_CHARS", defaultMaxTranscriptChars) + if err != nil { + return Config{}, err + } + if maxChars < 0 { + maxChars = 0 + } + c.MaxTranscriptChars = maxChars + onboard, err := intOr("TAPIR_ONBOARD_SUMMARIZE_COUNT", defaultOnboardSummarizeCount) if err != nil { return Config{}, err @@ -269,6 +306,17 @@ func envOr(key, fallback string) string { return fallback } +// lookupOr returns the env value when the key is PRESENT (even if empty), else +// fallback. Unlike envOr it lets an explicit empty value override the default — +// needed to DISABLE an optional fallback model (e.g. set the cloud fallback empty +// for a client deployment so content never leaves the local stack). +func lookupOr(key, fallback string) string { + if v, ok := os.LookupEnv(key); ok { + return v + } + return fallback +} + func intOr(key string, fallback int) (int, error) { v := os.Getenv(key) if v == "" { diff --git a/internal/config/config_test.go b/internal/config/config_test.go index ca0c99a..7b2c4db 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -1,11 +1,29 @@ package config import ( + "os" "strings" "testing" "time" ) +// unset removes an env key for the duration of the test, restoring it after. +// Needed to observe a default for a key read with LookupEnv (where present-empty +// means "explicitly disabled", not "use default"). +func unset(t *testing.T, key string) { + t.Helper() + if old, ok := os.LookupEnv(key); ok { + t.Cleanup(func() { + if err := os.Setenv(key, old); err != nil { + t.Fatalf("restore %s: %v", key, err) + } + }) + } + if err := os.Unsetenv(key); err != nil { + t.Fatalf("unset %s: %v", key, err) + } +} + // setEnv sets env vars for the test and clears them afterward, so cases don't // leak into one another. t.Setenv handles restoration. func setEnv(t *testing.T, kv map[string]string) { @@ -53,6 +71,46 @@ func TestLoad_AppliesDefaults(t *testing.T) { } } +func TestLoad_SummarizerChainDefaults(t *testing.T) { + setEnv(t, map[string]string{ + "TAPIR_SUMMARIZER_MODEL": "", + "TAPIR_SUMMARY_MAX_TOKENS": "", + "TAPIR_MAX_TRANSCRIPT_CHARS": "", + }) + unset(t, "TAPIR_FALLBACK_MODEL") + unset(t, "TAPIR_CLOUD_FALLBACK_MODEL") + + c, err := Load() + if err != nil { + t.Fatalf("Load: %v", err) + } + if c.FallbackModel != defaultFallbackModel { + t.Errorf("FallbackModel = %q, want %q", c.FallbackModel, defaultFallbackModel) + } + if c.CloudFallbackModel != defaultCloudFallbackModel { + t.Errorf("CloudFallbackModel = %q, want %q", c.CloudFallbackModel, defaultCloudFallbackModel) + } + if c.SummaryMaxTokens != defaultSummaryMaxTokens { + t.Errorf("SummaryMaxTokens = %d, want %d", c.SummaryMaxTokens, defaultSummaryMaxTokens) + } + if c.MaxTranscriptChars != defaultMaxTranscriptChars { + t.Errorf("MaxTranscriptChars = %d, want %d", c.MaxTranscriptChars, defaultMaxTranscriptChars) + } +} + +// An explicitly empty cloud-fallback env disables external routing — the lever a +// client deployment pulls so content never leaves the local stack. +func TestLoad_EmptyCloudFallbackDisables(t *testing.T) { + t.Setenv("TAPIR_CLOUD_FALLBACK_MODEL", "") + c, err := Load() + if err != nil { + t.Fatalf("Load: %v", err) + } + if c.CloudFallbackModel != "" { + t.Errorf("CloudFallbackModel = %q, want empty (disabled)", c.CloudFallbackModel) + } +} + func TestLoad_OnboardSummarizeCount(t *testing.T) { cases := []struct { name, env string diff --git a/test/acceptance/scenario_coverage_test.go b/test/acceptance/scenario_coverage_test.go index a4a4708..141131f 100644 --- a/test/acceptance/scenario_coverage_test.go +++ b/test/acceptance/scenario_coverage_test.go @@ -26,10 +26,11 @@ import ( // longer matches a real non-pending scenario. var scenarioCoverage = map[string]string{ // ai_routing.feature - "Local AI produces the summary": "TestSummarize_LocalSucceeds", - "Local AI fails and the user has a BYO provider configured": "TestSummarize_FallsBackToBYO", - "Local AI fails and the user has no BYO provider": "TestSummarize_LocalFailsNoBYO_NoExternalSend", - "A user without BYO never has content sent externally": "TestSummarize_NoBYO_ContentOnlyLocal", + "Local AI produces the summary": "TestSummarize_LocalSucceeds", + "Local AI fails and the user has a BYO provider configured": "TestSummarize_FallsBackToBYO", + "Local AI fails and the user has no BYO provider": "TestSummarize_LocalFailsNoBYO_NoExternalSend", + "A user without BYO never has content sent externally": "TestSummarize_NoBYO_ContentOnlyLocal", + "A model returns unparseable output and the next endpoint succeeds": "TestSummarize_FallsBackOnMalformedOutput", // landing_page.feature "An unauthenticated visit to the root is sent to the welcome page": "TestUnauthenticatedRootRedirectsToWelcome",