Compare commits

..
5 Commits
Author SHA1 Message Date
mathiasandClaude Opus 4.8 cc69a912f4 fix(scheduler): rotate lead user each pass so caption budget is shared
CI / Lint / Test / Vet (push) Successful in 12s
CI / Build & Import (push) Successful in 10s
Caption fetches share one per-egress-IP rate budget; whoever runs first each pass
spends the pre-throttle window before YouTube starts 429ing. ListAllUsers order
is unspecified and was stable, so the last-listed user was permanently starved —
a friendly-pilot user got 0 fetches in 12h (all rate_limited) while the
first-listed user got every successful fetch. rotateUsers left-rotates the user
order by pass index so each user leads 1/N passes and the lead slot is shared.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-10 18:07:01 +02:00
mathiasandClaude Opus 4.8 f4a0544903 fix(scheduler): cache transcripts on the scheduled path (ADR-021 regression)
buildUserRunner built the engine without engine.Transcripts = st, so the
scheduler — unlike the web "Summarize now" path — never read or wrote the shared
transcript cache. Every discovery pass re-fetched transcripts it had already
fetched, burning the scarce per-egress-IP caption budget (ADR-014) on redundant
work and starving other users' first-time fetches. The transcripts table was
empty despite summaries existing. Wire the cache on this path too.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-10 18:07:01 +02:00
mathiasandClaude Opus 4.8 e2a52789b9 feat(web): mode-aware backlog banner — stop telling manual users summaries auto-arrive
CI / Lint / Test / Vet (push) Successful in 12s
CI / Build & Import (push) Successful in 10s
The first pilot user sat in Manual mode reading "new summaries land gradually,
check back tomorrow" — copy that only makes sense in Automatic mode. Manual mode
never auto-summarizes, so the banner promised delivery that would never come.

The list page now reads the user's summarize mode and shows mode-correct copy:
- Auto: unchanged "land gradually" backlog note.
- Manual: "new videos appear here but are not summarized automatically — use the
  Summarize button" plus a "Switch to Automatic" link to /account.
