Compare commits

...
9 Commits
Author SHA1 Message Date
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
mathiasandClaude Opus 4.8 0ba78e8868 docs: refresh build-state for transcript persistence (ADR-021)
CI / Lint / Test / Vet (push) Successful in 11s
CI / Build & Import (push) Successful in 10s
Update the CLAUDE.md orientation block: last tag v0.14.0, migrations
001–015, and a transcript-persistence bullet (shared non-RLS store,
engine reads stored-first). The stale "v0.9.0" reference is corrected.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-09 23:39:27 +02:00
mathiasandClaude Opus 4.8 821d5f99cd docs(bdd): scenario for transcript reuse — re-analysis never re-fetches
Capture the ADR-021 promise as a mapped BDD scenario: re-analyzing a
stored video reads the stored transcript and does not fetch captions.
Since paste-a-URL and the onboarding burst summarize through the same
engine chokepoint (resolveTranscript, store-first), this one scenario
covers their reuse path too — there is exactly one gated caption entry
point (youtube.FetchTranscript → WaitFetchGate) and one engine caller in
front of it, so the dedup is structural, not per-feature.

Mapped to TestProcessNewVideo_SecondSummarizeDoesNotRefetch.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-09 23:37:54 +02:00
mathiasandClaude Opus 4.8 5c70408e75 feat(usecase): read stored transcript before fetching (ADR-021)
The engine now resolves transcripts store-first: a stored transcript —
including a stored SourceNone — is summarized without touching YouTube,
so re-analysis never re-fetches. On a miss it fetches through the source
(caption call still gated, ADR-014) and persists the terminal outcome for
the next analysis by any user. A transient SourceRateLimited is surfaced
to the runner for per-user backoff but never cached, so persistence can
never mask a 429 as a permanent "no transcript".

The TranscriptStore is optional (nil → fetch every time), keeping the
pure-core and scaffold wiring valid. cmd/tapir wires the store as both
summary sink and transcript cache, so `tapir run` and the web summarize
path (incl. paste + onboarding) all share the dedup.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-09 23:35:04 +02:00
mathiasandClaude Opus 4.8 cb6917ca59 feat(store): shared, video-keyed transcript persistence (ADR-021)
Reshape the dead per-user transcripts table (PK videos.id, user_id,
RLS-FORCEd — never read or written by app code) into the shared public
caption store ADR-021 specifies: keyed by (provider, provider_video_id),
no user_id, NOT RLS-scoped. Migration 015 (reversible). Add
ports.TranscriptStore + Store.GetTranscript/SaveTranscript via the raw
pool (no withUser): public content, shared across users by construction.
SaveTranscript persists only terminal outcomes (captions/none) and
refuses SourceRateLimited so a transient 429 can never be stored as a
false permanent absence (ADR-014).

Flip the isolation proof: transcripts leaves the RLS-scoped set;
TestTranscriptsTableIsSharedNotRLS asserts it is the SINGLE non-RLS
surface (writable/readable with no user scope, no user_id column, RLS off
on it alone, still on every user-owned table) — the proof the
public-content classification was applied exactly here and leaked nowhere.
appPool made idempotent so two tests can build it. Adjust the 010/011/014
up-down migration tests for the new HEAD. account.go: user deletion no
longer strips shared transcripts. Reconcile data-model.md + CLAUDE.md.

Wiring the engine to read-stored-first is the next commit.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-09 23:32:12 +02:00
mathiasandClaude Opus 4.8 099b2d4c68 docs(decisions): ADR-021 — shared, video-keyed transcript persistence
Persist transcripts in a single shared table keyed by
(provider, provider_video_id) — public caption content, NOT RLS-scoped —
so re-analysis (re-summarize, paste of an already-seen video, a second
user with overlapping subs) never re-fetches from YouTube. The avoided
cost is the rate-gated, reputation-risky caption fetch (ADR-010/014), not
LLM re-summarization, which is why this reopens the transcripts half of
the "no global cross-tenant table" rejection while videos stay per-user.
Summaries remain RLS-scoped (ADR-012 unchanged). The gate is neither
bypassed nor weakened — persistence reduces fetch frequency, not pacing.

Annotate the rejected-alternatives row to record the partial reopen.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-09 23:24:11 +02:00
mathiasandClaude Opus 4.8 f66c1bcdcc feat(web): real channel filter — multi-select of the user's channels
CI / Build & Import (push) Successful in 10s
CI / Lint / Test / Vet (push) Successful in 11s
The free-text 'channel' filter was dead: it exact-matched SummaryRow.Channel,
which is just the provider ('youtube'), because videos never stored their source
channel. Now they do.

- migration 014: videos.channel_title (nullable; existing rows backfill on the
  next discovery pass, pasted videos immediately).
- discovery (NewVideos) + paste (VideoByID) populate channel_title; UpsertVideo
  persists it, preserving an existing title when an update arrives empty.
- store.DistinctChannels lists a user's channels (RLS-scoped); SummaryRow carries
  ChannelTitle via the shared projection.
- Filter: single Channel -> Channels []string, matching on ChannelTitle; the feed
  renders a multi-select of DistinctChannels (hidden until channels exist).
- migrate tests: 014 reversibility + fixed the relative-step counts in the 010/011
  up/down tests (014 shifted the topology).

TDD throughout: channel persist + distinct, adapter channel wiring, multi-channel
filter match, handler channel filter, migration up/down.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-09 23:02:54 +02:00
40 changed files with 2062 additions and 723 deletions
+11 -4
View File
@@ -46,8 +46,9 @@ These caused real mistakes that were caught and corrected; the corrections are l
(See `DECISIONS.md` for full rationale. Listed here so you don't propose them.)
- **No Supabase** — reuse Dex / ESO+1Password / Postgres (ADR-002).
- **No global cross-tenant video/transcript table** — per-user isolation (data-model). Dedup
across users is a Future C concern, not a Stage 0/1 default.
- **No global cross-tenant *video* table** — videos stay per-user (data-model). Transcripts ARE
shared since ADR-021 (public caption content, keyed by `(provider, provider_video_id)`, non-RLS)
so re-analysis never re-fetches; the *videos* half of cross-tenant dedup stays a Future C concern.
- **No audio-download + speech-to-text in the core path** — captions-first (ADR-007). STT is a
deferred, bounded optional component.
- **No public SaaS / sign-up / billing / Google OAuth verification at scale** — Future C,
@@ -85,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)
The repo is **green and shipping** — last tag `v0.9.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`).
- Clean Architecture core is implemented: `internal/domain` (entities), `internal/ports`
@@ -94,13 +95,19 @@ The repo is **green and shipping** — last tag `v0.9.0`. `task check` passes (f
green against it.
- Adapters present under `internal/adapters/`: `youtube` (captions-first `VideoSource`,
timedtext/InnerTube acquisition per ADR-010), `summarizer` + `llm` (the copied AI router,
Primary→Fallback per ADR-004), `store` (Postgres, golang-migrate migrations 001006),
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
optional sink.
- Stage 1 is open (ADR-012): multi-user with **DB-enforced** isolation — Postgres RLS `FORCE`d
on all user-owned tables (migration 003), two-user isolation test in
`internal/adapters/store/rls_test.go`. Registration gate, per-user YouTube web connect, and
account management (disconnect / delete, ADR-013) all shipped.
- Transcript persistence (ADR-021, migration 015): transcripts are a **shared, non-RLS** store
keyed by `(provider, provider_video_id)` — the single exception to the isolation boundary
(`TestTranscriptsTableIsSharedNotRLS`). The engine reads stored transcripts before any caption
fetch (`usecase.resolveTranscript`), so re-analysis — re-summarize, paste-a-URL, onboarding
burst — never re-touches YouTube. Per-user summaries/videos stay RLS-scoped.
- `cmd/tapir` subcommands: `list`, `show`, `auth` (interactive host-side OAuth), `run` (batch
watch→summarize), `serve` (the HTMX+Templ web reader/writer under `internal/web`, a new
transport over the unchanged engine/ports — ADR-003). `tapir env` prints config.
+113 -1
View File
@@ -743,6 +743,118 @@ collapse keys off the same window (`App.RecencyWindow=0` → everything inline).
---
## ADR-021 — Persist transcripts as shared, video-keyed public content (re-analysis never re-fetches)
**Status:** Accepted (2026-06-09). **Reopens the transcripts half of** the "Global cross-tenant
`videos`/`transcripts` table" rejection (data-model.md). **Builds on ADR-007** (captions-first),
**ADR-010/ADR-014** (the per-IP caption rate gate), and **ADR-012** (per-user RLS isolation).
**Context.** Every summarization fetches the transcript fresh through the caption path, even when
the exact same transcript was fetched moments ago — for the same user re-summarizing, or for a
second user who happens to watch the same video. The caption fetch is the one genuinely scarce,
genuinely risky operation in the system: YouTube's timedtext endpoint is unofficial and per-IP
rate-limited (ADR-010), and tripping it risks the maintainer's Google standing (ADR-014). So the
operation we most want to *avoid repeating* is the one we currently repeat unconditionally. A
transcript is **public content** — the same words YouTube serves to anyone — and carries nothing
user-identifying. The per-user isolation that protects summaries, feeds, and tokens (ADR-012) is
the wrong shape for it: it forces a re-fetch per user for data that is identical across users.
The original rejection ("Global cross-tenant `videos`/`transcripts` table") bundled videos and
transcripts together and rejected both on the grounds that "at 15 users, re-summarizing is
cheaper than the coupling." That reasoning holds for **videos** (per-user feed rows, genuinely
user-scoped) but not for **transcripts**: the cost being avoided is not LLM re-summarization, it
is a *rate-gated, reputation-risky network fetch*, and that cost is paid per re-fetch regardless
of user count. One re-fetch avoided is strictly worth more than the coupling it removes.
**Decision.**
1. **A single shared `transcripts` table, keyed by the cross-user dedup key
`(provider, provider_video_id)`** — the stable public identity of the video, not Tapir's
internal per-user `videos.id`. Columns: the key, `source` (`captions`/`none`), `language`,
`content`, `fetched_at`. It holds **only public caption content + the video's public id**
nothing user-identifying — and is therefore **NOT RLS-scoped**: no `user_id`, no policy, no
`FORCE ROW LEVEL SECURITY`. This is the deliberate, single exception to the ADR-012 isolation
boundary, and the only one.
2. **Summarize path becomes read-stored-first.** Have a stored transcript for this video? →
summarize from the stored text, **no caption fetch**. No stored transcript? → fetch *through
the unchanged gate* (ADR-014) → store it → summarize. The gate is neither bypassed nor
weakened; persistence reduces how *often* we reach it, never how *fast*.
3. **De-facto cross-user dedup is the intended behaviour, not a feature with a switch.** Two
users who share a video share the one transcript row. A permanent `source = 'none'` (no
captions) is stored too, so a known-caption-less video is not re-fetched by anyone. A
transient 429 (`SourceRateLimited`) is **never** stored as terminal — it stays a per-user
retry via the existing `transcript_status` backoff (ADR-014), so persistence cannot mask a
rate-limit into a false "no transcript."
4. **Per-user `summaries` stay RLS-scoped (ADR-012 unchanged)** and reference the transcript by
video id. Videos stay per-user. Only transcripts go shared.
**Consequences.** Re-analysis (re-summarize, different model, paste of an already-seen video,
onboarding of a second user with overlapping subscriptions) never re-touches YouTube — the
primary win, and it *reduces* aggregate caption-gate pressure, reinforcing ADR-010/ADR-014 rather
than straining them. The isolation surface gains exactly one non-RLS table; an isolation test
asserts the boundary is *exactly* there and has not leaked to any user-owned table (this is the
proof the public-content classification was implemented as designed). It also unblocks
multi-model / customizable analysis (re-run analysis on stored text for free) — enabling that is
this ADR's point; building it is separate.
**Reversibility.** The read-stored-first check is the only behavioural coupling; removing it
restores fetch-every-time. The down-migration recreates the per-user RLS-scoped transcripts shape
(001/003). No user-facing surface depends on cross-user sharing — sharing is the *storage shape*,
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
@@ -757,7 +869,7 @@ maps to the ADR that settles it.
| Lifting shared packages into a `brain-common` module | Couples Tapir's release cycle to the monolith for negligible code savings | ADR-004 |
| Importing/replicating the filesystem `brain` package | Assumes co-location with the brain git checkout; wrong for a standalone networked service | ADR-005 |
| Reusing `ingestion`'s `oauth` package for YouTube/Vimeo | Same name, opposite direction — it's inbound MCP-server auth, not outbound provider OAuth | ADR-006 |
| Global cross-tenant `videos`/`transcripts` table (dedup) | Reintroduces the cross-domain DB coupling the homelab review is removing; at 15 users, re-summarizing is cheaper than the coupling | data-model.md |
| Global cross-tenant `videos`/`transcripts` table (dedup) | Reintroduces the cross-domain DB coupling the homelab review is removing; at 15 users, re-summarizing is cheaper than the coupling. **Transcripts half reopened by ADR-021** — the avoided cost there is a rate-gated, reputation-risky *caption fetch*, not LLM re-summarization, so it outweighs the coupling; **videos stay per-user.** | data-model.md, **ADR-021** (transcripts only) |
| Audio-download + Whisper STT in the core path | ToS-grey, breakage-prone (yt-dlp), contends for koala GPU with the JEPA PoC; captions alone test the core hypothesis | ADR-007 |
| Building multi-tenant SaaS / Google OAuth verification now | "Real users soon" was lowered to Future B; SaaS machinery before the Stage 0 self-use gate is the primary documented anti-goal | ADR-008, VISION |
| Delegating the S5 reuse spike to an agent swarm | A 1-hour sequential read-and-judge with a single coupled conclusion; orchestration overhead exceeds the work, and it's Diamond-1 judgment the maintainer wanted to own | (process note) |
+42 -8
View File
@@ -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,15 +90,14 @@ 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)
return usecase.NewEngine(src, sum, st), nil
// 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
// resolved ones back, so re-analysis never re-touches YouTube.
eng := usecase.NewEngine(src, sum, st)
eng.Transcripts = st
return eng, nil
}
// engineProcessor adapts the engine (which works in terms of a domain.Video) to
+1 -8
View File
@@ -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),
+21 -13
View File
@@ -10,12 +10,15 @@ only opaque references to them; the secret material lives in ESO/1Password (ADR-
## Design decisions baked into this model
- **Per-user isolation, not a shared global video table.** The earlier draft proposed a
global `videos`/`transcripts` table deduped across tenants. Rejected for Future B: it
reintroduces exactly the cross-domain coupling the homelab architecture review is
removing, and at 15 users the cost of occasionally re-summarizing the same video is
trivial compared to the isolation it would cost. Each user's data is self-contained.
(Revisit only if Future C makes GPU/transcription cost dominate — a new ADR, not a default.)
- **Per-user isolation for everything except transcripts.** The earlier draft proposed a
global `videos`/`transcripts` table deduped across tenants. **Videos** stay per-user and
RLS-scoped — a shared video table reintroduces exactly the cross-domain coupling the homelab
architecture review is removing. **Transcripts**, however, are now shared (ADR-021): keyed by
`(provider, provider_video_id)`, no `user_id`, **not** RLS-scoped. The cost avoided there is
not LLM re-summarization but a rate-gated, reputation-risky caption fetch (ADR-010/014), which
is paid per re-fetch regardless of user count — so persisting public caption content once and
sharing it strictly beats the coupling it removes. Everything else each user owns is
self-contained; `rls_test.go` proves transcripts is the single exception.
- **Secrets by reference only.** Tables hold a `secret_ref` (opaque string/UUID resolved via
the `SecretStore` port), never tokens or keys.
- **The brain sink is just a delivery target.** No brain-specific tables. Whether a summary
@@ -34,7 +37,7 @@ erDiagram
USER ||--o{ AI_CREDENTIAL : "has (planned)"
VIDEO_CONNECTION ||--o{ SUBSCRIPTION : "exposes (planned)"
SUBSCRIPTION ||--o{ VIDEO : "produces (per user)"
VIDEO ||--o| TRANSCRIPT : "has at most one"
VIDEO }o--o| TRANSCRIPT : "shares one by (provider, provider_video_id) — not FK (ADR-021)"
VIDEO ||--o| SUMMARY : "has at most one"
SUMMARY ||--o{ SINK_DELIVERY : "delivered via"
USER ||--o{ CHANNEL_ERROR : "reports unavailable channels"
@@ -92,12 +95,12 @@ erDiagram
timestamptz rate_limited_at "backoff clock for 429 retries (migration 007)"
}
TRANSCRIPT {
uuid video_id PK_FK
uuid user_id FK
text provider PK "part of shared key (ADR-021)"
text provider_video_id PK "part of shared key — the cross-user dedup key"
text source "captions | none"
text language
text content "null when source = none"
timestamptz resolved_at
timestamptz fetched_at
}
SUMMARY {
uuid id PK
@@ -173,8 +176,12 @@ mechanism.
`transcript_status` and `rate_limited_at` (migration 007) track caption-fetch outcomes for
rate-limit backoff: `NULL` = not attempted; `rate_limited` = 429 seen, skip until
`NOW() - rate_limited_at > TAPIR_FETCH_BACKOFF`; `fetched` = resolved; `none` = no transcript.
- **TRANSCRIPT** — at most one per video. `source = none` records "checked, no usable
transcript" so the watcher doesn't reprocess (ADR-007). `content` null in that case.
- **TRANSCRIPT** — shared public caption content, one row per `(provider, provider_video_id)`,
**not** RLS-scoped and carrying no `user_id` (ADR-021). Two users who watch the same video
share the one row; the summarize path reads it before any caption fetch, so re-analysis never
re-touches YouTube (ADR-010/014). `source = none` records "checked, no usable transcript" so
no one reprocesses (ADR-007); `content` null in that case. A transient 429 is never stored
here — it stays a per-user retry via `VIDEO.transcript_status`.
- **SUMMARY** — at most one per video. `fallback_used` + `ai_provider`/`ai_model` make the
"is local good enough?" question queryable (the Stage 0 quality signal). `highlights`/
`takeaways` as jsonb to stay schema-flexible while the output format settles.
@@ -231,7 +238,8 @@ queue, doesn't replace it). Deferred until there's a reason.
## Explicitly out of scope (Future C)
- Global cross-tenant video/transcript dedup (rejected above).
- Global cross-tenant *video* dedup (rejected above). Note: cross-tenant *transcript* sharing
is now in scope and shipped (ADR-021); only the videos half stays per-user.
- Sharding / per-tenant physical databases.
- Soft-delete + full audit trail on connections/credentials (a Stage 2 hardening item; add
via ADR when Stage 2 work starts).
+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.
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):**
+13 -1
View File
@@ -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.
@@ -31,5 +31,16 @@ Feature: Summarize new videos from subscribed channels
When the watcher sees "Designing for Attention" again
Then Tapir does not produce a second summary for it
Scenario: Re-analyzing a stored video does not re-fetch its transcript
Given a transcript for "Designing for Attention" is already stored
When the video is summarized again
Then Tapir reads the stored transcript
And Tapir does not fetch captions from YouTube
# Captions-first is the core path (ADR-007). Audio-download + speech-to-text is
# deferred and intentionally has no scenario here yet.
#
# Transcript persistence (ADR-021): the stored transcript is shared, keyed by
# (provider, provider_video_id) and read before any caption fetch, so the
# re-analysis scenario above also covers paste-a-URL and the onboarding burst —
# both summarize through the same engine chokepoint.
+22 -2
View File
@@ -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 {
+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) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
http.Error(w, "overloaded", http.StatusServiceUnavailable)
+7 -4
View File
@@ -10,10 +10,13 @@ import (
// DeleteUser permanently removes a user and all of their data. It runs through
// withUser so RLS confines every statement to the calling user's own rows.
//
// Deleting the users row cascades (ON DELETE CASCADE) to videos, transcripts,
// summaries (→ sink_deliveries), video_connections, and the user_identities map
// referential-integrity cascades bypass RLS, so a user's child rows are removed
// even though the deleting connection is scoped. summary_actions and login_events
// Deleting the users row cascades (ON DELETE CASCADE) to videos, summaries
// (→ sink_deliveries), video_connections, and the user_identities map
// referential-integrity cascades bypass RLS, so a user's child rows are removed
// even though the deleting connection is scoped. Transcripts are NOT removed:
// since ADR-021 they are shared public content keyed by (provider,
// provider_video_id) with no user_id, so another user may still reference the
// same row — a user deletion must not strip shared caption content. summary_actions and login_events
// are the exceptions: each carries a user_id but has NO foreign key to users
// (migrations 002 and 010), so the cascade does not reach them; they are deleted
// explicitly in the same scoped transaction. Deleting an absent user is a no-op
+37 -2
View File
@@ -53,7 +53,11 @@ func TestMigration010LoginEventsUpDown(t *testing.T) {
require.True(t, loginEventsExists(t), "login_events must exist at latest migration")
m := fileMigrator(t)
// 011, 012, 013 sit above 010; step them down first so 010 is exercised in isolation.
// 011..015 sit above 010; step them down first so 010 is exercised in isolation.
require.NoError(t, m.Steps(-1), "down 015 reshapes transcripts, login_events intact")
require.True(t, loginEventsExists(t), "015 down leaves login_events intact")
require.NoError(t, m.Steps(-1), "down 014 drops channel_title, login_events intact")
require.True(t, loginEventsExists(t), "014 down leaves login_events intact")
require.NoError(t, m.Steps(-1), "down 013 drops channel_errors, login_events intact")
require.True(t, loginEventsExists(t), "013 down leaves login_events intact")
require.NoError(t, m.Steps(-1), "down 012 is a no-op, login_events intact")
@@ -64,7 +68,7 @@ func TestMigration010LoginEventsUpDown(t *testing.T) {
require.NoError(t, m.Steps(-1), "down 010 must drop login_events")
require.False(t, loginEventsExists(t), "login_events must be gone after the down migration")
require.NoError(t, m.Steps(4), "up must recreate 010 then re-apply 011, 012, 013")
require.NoError(t, m.Steps(6), "up must recreate 010 then re-apply 011..015")
require.True(t, loginEventsExists(t), "login_events must be restored after the up migration")
}
@@ -87,6 +91,8 @@ func TestMigration011AutoSummarizeDefaultUpDown(t *testing.T) {
require.Equal(t, "true", autoSummarizeDefault(t), "011 sets the default to TRUE")
m := fileMigrator(t)
require.NoError(t, m.Steps(-1), "down 015 reshapes transcripts")
require.NoError(t, m.Steps(-1), "down 014 drops channel_title")
require.NoError(t, m.Steps(-1), "down 013 drops channel_errors")
require.NoError(t, m.Steps(-1), "down 012 is a no-op")
require.NoError(t, m.Steps(-1), "down 011 reverts the column default")
@@ -96,6 +102,35 @@ func TestMigration011AutoSummarizeDefaultUpDown(t *testing.T) {
require.Equal(t, "true", autoSummarizeDefault(t))
require.NoError(t, m.Steps(1), "up 012 runs clean (no FORCE RLS on fresh schema)")
require.NoError(t, m.Steps(1), "up 013 creates channel_errors")
require.NoError(t, m.Steps(1), "up 014 recreates channel_title")
require.NoError(t, m.Steps(1), "up 015 reshapes transcripts to shared")
}
// channelTitleExists reports whether videos.channel_title is present.
func channelTitleExists(t *testing.T) bool {
t.Helper()
var exists bool
require.NoError(t, rawPool(t).QueryRow(context.Background(),
`SELECT EXISTS (SELECT 1 FROM information_schema.columns
WHERE table_name = 'videos' AND column_name = 'channel_title')`).Scan(&exists))
return exists
}
// TestMigration014VideoChannelTitleUpDown proves 014 is reversible: down drops
// videos.channel_title, up recreates it.
func TestMigration014VideoChannelTitleUpDown(t *testing.T) {
newStore(t) // latest (014 applied)
require.True(t, channelTitleExists(t), "channel_title exists at latest migration")
m := fileMigrator(t)
require.NoError(t, m.Steps(-1), "down 015 reshapes transcripts, channel_title intact")
require.True(t, channelTitleExists(t), "015 down leaves channel_title intact")
require.NoError(t, m.Steps(-1), "down 014 must drop channel_title")
require.False(t, channelTitleExists(t), "channel_title must be gone after the down migration")
require.NoError(t, m.Steps(1), "up 014 must recreate channel_title")
require.True(t, channelTitleExists(t), "channel_title must be restored after the up migration")
require.NoError(t, m.Steps(1), "up 015 restores the shared transcripts shape (HEAD)")
}
// TestMigration012FixAutoSummarizeRLS proves 012 runs cleanly and flips any
@@ -0,0 +1 @@
ALTER TABLE videos DROP COLUMN channel_title;
@@ -0,0 +1,8 @@
-- Store the source channel's title per video so the list can offer a real
-- channel filter (multi-select of the user's channels) instead of the dead
-- free-text field that only ever matched the provider string. Nullable: existing
-- rows backfill on the next discovery pass (UpsertVideo writes it); pasted videos
-- get it immediately from videos.list. No FK to a channels table at Stage 0 — the
-- title is a denormalised display/filter value, consistent with the existing
-- subscription_id-stays-NULL stance (data-model.md).
ALTER TABLE videos ADD COLUMN channel_title TEXT;
@@ -0,0 +1,19 @@
-- Down 015: restore the per-user RLS-scoped transcripts shape (001 + 003).
DROP TABLE transcripts;
CREATE TABLE transcripts (
video_id UUID PRIMARY KEY REFERENCES videos(id) ON DELETE CASCADE,
user_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE,
source TEXT NOT NULL,
language TEXT,
content TEXT,
resolved_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE INDEX idx_transcripts_user_id ON transcripts(user_id);
ALTER TABLE transcripts ENABLE ROW LEVEL SECURITY;
ALTER TABLE transcripts FORCE ROW LEVEL SECURITY;
CREATE POLICY transcripts_isolation ON transcripts
FOR ALL
USING (user_id = current_setting('tapir.current_user_id', true)::uuid);
@@ -0,0 +1,30 @@
-- Migration 015: transcripts become SHARED public-content storage (ADR-021).
--
-- The per-user transcripts table from 001 (PK videos.id, user_id NOT NULL, RLS
-- FORCEd in 003) was dead: no application code ever read or wrote it — only the
-- transcript_status columns on `videos` (007) carried fetch outcomes. ADR-021
-- repurposes it as the single shared store of public caption content, keyed by
-- the cross-user dedup key (provider, provider_video_id) — the video's public
-- identity, not Tapir's per-user videos.id — so re-analysis never re-fetches
-- from YouTube (ADR-010/014).
--
-- It holds ONLY public caption content + the video's public id (nothing
-- user-identifying), so it is deliberately NOT RLS-scoped: no user_id, no
-- policy, no FORCE. This is the single, intentional exception to the ADR-012
-- isolation boundary; rls_test.go asserts the boundary is exactly here and
-- nowhere else. Dropping the old table drops its RLS policy with it; it held no
-- real data, so drop+recreate loses nothing.
DROP TABLE transcripts;
CREATE TABLE transcripts (
provider TEXT NOT NULL,
provider_video_id TEXT NOT NULL,
source TEXT NOT NULL, -- 'captions' (content set) | 'none' (no captions; content NULL)
language TEXT,
content TEXT,
fetched_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
PRIMARY KEY (provider, provider_video_id)
);
COMMENT ON TABLE transcripts IS
'Shared public caption content keyed by (provider, provider_video_id). NOT RLS-scoped — public content only, de-facto cross-user dedup (ADR-021).';
+4 -1
View File
@@ -29,6 +29,7 @@ type SummaryRow struct {
ProviderVideoID string // videos.provider_video_id; empty when no videos row
Title string // videos.title; empty when no videos row
Channel string // videos.provider for now; empty when no videos row
ChannelTitle string // videos.channel_title; the source channel, for display + filtering
URL string // videos.url; empty when no videos row
PublishedAt time.Time // videos.published_at; zero when absent
Summary string
@@ -137,7 +138,8 @@ const selectVideo = `
COALESCE(s.created_at, v.seen_at),
(s.id IS NOT NULL) AS summarized,
v.summarize_requested,
COALESCE(v.transcript_status, '')
COALESCE(v.transcript_status, ''),
COALESCE(v.channel_title, '')
FROM videos v
LEFT JOIN summaries s ON s.video_id = v.id AND s.user_id = v.user_id`
@@ -254,6 +256,7 @@ func scanVideoRow(rows pgx.Row) (SummaryRow, error) {
&row.Summarized,
&row.SummarizeRequested,
&row.TranscriptStatus,
&row.ChannelTitle,
); err != nil {
return SummaryRow{}, fmt.Errorf("store: scan video: %w", err)
}
+70 -14
View File
@@ -22,9 +22,11 @@ import (
// (no GUC set → zero rows) proves the enforcement path is live, not bypassed.
// userIsolatedTables are the tables that carry a user_id and whose policy keys
// directly off the tapir.current_user_id GUC.
// directly off the tapir.current_user_id GUC. transcripts is deliberately ABSENT
// — ADR-021 made it shared public content (non-RLS); TestTranscriptsTableIsSharedNotRLS
// proves that is the only place the isolation boundary moved.
var userIsolatedTables = []string{
"users", "videos", "transcripts", "summaries", "summary_actions", "login_events", "video_connections",
"users", "videos", "summaries", "summary_actions", "login_events", "video_connections",
}
// allIsolatedTables adds sink_deliveries, whose ownership is derived from its
@@ -38,9 +40,10 @@ type seeded struct {
summaryID string
}
// seedUser inserts one full chain (user → video → transcript → summary →
// action → delivery) as the superuser pool, which bypasses RLS so both users'
// data lands regardless of the GUC.
// seedUser inserts one full chain (user → video → summary → action → delivery)
// as the superuser pool, which bypasses RLS so both users' data lands regardless
// of the GUC. Transcripts are NOT seeded here: they are shared, non-RLS public
// content (ADR-021), so they have no place in a per-user isolation chain.
func seedUser(t *testing.T, p *pgxpool.Pool, userID string) seeded {
t.Helper()
ctx := context.Background()
@@ -54,11 +57,6 @@ func seedUser(t *testing.T, p *pgxpool.Pool, userID string) seeded {
VALUES ($1, 'youtube', $2, 'title') RETURNING id`,
userID, "vid-"+userID).Scan(&videoID))
_, err = p.Exec(ctx,
`INSERT INTO transcripts (video_id, user_id, source, content)
VALUES ($1, $2, 'captions', 'words')`, videoID, userID)
require.NoError(t, err)
var summaryID string
require.NoError(t, p.QueryRow(ctx,
`INSERT INTO summaries (user_id, video_id, summary) VALUES ($1, $2, 'sum')
@@ -92,9 +90,16 @@ func appPool(t *testing.T, super *pgxpool.Pool) *pgxpool.Pool {
t.Helper()
ctx := context.Background()
// Idempotent across test runs (schema/role persist for the TestMain PG).
_, _ = super.Exec(ctx, `DROP ROLE IF EXISTS app`)
_, err := super.Exec(ctx, `CREATE ROLE app LOGIN PASSWORD 'app'`)
// Idempotent across tests AND runs: the role persists for the TestMain PG and
// owns granted privileges, so a plain DROP ROLE fails once any GRANT exists
// (and more than one test now builds an app pool). Create only if absent; the
// GRANTs below are themselves idempotent.
_, err := super.Exec(ctx,
`DO $$ BEGIN
IF NOT EXISTS (SELECT FROM pg_roles WHERE rolname = 'app') THEN
CREATE ROLE app LOGIN PASSWORD 'app';
END IF;
END $$`)
require.NoError(t, err)
_, err = super.Exec(ctx, `GRANT USAGE ON SCHEMA public TO app`)
require.NoError(t, err)
@@ -186,7 +191,6 @@ func TestRLSEnforcesPerUserIsolation(t *testing.T) {
{"update users", `UPDATE users SET display_name = 'hacked' WHERE id = $1`, b.userID},
{"update videos", `UPDATE videos SET title = 'hacked' WHERE user_id = $1`, b.userID},
{"queue videos summarize", `UPDATE videos SET summarize_requested = TRUE WHERE id = $1`, b.videoID},
{"update transcripts", `UPDATE transcripts SET content = 'hacked' WHERE user_id = $1`, b.userID},
{"update summaries", `UPDATE summaries SET summary = 'hacked' WHERE user_id = $1`, b.userID},
{"update summary_actions", `UPDATE summary_actions SET action = 'skipped' WHERE user_id = $1`, b.userID},
{"update login_events", `UPDATE login_events SET seen_at = NOW() WHERE user_id = $1`, b.userID},
@@ -235,3 +239,55 @@ func TestRLSEnforcesPerUserIsolation(t *testing.T) {
_ = a // a's ids are seeded for the symmetric read assertions above
}
// TestTranscriptsTableIsSharedNotRLS is the ADR-021 isolation proof: transcripts
// is the ONE shared, non-RLS surface, and the public-content classification
// leaked to nothing else. It is the inverse of TestRLSEnforcesPerUserIsolation —
// where that asserts deny-all on every user-owned table, this asserts transcripts
// is readable and writable with no user scope at all, holds no user_id, and is
// the single table with row-level security switched off.
func TestTranscriptsTableIsSharedNotRLS(t *testing.T) {
newStore(t)
super := rawPool(t)
resetDB(t, super)
app := appPool(t, super)
ctx := context.Background()
// 1. Shared + non-RLS: with NO GUC set, the app role both writes and reads a
// transcript. On an RLS table this would be deny-all (zero rows), exactly as
// the main isolation test asserts for every user-owned table.
_, err := app.Exec(ctx,
`INSERT INTO transcripts (provider, provider_video_id, source, content)
VALUES ('youtube', 'shared-vid', 'captions', 'public words')`)
require.NoError(t, err, "app role must write shared transcript content with no user scope")
require.Equal(t, 1, scopedCount(t, app, "", "transcripts"),
"transcripts must be readable with NO user scope — it is shared, non-RLS (ADR-021)")
// 2. No user_id column: the table holds only public caption content + the
// video's public id, nothing user-identifying.
var hasUserID bool
require.NoError(t, super.QueryRow(ctx,
`SELECT EXISTS (SELECT 1 FROM information_schema.columns
WHERE table_name = 'transcripts' AND column_name = 'user_id')`).Scan(&hasUserID))
require.False(t, hasUserID, "transcripts must carry no user_id (ADR-021 public content)")
// 3. The boundary is EXACTLY here: every user-owned table still has row-level
// security enabled; transcripts alone has it off. This is the proof the
// non-RLS classification was applied to transcripts and leaked nowhere else.
for _, table := range allIsolatedTables {
require.True(t, rlsEnabled(t, super, table),
"%s must still enforce row-level security — isolation must not have regressed", table)
}
require.False(t, rlsEnabled(t, super, "transcripts"),
"transcripts must be the single table with row-level security OFF (the one shared surface)")
}
// rlsEnabled reports whether a public table has ROW LEVEL SECURITY enabled.
func rlsEnabled(t *testing.T, p *pgxpool.Pool, table string) bool {
t.Helper()
var enabled bool
require.NoError(t, p.QueryRow(context.Background(),
`SELECT relrowsecurity FROM pg_class
WHERE relname = $1 AND relnamespace = 'public'::regnamespace`, table).Scan(&enabled))
return enabled
}
+66
View File
@@ -0,0 +1,66 @@
package store
import (
"context"
"errors"
"fmt"
"github.com/jackc/pgx/v5"
"gitea.d-ma.be/mathias/tapir/internal/domain"
)
// GetTranscript returns the shared, stored transcript for a video keyed by the
// cross-user dedup key (provider, providerVideoID), and whether one exists
// (ADR-021). It reads via the raw pool, NOT withUser: the table holds public
// content with no user_id and no RLS policy, so it is shared across users by
// construction. A stored SourceNone is a real hit (ok == true, HasText() ==
// false) — a known caption-less video, so the caller skips without re-fetching.
func (s *Store) GetTranscript(ctx context.Context, provider, providerVideoID string) (domain.Transcript, bool, error) {
var source, lang, content string
err := s.pool.QueryRow(ctx,
`SELECT source, COALESCE(language, ''), COALESCE(content, '')
FROM transcripts WHERE provider = $1 AND provider_video_id = $2`,
provider, providerVideoID).Scan(&source, &lang, &content)
if errors.Is(err, pgx.ErrNoRows) {
return domain.Transcript{}, false, nil
}
if err != nil {
return domain.Transcript{}, false, fmt.Errorf("store: get transcript: %w", err)
}
return domain.Transcript{
Source: domain.TranscriptSource(source),
Language: lang,
Content: content,
}, true, nil
}
// SaveTranscript upserts the shared transcript for (provider, providerVideoID).
// Only terminal outcomes belong here: SourceCaptions (with text) or SourceNone
// (no captions). A transient SourceRateLimited is rejected so persistence never
// masks a 429 as a permanent absence — that stays a per-user retry (ADR-014).
// Last write wins on conflict (a later re-fetch may correct an entry). It writes
// via the raw pool, NOT withUser — public content, shared, non-RLS (ADR-021).
func (s *Store) SaveTranscript(ctx context.Context, provider, providerVideoID string, t domain.Transcript) error {
switch t.Source {
case domain.SourceCaptions, domain.SourceNone:
// terminal — persist
case domain.SourceRateLimited:
return fmt.Errorf("store: refusing to persist transient rate-limited transcript for %s/%s", provider, providerVideoID)
default:
return fmt.Errorf("store: invalid transcript source %q", t.Source)
}
_, err := s.pool.Exec(ctx,
`INSERT INTO transcripts (provider, provider_video_id, source, language, content)
VALUES ($1, $2, $3, NULLIF($4, ''), NULLIF($5, ''))
ON CONFLICT (provider, provider_video_id)
DO UPDATE SET source = EXCLUDED.source,
language = EXCLUDED.language,
content = EXCLUDED.content,
fetched_at = NOW()`,
provider, providerVideoID, string(t.Source), t.Language, t.Content)
if err != nil {
return fmt.Errorf("store: save transcript: %w", err)
}
return nil
}
@@ -0,0 +1,87 @@
package store_test
import (
"context"
"testing"
"github.com/stretchr/testify/require"
"gitea.d-ma.be/mathias/tapir/internal/adapters/store"
"gitea.d-ma.be/mathias/tapir/internal/domain"
"gitea.d-ma.be/mathias/tapir/internal/ports"
)
// Static check: Store satisfies the shared TranscriptStore port (ADR-021).
var _ ports.TranscriptStore = (*store.Store)(nil)
func TestSaveAndGetTranscript_RoundTrip(t *testing.T) {
s := newStore(t)
resetDB(t, rawPool(t))
ctx := context.Background()
want := domain.Transcript{Source: domain.SourceCaptions, Language: "en", Content: "the words"}
require.NoError(t, s.SaveTranscript(ctx, "youtube", "vid-1", want))
got, ok, err := s.GetTranscript(ctx, "youtube", "vid-1")
require.NoError(t, err)
require.True(t, ok, "a saved transcript must be found")
require.Equal(t, domain.SourceCaptions, got.Source)
require.Equal(t, "en", got.Language)
require.Equal(t, "the words", got.Content)
require.True(t, got.HasText())
}
func TestGetTranscript_Miss(t *testing.T) {
s := newStore(t)
resetDB(t, rawPool(t))
_, ok, err := s.GetTranscript(context.Background(), "youtube", "absent")
require.NoError(t, err, "a miss is not an error")
require.False(t, ok)
}
// A stored "no captions" outcome is a real hit: callers must skip without
// re-fetching, so ok is true even though there is no text (ADR-021 / ADR-007).
func TestSaveAndGetTranscript_NoneIsAStoredHit(t *testing.T) {
s := newStore(t)
resetDB(t, rawPool(t))
ctx := context.Background()
require.NoError(t, s.SaveTranscript(ctx, "youtube", "vid-none", domain.Transcript{Source: domain.SourceNone}))
got, ok, err := s.GetTranscript(ctx, "youtube", "vid-none")
require.NoError(t, err)
require.True(t, ok, "a stored SourceNone is a hit, not a miss")
require.Equal(t, domain.SourceNone, got.Source)
require.False(t, got.HasText())
}
// A transient 429 must never be persisted as a terminal transcript, or a later
// read would mask the rate-limit as a permanent "no transcript" (ADR-014).
func TestSaveTranscript_RejectsRateLimited(t *testing.T) {
s := newStore(t)
resetDB(t, rawPool(t))
err := s.SaveTranscript(context.Background(), "youtube", "vid-429",
domain.Transcript{Source: domain.SourceRateLimited})
require.Error(t, err)
_, ok, _ := s.GetTranscript(context.Background(), "youtube", "vid-429")
require.False(t, ok, "a rejected rate-limited save must leave nothing stored")
}
func TestSaveTranscript_UpsertLastWriteWins(t *testing.T) {
s := newStore(t)
resetDB(t, rawPool(t))
ctx := context.Background()
require.NoError(t, s.SaveTranscript(ctx, "youtube", "vid-up", domain.Transcript{Source: domain.SourceNone}))
require.NoError(t, s.SaveTranscript(ctx, "youtube", "vid-up",
domain.Transcript{Source: domain.SourceCaptions, Language: "en", Content: "now resolved"}))
got, ok, err := s.GetTranscript(ctx, "youtube", "vid-up")
require.NoError(t, err)
require.True(t, ok)
require.Equal(t, domain.SourceCaptions, got.Source)
require.Equal(t, "now resolved", got.Content)
}
+35 -6
View File
@@ -46,14 +46,15 @@ func (s *Store) UpsertVideo(ctx context.Context, v domain.Video) (string, error)
}
if err := tx.QueryRow(ctx,
`INSERT INTO videos (user_id, provider, provider_video_id, title, url, published_at)
VALUES ($1, $2, $3, $4, $5, $6)
`INSERT INTO videos (user_id, provider, provider_video_id, title, url, published_at, channel_title)
VALUES ($1, $2, $3, $4, $5, $6, $7)
ON CONFLICT (user_id, provider, provider_video_id) DO UPDATE SET
title = EXCLUDED.title,
url = EXCLUDED.url,
published_at = EXCLUDED.published_at
title = EXCLUDED.title,
url = EXCLUDED.url,
published_at = EXCLUDED.published_at,
channel_title = COALESCE(NULLIF(EXCLUDED.channel_title, ''), videos.channel_title)
RETURNING id`,
v.UserID, provider, v.ProviderVideoID, v.Title, v.URL, nullTime(v.PublishedAt),
v.UserID, provider, v.ProviderVideoID, v.Title, v.URL, nullTime(v.PublishedAt), v.ChannelTitle,
).Scan(&id); err != nil {
return fmt.Errorf("store: upsert video: %w", err)
}
@@ -110,3 +111,31 @@ func (s *Store) NewestUnsummarizedVideoIDs(ctx context.Context, userID string, l
}
return ids, nil
}
// DistinctChannels returns the user's distinct, non-empty source channel titles
// (the channels they have videos from), alphabetically — the option list for the
// feed's channel filter. RLS-scoped via withUser.
func (s *Store) DistinctChannels(ctx context.Context, userID string) ([]string, error) {
var out []string
if err := s.withUser(ctx, userID, func(tx pgx.Tx) error {
rows, err := tx.Query(ctx,
`SELECT DISTINCT channel_title FROM videos
WHERE user_id = $1 AND channel_title IS NOT NULL AND channel_title <> ''
ORDER BY channel_title`, userID)
if err != nil {
return fmt.Errorf("store: distinct channels: %w", err)
}
defer rows.Close()
for rows.Next() {
var c string
if err := rows.Scan(&c); err != nil {
return fmt.Errorf("store: scan channel: %w", err)
}
out = append(out, c)
}
return rows.Err()
}); err != nil {
return nil, err
}
return out, nil
}
+26
View File
@@ -113,3 +113,29 @@ func TestNewestUnsummarizedVideoIDs(t *testing.T) {
require.NoError(t, err)
require.Empty(t, none, "limit 0 returns nothing")
}
func TestUpsertVideoPersistsChannelAndDistinctChannels(t *testing.T) {
ctx := context.Background()
s := newStore(t)
resetDB(t, rawPool(t))
mk := func(pid, channel string) {
v := ytVideo(userA, pid, pid)
v.ChannelTitle = channel
_, err := s.UpsertVideo(ctx, v)
require.NoError(t, err)
}
mk("aa11111aaaa", "Acme Talks")
mk("bb22222bbbb", "Acme Talks") // same channel
mk("cc33333cccc", "Zeta Channel")
// userB's channel must not leak.
vb := ytVideo(userB, "dd44444dddd", "x")
vb.ChannelTitle = "Bravo Only"
_, err := s.UpsertVideo(ctx, vb)
require.NoError(t, err)
got, err := s.DistinctChannels(ctx, userA)
require.NoError(t, err)
require.Equal(t, []string{"Acme Talks", "Zeta Channel"}, got,
"distinct, alphabetical, user-scoped (no Bravo Only)")
}
+144 -40
View File
@@ -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,
@@ -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)
}
+4 -1
View File
@@ -21,7 +21,7 @@ func TestVideoByID(t *testing.T) {
if got := r.URL.Query().Get("part"); got != "snippet" {
t.Errorf("expected part=snippet, got %q", got)
}
_, _ = w.Write([]byte(`{"items":[{"snippet":{"title":"Never Gonna Give You Up","publishedAt":"2026-05-20T09:00:00Z"}}]}`))
_, _ = w.Write([]byte(`{"items":[{"snippet":{"title":"Never Gonna Give You Up","channelTitle":"Rick Astley","publishedAt":"2026-05-20T09:00:00Z"}}]}`))
})
v, err := a.VideoByID(context.Background(), "u1", id)
@@ -34,6 +34,9 @@ func TestVideoByID(t *testing.T) {
if v.ProviderVideoID != id || v.Title != "Never Gonna Give You Up" {
t.Errorf("unexpected video: %+v", v)
}
if v.ChannelTitle != "Rick Astley" {
t.Errorf("ChannelTitle = %q, want Rick Astley", v.ChannelTitle)
}
if v.Provider != domain.ProviderYouTube || v.URL != "https://www.youtube.com/watch?v="+id {
t.Errorf("video not wired correctly: %+v", v)
}
+5 -2
View File
@@ -234,6 +234,7 @@ func (a *Adapter) NewVideos(ctx context.Context, sub domain.Subscription) ([]dom
Provider: domain.ProviderYouTube,
ProviderVideoID: vid,
Title: item.Snippet.Title,
ChannelTitle: sub.ChannelTitle,
URL: "https://www.youtube.com/watch?v=" + vid,
PublishedAt: item.Snippet.PublishedAt,
})
@@ -270,6 +271,7 @@ func (a *Adapter) VideoByID(ctx context.Context, userID, videoID string) (domain
Provider: domain.ProviderYouTube,
ProviderVideoID: videoID,
Title: it.Snippet.Title,
ChannelTitle: it.Snippet.ChannelTitle,
URL: "https://www.youtube.com/watch?v=" + videoID,
PublishedAt: it.Snippet.PublishedAt,
}, nil
@@ -376,8 +378,9 @@ type playlistItemListResponse struct {
type videoListResponse struct {
Items []struct {
Snippet struct {
Title string `json:"title"`
PublishedAt time.Time `json:"publishedAt"`
Title string `json:"title"`
ChannelTitle string `json:"channelTitle"`
PublishedAt time.Time `json:"publishedAt"`
} `json:"snippet"`
} `json:"items"`
}
+4 -1
View File
@@ -131,7 +131,7 @@ func TestNewVideos(t *testing.T) {
}`))
})
sub := domain.Subscription{ID: "s1", UserID: "u1", ChannelID: "UC_acme"}
sub := domain.Subscription{ID: "s1", UserID: "u1", ChannelID: "UC_acme", ChannelTitle: "Acme Channel"}
vids, err := a.NewVideos(context.Background(), sub)
if err != nil {
t.Fatalf("NewVideos: %v", err)
@@ -143,6 +143,9 @@ func TestNewVideos(t *testing.T) {
if v.ProviderVideoID != "vid1" || v.Title != "Designing for Attention" {
t.Errorf("unexpected video: %+v", v)
}
if v.ChannelTitle != "Acme Channel" {
t.Errorf("ChannelTitle = %q, want Acme Channel", v.ChannelTitle)
}
if v.Provider != domain.ProviderYouTube || v.URL != "https://www.youtube.com/watch?v=vid1" {
t.Errorf("video not wired correctly: %+v", v)
}
+49 -1
View File
@@ -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 == "" {
+58
View File
@@ -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
+1
View File
@@ -73,6 +73,7 @@ type Video struct {
Provider Provider
ProviderVideoID string
Title string
ChannelTitle string
URL string
PublishedAt time.Time
SeenAt time.Time
+20
View File
@@ -27,6 +27,26 @@ type Summarizer interface {
Summarize(ctx context.Context, v domain.Video, t domain.Transcript) (domain.Summary, error)
}
// TranscriptStore persists transcripts as shared, video-keyed public content
// (ADR-021). It is keyed by the cross-user dedup key (provider, providerVideoID)
// — the video's public identity, NOT Tapir's per-user videos.id — and holds only
// public caption content, so it is deliberately NOT user-scoped: two users who
// share a video share the one row. The engine reads it before any caption fetch
// so re-analysis never re-touches YouTube (ADR-010/014).
type TranscriptStore interface {
// GetTranscript returns the stored transcript for a video and whether one
// exists. A stored Source == SourceNone (captions permanently absent) is a
// real hit: ok is true and HasText() is false, so callers skip without
// re-fetching. A transient rate-limit is never stored, so it never appears
// here as a false absence.
GetTranscript(ctx context.Context, provider, providerVideoID string) (t domain.Transcript, ok bool, err error)
// SaveTranscript upserts the transcript for (provider, providerVideoID). Only
// terminal outcomes are persisted: SourceCaptions (with text) or SourceNone.
// SourceRateLimited must NOT be passed — it is a per-user retry (ADR-014), not
// a shared terminal state.
SaveTranscript(ctx context.Context, provider, providerVideoID string, t domain.Transcript) error
}
// Sink delivers a summary to a destination (user store, brain, ...).
// Implementations fail independently of one another.
type Sink interface {
+42 -2
View File
@@ -27,6 +27,13 @@ type Engine struct {
AI ports.Summarizer
Sinks []ports.Sink
// Transcripts, when set, is the shared transcript cache (ADR-021): the engine
// reads it before any caption fetch and writes resolved transcripts back, so
// re-analysis — the same user re-summarizing, or a second user with the same
// video — never re-touches YouTube (ADR-010/014). Optional: nil disables
// persistence (fetch every time), keeping the pure-core/scaffold wiring valid.
Transcripts ports.TranscriptStore
// processed dedups videos within this engine's lifetime so a video is not
// summarized twice when the watcher sees it again. Durable cross-restart
// dedup is the store's concern (a resolved TRANSCRIPT / existing SUMMARY,
@@ -57,9 +64,9 @@ type ProcessResult struct {
// resolve transcript -> (summarize -> deliver) | skip.
// See docs/use-cases/summarize_new_video.feature.
func (e *Engine) ProcessNewVideo(ctx context.Context, v domain.Video) (ProcessResult, error) {
t, err := e.Source.FetchTranscript(ctx, v)
t, err := e.resolveTranscript(ctx, v)
if err != nil {
return ProcessResult{Video: v}, fmt.Errorf("fetch transcript: %w", err)
return ProcessResult{Video: v}, err
}
if !t.HasText() {
// No usable transcript: record the skip, produce no summary, deliver nothing
@@ -86,6 +93,39 @@ func (e *Engine) ProcessNewVideo(ctx context.Context, v domain.Video) (ProcessRe
return ProcessResult{Video: v, Summary: &sum, TranscriptSource: string(t.Source)}, errors.Join(errs...)
}
// resolveTranscript returns v's transcript, reading the shared store first
// (ADR-021): a stored transcript — including a stored SourceNone (captions
// permanently absent) — is returned without touching YouTube, so re-analysis
// never re-fetches. On a store miss it fetches through the source (which gates
// the caption call, ADR-014) and persists the terminal outcome so the next
// analysis, for any user, reads from the store. A transient SourceRateLimited is
// returned to the caller (the runner stamps a per-user backoff) but never stored,
// so persistence can never mask a 429 as a permanent "no transcript". When no
// TranscriptStore is wired the engine simply fetches every time.
func (e *Engine) resolveTranscript(ctx context.Context, v domain.Video) (domain.Transcript, error) {
if e.Transcripts != nil {
stored, ok, err := e.Transcripts.GetTranscript(ctx, string(v.Provider), v.ProviderVideoID)
if err != nil {
return domain.Transcript{}, fmt.Errorf("get stored transcript: %w", err)
}
if ok {
return stored, nil
}
}
t, err := e.Source.FetchTranscript(ctx, v)
if err != nil {
return domain.Transcript{}, fmt.Errorf("fetch transcript: %w", err)
}
if e.Transcripts != nil && t.Source != domain.SourceRateLimited {
if err := e.Transcripts.SaveTranscript(ctx, string(v.Provider), v.ProviderVideoID, t); err != nil {
return domain.Transcript{}, fmt.Errorf("save transcript: %w", err)
}
}
return t, nil
}
// ProcessNewVideos walks a user's subscriptions and processes each newly seen
// video. Only videos surfaced via the user's subscriptions are considered, so a
// channel the user is not subscribed to is never processed. A video already
+187
View File
@@ -0,0 +1,187 @@
package usecase
import (
"context"
"testing"
"gitea.d-ma.be/mathias/tapir/internal/domain"
)
// These tests pin the ADR-021 read-stored-first behaviour at the engine core:
// a stored transcript is summarized without re-touching the source, a miss
// fetches once and persists, and a transient rate-limit is never cached.
type recordingSource struct {
transcript domain.Transcript
fetchCalls int
}
func (s *recordingSource) ListSubscriptions(context.Context, string) ([]domain.Subscription, error) {
return nil, nil
}
func (s *recordingSource) NewVideos(context.Context, domain.Subscription) ([]domain.Video, error) {
return nil, nil
}
func (s *recordingSource) FetchTranscript(context.Context, domain.Video) (domain.Transcript, error) {
s.fetchCalls++
return s.transcript, nil
}
type fakeTranscriptStore struct {
stored map[string]domain.Transcript
saves int
}
func newFakeTranscriptStore() *fakeTranscriptStore {
return &fakeTranscriptStore{stored: make(map[string]domain.Transcript)}
}
func (f *fakeTranscriptStore) key(provider, id string) string { return provider + "|" + id }
func (f *fakeTranscriptStore) GetTranscript(_ context.Context, provider, id string) (domain.Transcript, bool, error) {
t, ok := f.stored[f.key(provider, id)]
return t, ok, nil
}
func (f *fakeTranscriptStore) SaveTranscript(_ context.Context, provider, id string, t domain.Transcript) error {
f.saves++
f.stored[f.key(provider, id)] = t
return nil
}
type countingSummarizer struct{ calls int }
func (c *countingSummarizer) Summarize(_ context.Context, v domain.Video, _ domain.Transcript) (domain.Summary, error) {
c.calls++
return domain.Summary{VideoID: v.ID, UserID: v.UserID, Summary: "s", AIProvider: "local"}, nil
}
type nopSink struct{}
func (nopSink) Name() string { return "nop" }
func (nopSink) Deliver(context.Context, domain.Summary) error { return nil }
func testVideo() domain.Video {
return domain.Video{ID: "v1", UserID: "u1", Provider: domain.ProviderYouTube, ProviderVideoID: "yt1"}
}
func TestProcessNewVideo_StoredTranscriptSkipsFetch(t *testing.T) {
src := &recordingSource{}
ts := newFakeTranscriptStore()
ts.stored[ts.key("youtube", "yt1")] = domain.Transcript{Source: domain.SourceCaptions, Content: "stored words"}
sum := &countingSummarizer{}
eng := NewEngine(src, sum, nopSink{})
eng.Transcripts = ts
res, err := eng.ProcessNewVideo(context.Background(), testVideo())
if err != nil {
t.Fatalf("ProcessNewVideo: %v", err)
}
if src.fetchCalls != 0 {
t.Fatalf("stored transcript must not re-fetch from source; got %d fetches", src.fetchCalls)
}
if ts.saves != 0 {
t.Fatalf("a store hit must not re-save; got %d saves", ts.saves)
}
if sum.calls != 1 || res.Summary == nil {
t.Fatalf("expected a summary from the stored transcript; calls=%d summary=%v", sum.calls, res.Summary)
}
}
func TestProcessNewVideo_StoreMissFetchesAndPersists(t *testing.T) {
src := &recordingSource{transcript: domain.Transcript{Source: domain.SourceCaptions, Language: "en", Content: "fetched words"}}
ts := newFakeTranscriptStore()
sum := &countingSummarizer{}
eng := NewEngine(src, sum, nopSink{})
eng.Transcripts = ts
if _, err := eng.ProcessNewVideo(context.Background(), testVideo()); err != nil {
t.Fatalf("ProcessNewVideo: %v", err)
}
if src.fetchCalls != 1 {
t.Fatalf("a store miss must fetch exactly once; got %d", src.fetchCalls)
}
if ts.saves != 1 {
t.Fatalf("a fetched transcript must be persisted; got %d saves", ts.saves)
}
got, ok, _ := ts.GetTranscript(context.Background(), "youtube", "yt1")
if !ok || got.Content != "fetched words" {
t.Fatalf("persisted transcript not readable back: ok=%v content=%q", ok, got.Content)
}
}
// The second summarize of the same video reads the persisted transcript and does
// NOT re-fetch — the primary ADR-021 win, proven end to end at the engine.
func TestProcessNewVideo_SecondSummarizeDoesNotRefetch(t *testing.T) {
src := &recordingSource{transcript: domain.Transcript{Source: domain.SourceCaptions, Content: "words"}}
ts := newFakeTranscriptStore()
eng := NewEngine(src, &countingSummarizer{}, nopSink{})
eng.Transcripts = ts
for i := 0; i < 2; i++ {
if _, err := eng.ProcessNewVideo(context.Background(), testVideo()); err != nil {
t.Fatalf("pass %d: %v", i, err)
}
}
if src.fetchCalls != 1 {
t.Fatalf("the second summarize must reuse the stored transcript; got %d fetches", src.fetchCalls)
}
}
// A stored "no captions" outcome short-circuits before both fetch and summarize.
func TestProcessNewVideo_StoredNoneSkipsFetchAndSummarize(t *testing.T) {
src := &recordingSource{}
ts := newFakeTranscriptStore()
ts.stored[ts.key("youtube", "yt1")] = domain.Transcript{Source: domain.SourceNone}
sum := &countingSummarizer{}
eng := NewEngine(src, sum, nopSink{})
eng.Transcripts = ts
res, err := eng.ProcessNewVideo(context.Background(), testVideo())
if err != nil {
t.Fatalf("ProcessNewVideo: %v", err)
}
if !res.Skipped {
t.Fatal("a stored SourceNone must skip")
}
if src.fetchCalls != 0 || sum.calls != 0 {
t.Fatalf("stored none must neither fetch nor summarize; fetches=%d calls=%d", src.fetchCalls, sum.calls)
}
}
// A transient 429 is surfaced (so the runner backs off per-user) but never cached
// as a shared terminal state — otherwise it would mask a rate-limit as permanent.
func TestProcessNewVideo_RateLimitedIsNotPersisted(t *testing.T) {
src := &recordingSource{transcript: domain.Transcript{Source: domain.SourceRateLimited}}
ts := newFakeTranscriptStore()
eng := NewEngine(src, &countingSummarizer{}, nopSink{})
eng.Transcripts = ts
res, err := eng.ProcessNewVideo(context.Background(), testVideo())
if err != nil {
t.Fatalf("ProcessNewVideo: %v", err)
}
if !res.Skipped || res.TranscriptSource != string(domain.SourceRateLimited) {
t.Fatalf("expected a rate-limited skip; skipped=%v source=%q", res.Skipped, res.TranscriptSource)
}
if ts.saves != 0 {
t.Fatalf("a transient rate-limit must not be persisted; got %d saves", ts.saves)
}
}
// With no TranscriptStore wired the engine fetches every time (back-compat).
func TestProcessNewVideo_NilStoreFetchesEveryTime(t *testing.T) {
src := &recordingSource{transcript: domain.Transcript{Source: domain.SourceCaptions, Content: "words"}}
eng := NewEngine(src, &countingSummarizer{}, nopSink{})
for i := 0; i < 2; i++ {
if _, err := eng.ProcessNewVideo(context.Background(), testVideo()); err != nil {
t.Fatalf("pass %d: %v", i, err)
}
}
if src.fetchCalls != 2 {
t.Fatalf("nil store must fetch every time; got %d", src.fetchCalls)
}
}
+23 -4
View File
@@ -21,6 +21,9 @@ import (
// fake without a database.
type Store interface {
ListVideos(ctx context.Context, userID string, limit int) ([]store.SummaryRow, error)
// DistinctChannels lists the user's source channels — the options for the
// feed's channel multi-select filter.
DistinctChannels(ctx context.Context, userID string) ([]string, error)
GetSummaryByVideo(ctx context.Context, userID, videoID string) (*store.SummaryRow, error)
GetVideoRow(ctx context.Context, userID, videoID string) (*store.SummaryRow, error)
ActionsFor(ctx context.Context, userID string, videoIDs []string) (map[string][]string, error)
@@ -196,7 +199,7 @@ func (a *App) handleList(w http.ResponseWriter, r *http.Request) {
}
q := r.URL.Query()
f := Filter{
Channel: q.Get("channel"),
Channels: nonEmptyStrings(q["channel"]),
From: q.Get("from"),
To: q.Get("to"),
OnlySummarized: q.Get("summarized") == "1",
@@ -211,6 +214,13 @@ func (a *App) handleList(w http.ResponseWriter, r *http.Request) {
rows := f.apply(allRows)
buckets := bucketRows(rows, a.recencyCutoff())
// Channel options for the multi-select filter (the user's source channels).
channels, err := a.Store.DistinctChannels(r.Context(), userID)
if err != nil {
a.serverError(w, r, "distinct channels", err)
return
}
// hasConnected drives both the paste box (shown to ANY connected user, #2) and
// the empty-state copy (a fresh account with a connection but no discovery pass
// yet reads "connected, summaries land gradually" rather than "nothing here").
@@ -223,11 +233,20 @@ func (a *App) handleList(w http.ResponseWriter, r *http.Request) {
}
hasConnected := len(conns) > 0
if isHTMX(r) {
a.render(w, r, summaryList(buckets, hasConnected))
// Summarization mode drives the backlog copy: an auto user is told summaries
// 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
}
a.render(w, r, ListPage(buckets, f, stats, takeFlash(w, r), hasConnected))
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).
+38 -3
View File
@@ -207,13 +207,15 @@ func TestListChannelFilter(t *testing.T) {
resetDB(t, p)
require.NoError(t, deliver(ctx, app, videoX, "body x"))
seedVideo(t, p, videoX, "X Title", "https://x", time.Time{})
_, err := p.Exec(ctx, `UPDATE videos SET channel_title = 'Acme Channel' WHERE id = $1`, videoX)
require.NoError(t, err)
// Channel is "youtube" for seeded rows; a non-matching filter hides them.
rec := do(t, app, httptest.NewRequest(http.MethodGet, "/?channel=vimeo", nil))
// Selecting a different channel hides the row; selecting its channel shows it.
rec := do(t, app, httptest.NewRequest(http.MethodGet, "/?channel=Other+Channel", nil))
require.Equal(t, http.StatusOK, rec.Code)
require.NotContains(t, body(t, rec), "X Title")
rec = do(t, app, httptest.NewRequest(http.MethodGet, "/?channel=youtube", nil))
rec = do(t, app, httptest.NewRequest(http.MethodGet, "/?channel=Acme+Channel", nil))
require.Contains(t, body(t, rec), "X Title")
}
@@ -475,3 +477,36 @@ func postAction(t *testing.T, app *web.App, videoID, action string, htmx bool) *
}
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")
}
+18 -1
View File
@@ -3,6 +3,7 @@ package web
import (
"bytes"
"context"
"gitea.d-ma.be/mathias/tapir/internal/adapters/store"
"strings"
"testing"
)
@@ -64,7 +65,7 @@ func TestParseYouTubeVideoID(t *testing.T) {
func TestListPageShowsPasteFormOnlyWhenConnected(t *testing.T) {
render := func(connected bool) string {
var buf bytes.Buffer
if err := ListPage(listBuckets{}, Filter{}, PipelineStats{}, "", connected).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)
}
return buf.String()
@@ -78,3 +79,19 @@ func TestListPageShowsPasteFormOnlyWhenConnected(t *testing.T) {
t.Errorf("disconnected feed must not show the paste form")
}
}
func TestFilterMatchesMultipleChannels(t *testing.T) {
f := Filter{Channels: []string{"Acme", "Zeta"}}
row := func(ch string) store.SummaryRow { return store.SummaryRow{ChannelTitle: ch, Summarized: true} }
rows := []store.SummaryRow{row("Acme"), row("Beta"), row("Zeta")}
got := f.apply(rows)
if len(got) != 2 || got[0].ChannelTitle != "Acme" || got[1].ChannelTitle != "Zeta" {
t.Fatalf("multi-channel filter = %+v, want Acme+Zeta only", got)
}
// Empty selection = no channel constraint (all pass).
if n := len(Filter{}.apply(rows)); n != 3 {
t.Fatalf("no channel filter should pass all rows, got %d", n)
}
}
+26 -4
View File
@@ -2,6 +2,7 @@ package web
import (
"regexp"
"slices"
"strings"
"time"
"unicode/utf8"
@@ -469,7 +470,7 @@ func (b listBuckets) empty() bool {
// Dates are kept as the raw YYYY-MM-DD strings so the form re-renders the user's
// input verbatim; parsing happens in matchFilter.
type Filter struct {
Channel string
Channels []string // selected channel titles; empty = all channels
From string
To string
OnlySummarized bool // show only videos that have a summary
@@ -480,7 +481,28 @@ type Filter struct {
// filter) the bar is hidden so the connect CTA stands alone (UX review C1); a
// filter that happens to match nothing still shows the bar so it can be cleared.
func (f Filter) active() bool {
return f.Channel != "" || f.From != "" || f.To != "" || f.OnlySummarized
return len(f.Channels) > 0 || f.From != "" || f.To != "" || f.OnlySummarized
}
// HasChannel reports whether a channel is currently selected (drives the
// multi-select's selected state in the view).
func (f Filter) HasChannel(c string) bool {
return slices.Contains(f.Channels, c)
}
// nonEmptyStrings drops blank entries. A channel multi-select submits real
// channel titles; this guards against a stray empty value reaching the filter.
func nonEmptyStrings(ss []string) []string {
out := ss[:0:0]
for _, s := range ss {
if strings.TrimSpace(s) != "" {
out = append(out, s)
}
}
if len(out) == 0 {
return nil
}
return out
}
// matches reports whether a row satisfies the filter. Channel is an exact match;
@@ -491,7 +513,7 @@ func (f Filter) matches(r store.SummaryRow) bool {
if f.OnlySummarized && !r.Summarized {
return false
}
if f.Channel != "" && r.Channel != f.Channel {
if len(f.Channels) > 0 && !slices.Contains(f.Channels, r.ChannelTitle) {
return false
}
if from, ok := parseDate(f.From); ok {
@@ -521,7 +543,7 @@ func parseDate(s string) (time.Time, bool) {
// apply returns the subset of rows matching the filter, preserving order.
func (f Filter) apply(rows []store.SummaryRow) []store.SummaryRow {
if f.Channel == "" && f.From == "" && f.To == "" && !f.OnlySummarized {
if len(f.Channels) == 0 && f.From == "" && f.To == "" && !f.OnlySummarized {
return rows
}
out := rows[:0:0]
+30 -8
View File
@@ -104,26 +104,35 @@ templ flashBanner(code string) {
// #summary-list region; a non-HTMX request renders the whole page. flash carries
// a one-shot notification (e.g. "connected", "registered") surfaced on arrival
// after a POST→redirect.
templ ListPage(b listBuckets, f Filter, stats PipelineStats, flash string, hasConnected bool) {
templ ListPage(b listBuckets, f Filter, stats PipelineStats, flash string, hasConnected bool, channels []string, autoSummarize bool) {
@Layout("Tapir — Summaries") {
@flashBanner(flash)
if hasConnected {
@pasteForm()
}
if !b.empty() || f.active() {
@filterForm(f)
@filterForm(f, channels)
}
if stats.RateLimited > 0 || stats.Pending > 0 || stats.NoText > 0 {
@pipelineBar(stats)
}
if stats.RateLimited+stats.Pending > 0 {
if (stats.RateLimited+stats.Pending) > 0 && autoSummarize {
<p class="pipeline-note muted">
Tapir fetches captions slowly on purpose, to respect YouTube's limits
new summaries land gradually. Check back tomorrow.
</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">
@summaryList(b, hasConnected)
@summaryList(b, hasConnected, autoSummarize)
</div>
}
}
@@ -168,7 +177,7 @@ templ pasteForm() {
<div id="paste-result"></div>
}
templ filterForm(f Filter) {
templ filterForm(f Filter, channels []string) {
<form
class="filters"
method="get"
@@ -178,7 +187,16 @@ templ filterForm(f Filter) {
hx-swap="innerHTML"
hx-indicator="#filter-indicator"
>
<label>Channel <input type="text" name="channel" value={ f.Channel } placeholder="any"/></label>
if len(channels) > 0 {
<label>
Channels
<select name="channel" multiple size="4">
for _, c := range channels {
<option value={ c } selected?={ f.HasChannel(c) }>{ c }</option>
}
</select>
</label>
}
<label class="filter-check">
<input type="checkbox" name="summarized" value="1" if f.OnlySummarized { checked }/>
Summarized only
@@ -194,12 +212,16 @@ templ filterForm(f Filter) {
// 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
// state instead of a blank table.
templ summaryList(b listBuckets, hasConnected bool) {
templ summaryList(b listBuckets, hasConnected bool, autoSummarize bool) {
if b.empty() {
if hasConnected {
<div class="empty empty-connected">
<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>
} else {
<div class="empty">
File diff suppressed because it is too large Load Diff
+6 -4
View File
@@ -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",
@@ -67,6 +68,7 @@ var scenarioCoverage = map[string]string{
"A subscribed channel posts a video with no usable transcript": "TestVideoWithNoTranscriptIsSkipped",
"A channel I am not subscribed to posts a video": "TestUnsubscribedChannelVideoIsNotProcessed",
"The same video is not summarized twice": "TestAlreadySummarizedVideoIsNotReprocessed",
"Re-analyzing a stored video does not re-fetch its transcript": "TestProcessNewVideo_SecondSummarizeDoesNotRefetch",
}
var (