The connected-but-empty first-run state is likewise mode-aware.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-10 09:14:52 +02:00
mathias e696b6405b docs(build-state): v0.15.0, summarizer is now a resilient chain (ADR-022)
CI / Lint / Test / Vet (push) Successful in 12s
CI / Build & Import (push) Successful in 10s
2026-06-10 08:59:26 +02:00
mathiasandClaude Opus 4.8 e9b5a3f3e7 feat(summarizer): resilient endpoint chain with local→cloud fallback (ADR-022)
CI / Lint / Test / Vet (push) Successful in 11s
CI / Build & Import (push) Successful in 10s
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) <noreply@anthropic.com>
2026-06-10 08:36:16 +02:00
19 changed files with 914 additions and 328 deletions
+3 -2
View File
@@ -86,7 +86,7 @@ Skills live in the canonical library `mathias/skills` and are wired into this re
## Current build state (start here for the first task) ## Current build state (start here for the first task)
The repo is **green and shipping** — last tag `v0.14.0`. `task check` passes (fmt, vet, lint, The repo is **green and shipping** — last tag `v0.15.0`. `task check` passes (fmt, vet, lint,
`go test -p 1 ./...`). Go is `1.26.1` (see `go.mod`). `go test -p 1 ./...`). Go is `1.26.1` (see `go.mod`).
- Clean Architecture core is implemented: `internal/domain` (entities), `internal/ports` - Clean Architecture core is implemented: `internal/domain` (entities), `internal/ports`
@@ -95,7 +95,8 @@ The repo is **green and shipping** — last tag `v0.14.0`. `task check` passes (
green against it. green against it.
- Adapters present under `internal/adapters/`: `youtube` (captions-first `VideoSource`, - Adapters present under `internal/adapters/`: `youtube` (captions-first `VideoSource`,
timedtext/InnerTube acquisition per ADR-010), `summarizer` + `llm` (the copied AI router, timedtext/InnerTube acquisition per ADR-010), `summarizer` + `llm` (the copied AI router,
Primary→Fallback per ADR-004), `store` (Postgres, golang-migrate migrations 001015), now a resilient endpoint chain — local primary → local fallback → external worst-case,
parse-failure-aware, ADR-004 + ADR-022), `store` (Postgres, golang-migrate migrations 001015),
`secrets` (file-backed `SecretStore`). The brain HTTP sink (ADR-005) is the remaining `secrets` (file-backed `SecretStore`). The brain HTTP sink (ADR-005) is the remaining
optional sink. optional sink.
- Stage 1 is open (ADR-012): multi-user with **DB-enforced** isolation — Postgres RLS `FORCE`d - Stage 1 is open (ADR-012): multi-user with **DB-enforced** isolation — Postgres RLS `FORCE`d
+52
View File
@@ -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 ## Rejected alternatives
Approaches considered during the 2026-06-02 planning + grill session and **deliberately not Approaches considered during the 2026-06-02 planning + grill session and **deliberately not
+36 -7
View File
@@ -3,6 +3,7 @@ package main
import ( import (
"context" "context"
"fmt" "fmt"
"strings"
"gitea.d-ma.be/mathias/tapir/internal/adapters/llm" "gitea.d-ma.be/mathias/tapir/internal/adapters/llm"
"gitea.d-ma.be/mathias/tapir/internal/adapters/secrets" "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 // 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 // `tapir run` reports the gap via its own ValidateForRun. Missing engine config
// is never an error here. // 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) { func buildProcessor(cfg config.Config, st *store.Store) (*usecase.Engine, error) {
if cfg.GatewayURL == "" || cfg.YTClientID == "" || cfg.YTClientSecret == "" || cfg.SecretsFile == "" { if cfg.GatewayURL == "" || cfg.YTClientID == "" || cfg.YTClientSecret == "" || cfg.SecretsFile == "" {
return nil, nil return nil, nil
@@ -55,13 +90,7 @@ func buildProcessor(cfg config.Config, st *store.Store) (*usecase.Engine, error)
PreferredLanguages: []string{"en"}, PreferredLanguages: []string{"en"},
}, secretStore) }, secretStore)
// Local Primary only; no BYO fallback for the demo (fallback nil). sum := buildSummarizer(cfg)
primary := summarizer.Endpoint{
Client: llm.New(cfg.GatewayURL, cfg.GatewayKey, cfg.SummarizerModel, cfg.SummarizerTimeout),
Provider: "local",
Model: cfg.SummarizerModel,
}
sum := summarizer.New(primary, nil)
// The store is both the summary sink and the shared transcript cache (ADR-021): // 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 // the engine reads stored transcripts before any caption fetch and writes
+38 -10
View File
@@ -6,9 +6,7 @@ import (
"log/slog" "log/slog"
"time" "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/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/adapters/youtube"
"gitea.d-ma.be/mathias/tapir/internal/config" "gitea.d-ma.be/mathias/tapir/internal/config"
"gitea.d-ma.be/mathias/tapir/internal/ports" "gitea.d-ma.be/mathias/tapir/internal/ports"
@@ -37,12 +35,12 @@ func buildUserRunner(cfg config.Config, st *store.Store, secretStore ports.Secre
PreferredLanguages: []string{"en"}, PreferredLanguages: []string{"en"},
}, secretStore) }, secretStore)
primary := summarizer.Endpoint{ engine := usecase.NewEngine(src, buildSummarizer(cfg), st)
Client: llm.New(cfg.GatewayURL, cfg.GatewayKey, cfg.SummarizerModel, cfg.SummarizerTimeout), // Share the transcript cache (ADR-021) on the scheduler path too — without
Provider: "local", // this every scheduled pass re-fetches transcripts it already had, burning the
Model: cfg.SummarizerModel, // 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 := usecase.NewEngine(src, summarizer.New(primary, nil), st) 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),
@@ -64,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,
@@ -74,6 +73,13 @@ func runDiscoveryPass(
return runner.Stats{} return runner.Stats{}
} }
// 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 user 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 gives every user the lead slot in turn.
users = rotateUsers(users, pass)
log.Info("scheduler: starting discovery pass", "users", len(users)) log.Info("scheduler: starting discovery pass", "users", len(users))
var total runner.Stats var total runner.Stats
for _, u := range users { for _, u := range users {
@@ -123,7 +129,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()
@@ -132,11 +139,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 {
+30 -4
View File
@@ -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,31 @@ 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"))
}
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 +147,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 +160,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)
+12
View File
@@ -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. `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 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. 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 - **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 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):** parser treats an empty summary as an error for exactly this reason. **Done (2026-06-02, Worker F):**
+13 -1
View File
@@ -31,5 +31,17 @@ Feature: Local-first AI with optional BYO fallback
When any transcript is summarized When any transcript is summarized
Then my content is only ever sent to the local AI stack 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. # Quality scoring may be added later without changing these scenarios.
+22 -2
View File
@@ -34,15 +34,35 @@ type Client struct {
httpClient *http.Client 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. // New constructs a Client.
func New(baseURL, apiKey, model string, timeout time.Duration) *Client { func New(baseURL, apiKey, model string, timeout time.Duration, opts ...Option) *Client {
return &Client{ c := &Client{
baseURL: strings.TrimRight(baseURL, "/"), baseURL: strings.TrimRight(baseURL, "/"),
apiKey: apiKey, apiKey: apiKey,
model: model, model: model,
maxTokens: defaultMaxTokens, maxTokens: defaultMaxTokens,
httpClient: &http.Client{Timeout: timeout}, httpClient: &http.Client{Timeout: timeout},
} }
for _, opt := range opts {
opt(c)
}
return c
} }
type chatRequest struct { type chatRequest struct {
+21
View File
@@ -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) { func TestClient_ReturnsErrorOnNon200(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
http.Error(w, "overloaded", http.StatusServiceUnavailable) http.Error(w, "overloaded", http.StatusServiceUnavailable)
+144 -40
View File
@@ -1,17 +1,21 @@
// Package summarizer implements ports.Summarizer backed by the copied llm // Package summarizer implements ports.Summarizer backed by the copied llm
// package's local-Primary -> BYO-Fallback routing (ADR-004). It is the only // package's routing (ADR-004, extended by ADR-022). It is the only place content
// place content ever leaves the engine toward an AI model, so it is also the // ever leaves the engine toward an AI model, so it is also the enforcement point
// enforcement point for the local-first guarantee in // for the local-first guarantee in docs/use-cases/ai_routing.feature: endpoints
// docs/use-cases/ai_routing.feature: a user with no BYO provider configured has // are tried in order, locals first, so content only reaches an external model
// their content sent to the local stack and nowhere else. // after every local endpoint has failed — and never at all when no external
// endpoint is configured.
package summarizer package summarizer
import ( import (
"bytes"
"context" "context"
"encoding/json" "encoding/json"
"errors"
"fmt" "fmt"
"strings" "strings"
"time" "time"
"unicode/utf8"
"gitea.d-ma.be/mathias/tapir/internal/domain" "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" Model string // resolved alias, e.g. "iguana/deepseek-r1-14b"
} }
// Summarizer routes a transcript through the local endpoint first, then the // Summarizer routes a transcript through an ordered chain of endpoints, trying
// optional BYO endpoint. It owns its routing (rather than delegating to // each in turn until one returns a parseable summary. It owns its routing
// llm.Router) so it can record which provider answered and whether the fallback // (rather than delegating to llm.Router) so it can record which provider answered
// was used — information llm.Router collapses away. // 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 { type Summarizer struct {
primary Endpoint endpoints []Endpoint
fallback *Endpoint // nil => no BYO; primary errors are returned, never sent externally maxInputChars int // transcript truncation budget; 0 = no limit
now func() time.Time 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 { 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. 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. - "takeaways": the actionable conclusions a viewer should leave with.
Output the JSON object and nothing else.` 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) { func (s *Summarizer) Summarize(ctx context.Context, v domain.Video, t domain.Transcript) (domain.Summary, error) {
if !t.HasText() { if !t.HasText() {
return domain.Summary{}, fmt.Errorf("summarize: transcript for video %s has no text", v.ID) 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, var errs []error
// and only when a BYO fallback is configured. for i, ep := range s.endpoints {
out, err := s.primary.Client.Complete(ctx, systemPrompt, user) out, err := ep.Client.Complete(ctx, systemPrompt, user)
if err == nil { if err != nil {
return s.build(v, s.primary, false, out) 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
} }
return domain.Summary{}, fmt.Errorf("summarize: all %d endpoint(s) failed: %w", len(s.endpoints), errors.Join(errs...))
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)
} }
func (s *Summarizer) build(v domain.Video, ep Endpoint, fallbackUsed bool, raw string) (domain.Summary, error) { 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, UserID: v.UserID,
VideoID: v.ID, VideoID: v.ID,
Summary: parsed.Summary, Summary: parsed.Summary,
Highlights: parsed.Highlights, Highlights: []string(parsed.Highlights),
Takeaways: parsed.Takeaways, Takeaways: []string(parsed.Takeaways),
AIProvider: ep.Provider, AIProvider: ep.Provider,
AIModel: ep.Model, AIModel: ep.Model,
FallbackUsed: fallbackUsed, FallbackUsed: fallbackUsed,
@@ -98,20 +128,94 @@ func (s *Summarizer) build(v domain.Video, ep Endpoint, fallbackUsed bool, raw s
}, nil }, 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 var b strings.Builder
fmt.Fprintf(&b, "Title: %s\n", v.Title) fmt.Fprintf(&b, "Title: %s\n", v.Title)
if v.URL != "" { if v.URL != "" {
fmt.Fprintf(&b, "URL: %s\n", 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() 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 { type parsedSummary struct {
Summary string `json:"summary"` Summary string `json:"summary"`
Highlights []string `json:"highlights"` Highlights flexStrings `json:"highlights"`
Takeaways []string `json:"takeaways"` Takeaways flexStrings `json:"takeaways"`
} }
// parse extracts the JSON object from a model reply. Thinking models (qwen3, // parse extracts the JSON object from a model reply. Thinking models (qwen3,
@@ -129,8 +129,8 @@ func TestSummarize_NoBYO_ContentOnlyLocal(t *testing.T) {
local := &fakeClient{reply: goodReply} local := &fakeClient{reply: goodReply}
s := New(Endpoint{Client: local, Provider: "local", Model: "iguana/deepseek-r1-14b"}, nil) s := New(Endpoint{Client: local, Provider: "local", Model: "iguana/deepseek-r1-14b"}, nil)
if s.fallback != nil { if len(s.endpoints) != 1 {
t.Fatal("no BYO configured but fallback endpoint is non-nil") t.Fatalf("no BYO configured but chain has %d endpoints, want 1", len(s.endpoints))
} }
for i := 0; i < 3; i++ { for i := 0; i < 3; i++ {
sum, err := s.Summarize(context.Background(), testVideo(), testTranscript()) 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)") 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)
}
+49 -1
View File
@@ -28,8 +28,24 @@ type Config struct {
GatewayURL string GatewayURL string
// GatewayKey authorizes the gateway. Read from env, never committed. // GatewayKey authorizes the gateway. Read from env, never committed.
GatewayKey string 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 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 // SummarizerTimeout bounds a single completion call. Thinking models are
// slow, so the default is generous. // slow, so the default is generous.
SummarizerTimeout time.Duration SummarizerTimeout time.Duration
@@ -118,6 +134,10 @@ func (c Config) DexConfigured() bool { return strings.TrimSpace(c.OIDCIssuer) !=
const ( const (
defaultGatewayURL = "http://koala:30401/v1" defaultGatewayURL = "http://koala:30401/v1"
defaultSummarizerModel = "koala/phi4-mini" defaultSummarizerModel = "koala/phi4-mini"
defaultFallbackModel = "koala/phi4-14b"
defaultCloudFallbackModel = "berget/mistral-small"
defaultSummaryMaxTokens = 1500
defaultMaxTranscriptChars = 18000
defaultSummarizerTimeout = 5 * time.Minute defaultSummarizerTimeout = 5 * time.Minute
defaultYTTokenRef = "youtube/refresh_token" defaultYTTokenRef = "youtube/refresh_token"
defaultYTConnectRedirectURL = "https://tapir.d-ma.be/oauth/youtube/callback" defaultYTConnectRedirectURL = "https://tapir.d-ma.be/oauth/youtube/callback"
@@ -141,6 +161,8 @@ func Load() (Config, error) {
GatewayURL: envOr("TAPIR_GATEWAY_URL", defaultGatewayURL), GatewayURL: envOr("TAPIR_GATEWAY_URL", defaultGatewayURL),
GatewayKey: os.Getenv("TAPIR_GATEWAY_KEY"), GatewayKey: os.Getenv("TAPIR_GATEWAY_KEY"),
SummarizerModel: envOr("TAPIR_SUMMARIZER_MODEL", defaultSummarizerModel), 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"), DBDSN: os.Getenv("TAPIR_DB_DSN"),
YTClientID: os.Getenv("TAPIR_YT_CLIENT_ID"), YTClientID: os.Getenv("TAPIR_YT_CLIENT_ID"),
YTClientSecret: os.Getenv("TAPIR_YT_CLIENT_SECRET"), YTClientSecret: os.Getenv("TAPIR_YT_CLIENT_SECRET"),
@@ -193,6 +215,21 @@ func Load() (Config, error) {
} }
c.AutoSummarizeWindow = autoWindow 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) onboard, err := intOr("TAPIR_ONBOARD_SUMMARIZE_COUNT", defaultOnboardSummarizeCount)
if err != nil { if err != nil {
return Config{}, err return Config{}, err
@@ -269,6 +306,17 @@ func envOr(key, fallback string) string {
return fallback 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) { func intOr(key string, fallback int) (int, error) {
v := os.Getenv(key) v := os.Getenv(key)
if v == "" { if v == "" {
+58
View File
@@ -1,11 +1,29 @@
package config package config
import ( import (
"os"
"strings" "strings"
"testing" "testing"
"time" "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 // setEnv sets env vars for the test and clears them afterward, so cases don't
// leak into one another. t.Setenv handles restoration. // leak into one another. t.Setenv handles restoration.
func setEnv(t *testing.T, kv map[string]string) { 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) { func TestLoad_OnboardSummarizeCount(t *testing.T) {
cases := []struct { cases := []struct {
name, env string name, env string
+12 -3
View File
@@ -233,11 +233,20 @@ func (a *App) handleList(w http.ResponseWriter, r *http.Request) {
} }
hasConnected := len(conns) > 0 hasConnected := len(conns) > 0
if isHTMX(r) { // Summarization mode drives the backlog copy: an auto user is told summaries
a.render(w, r, summaryList(buckets, hasConnected)) // land gradually; a manual user is told to click Summarize (the first pilot
// user sat in manual mode reading "land automatically" and waited forever).
autoSummarize, err := a.Store.GetAutoSummarize(r.Context(), userID)
if err != nil {
a.serverError(w, r, "summarize mode", err)
return return
} }
a.render(w, r, ListPage(buckets, f, stats, takeFlash(w, r), hasConnected, channels))
if isHTMX(r) {
a.render(w, r, summaryList(buckets, hasConnected, autoSummarize))
return
}
a.render(w, r, ListPage(buckets, f, stats, takeFlash(w, r), hasConnected, channels, autoSummarize))
} }
// handleDetail renders one summary in full (highlights, takeaways, action group). // handleDetail renders one summary in full (highlights, takeaways, action group).
+33
View File
@@ -477,3 +477,36 @@ func postAction(t *testing.T, app *web.App, videoID, action string, htmx bool) *
} }
return do(t, app, req) return do(t, app, req)
} }
// TestListManualModeBannerCopy: a manual-mode user with un-summarized videos
// sees the manual prompt (click Summarize), NOT the "summaries land
// automatically" copy that misled the first pilot user into waiting forever.
func TestListManualModeBannerCopy(t *testing.T) {
ctx := context.Background()
app := newApp(t)
p := rawPool(t)
resetDB(t, p)
seedVideo(t, p, videoX, "X Title", "https://x", time.Time{}) // pending, un-summarized
require.NoError(t, app.Store.SetAutoSummarize(ctx, userID, false))
html := body(t, do(t, app, httptest.NewRequest(http.MethodGet, "/", nil)))
require.Contains(t, html, "Manual mode")
require.Contains(t, html, "are not summarized automatically")
require.NotContains(t, html, "land gradually",
"manual-mode user must not be told summaries arrive automatically")
}
// TestListAutoModeBannerCopy: an auto-mode user with a backlog sees the
// gradual-delivery copy, not the manual prompt.
func TestListAutoModeBannerCopy(t *testing.T) {
ctx := context.Background()
app := newApp(t)
p := rawPool(t)
resetDB(t, p)
seedVideo(t, p, videoX, "X Title", "https://x", time.Time{})
require.NoError(t, app.Store.SetAutoSummarize(ctx, userID, true))
html := body(t, do(t, app, httptest.NewRequest(http.MethodGet, "/", nil)))
require.Contains(t, html, "land gradually")
require.NotContains(t, html, "are not summarized automatically")
}
+1 -1
View File
@@ -65,7 +65,7 @@ func TestParseYouTubeVideoID(t *testing.T) {
func TestListPageShowsPasteFormOnlyWhenConnected(t *testing.T) { func TestListPageShowsPasteFormOnlyWhenConnected(t *testing.T) {
render := func(connected bool) string { render := func(connected bool) string {
var buf bytes.Buffer var buf bytes.Buffer
if err := ListPage(listBuckets{}, Filter{}, PipelineStats{}, "", connected, nil).Render(context.Background(), &buf); err != nil { if err := ListPage(listBuckets{}, Filter{}, PipelineStats{}, "", connected, nil, true).Render(context.Background(), &buf); err != nil {
t.Fatalf("render: %v", err) t.Fatalf("render: %v", err)
} }
return buf.String() return buf.String()
+18 -5
View File
@@ -104,7 +104,7 @@ templ flashBanner(code string) {
// #summary-list region; a non-HTMX request renders the whole page. flash carries // #summary-list region; a non-HTMX request renders the whole page. flash carries
// a one-shot notification (e.g. "connected", "registered") surfaced on arrival // a one-shot notification (e.g. "connected", "registered") surfaced on arrival
// after a POST→redirect. // after a POST→redirect.
templ ListPage(b listBuckets, f Filter, stats PipelineStats, flash string, hasConnected bool, channels []string) { templ ListPage(b listBuckets, f Filter, stats PipelineStats, flash string, hasConnected bool, channels []string, autoSummarize bool) {
@Layout("Tapir — Summaries") { @Layout("Tapir — Summaries") {
@flashBanner(flash) @flashBanner(flash)
if hasConnected { if hasConnected {
@@ -116,14 +116,23 @@ templ ListPage(b listBuckets, f Filter, stats PipelineStats, flash string, hasCo
if stats.RateLimited > 0 || stats.Pending > 0 || stats.NoText > 0 { if stats.RateLimited > 0 || stats.Pending > 0 || stats.NoText > 0 {
@pipelineBar(stats) @pipelineBar(stats)
} }
if stats.RateLimited+stats.Pending > 0 { if (stats.RateLimited+stats.Pending) > 0 && autoSummarize {
<p class="pipeline-note muted"> <p class="pipeline-note muted">
Tapir fetches captions slowly on purpose, to respect YouTube's limits Tapir fetches captions slowly on purpose, to respect YouTube's limits
new summaries land gradually. Check back tomorrow. new summaries land gradually. Check back tomorrow.
</p> </p>
} }
if (stats.RateLimited+stats.Pending) > 0 && !autoSummarize {
<p class="pipeline-note muted">
You are in Manual mode: new videos appear here but are not summarized
automatically. Use the Summarize button on the ones you want.
</p>
<p class="pipeline-note muted">
<a href="/account">Switch to Automatic</a> to have new videos summarized for you.
</p>
}
<div id="summary-list"> <div id="summary-list">
@summaryList(b, hasConnected) @summaryList(b, hasConnected, autoSummarize)
</div> </div>
} }
} }
@@ -203,12 +212,16 @@ templ filterForm(f Filter, channels []string) {
// and a single disclosure holding the older un-summarized back-catalogue. Cards // and a single disclosure holding the older un-summarized back-catalogue. Cards
// reflow to a single column on mobile; an empty list shows a friendly first-run // reflow to a single column on mobile; an empty list shows a friendly first-run
// state instead of a blank table. // state instead of a blank table.
templ summaryList(b listBuckets, hasConnected bool) { templ summaryList(b listBuckets, hasConnected bool, autoSummarize bool) {
if b.empty() { if b.empty() {
if hasConnected { if hasConnected {
<div class="empty empty-connected"> <div class="empty empty-connected">
<strong>Your account is connected</strong> <strong>Your account is connected</strong>
<span>Tapir is finding your subscriptions and fetching captions summaries appear here gradually. Check back later.</span> if autoSummarize {
<span>Tapir is finding your subscriptions and fetching captions summaries appear here gradually. Check back later.</span>
} else {
<span>Tapir is finding your subscriptions. You are in Manual mode, so videos appear here with a Summarize button pick the ones you want, or switch to Automatic in your account.</span>
}
</div> </div>
} else { } else {
<div class="empty"> <div class="empty">
File diff suppressed because it is too large Load Diff
+5 -4
View File
@@ -26,10 +26,11 @@ import (
// longer matches a real non-pending scenario. // longer matches a real non-pending scenario.
var scenarioCoverage = map[string]string{ var scenarioCoverage = map[string]string{
// ai_routing.feature // ai_routing.feature
"Local AI produces the summary": "TestSummarize_LocalSucceeds", "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 a BYO provider configured": "TestSummarize_FallsBackToBYO",
"Local AI fails and the user has no BYO provider": "TestSummarize_LocalFailsNoBYO_NoExternalSend", "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 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 // landing_page.feature
"An unauthenticated visit to the root is sent to the welcome page": "TestUnauthenticatedRootRedirectsToWelcome", "An unauthenticated visit to the root is sent to the welcome page": "TestUnauthenticatedRootRedirectsToWelcome",