Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0ba78e8868 | ||
|
|
821d5f99cd | ||
|
|
5c70408e75 | ||
|
|
cb6917ca59 | ||
|
|
099b2d4c68 | ||
|
|
f66c1bcdcc | ||
|
|
1e65c3b413 |
@@ -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.)
|
(See `DECISIONS.md` for full rationale. Listed here so you don't propose them.)
|
||||||
|
|
||||||
- **No Supabase** — reuse Dex / ESO+1Password / Postgres (ADR-002).
|
- **No Supabase** — reuse Dex / ESO+1Password / Postgres (ADR-002).
|
||||||
- **No global cross-tenant video/transcript table** — per-user isolation (data-model). Dedup
|
- **No global cross-tenant *video* table** — videos stay per-user (data-model). Transcripts ARE
|
||||||
across users is a Future C concern, not a Stage 0/1 default.
|
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
|
- **No audio-download + speech-to-text in the core path** — captions-first (ADR-007). STT is a
|
||||||
deferred, bounded optional component.
|
deferred, bounded optional component.
|
||||||
- **No public SaaS / sign-up / billing / Google OAuth verification at scale** — Future C,
|
- **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)
|
## 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.14.0`. `task check` passes (fmt, vet, lint,
|
||||||
`go test -p 1 ./...`). Go is `1.26.1` (see `go.mod`).
|
`go test -p 1 ./...`). Go is `1.26.1` (see `go.mod`).
|
||||||
|
|
||||||
- Clean Architecture core is implemented: `internal/domain` (entities), `internal/ports`
|
- Clean Architecture core is implemented: `internal/domain` (entities), `internal/ports`
|
||||||
@@ -94,13 +95,18 @@ The repo is **green and shipping** — last tag `v0.9.0`. `task check` passes (f
|
|||||||
green against it.
|
green against it.
|
||||||
- Adapters present under `internal/adapters/`: `youtube` (captions-first `VideoSource`,
|
- Adapters present under `internal/adapters/`: `youtube` (captions-first `VideoSource`,
|
||||||
timedtext/InnerTube acquisition per ADR-010), `summarizer` + `llm` (the copied AI router,
|
timedtext/InnerTube acquisition per ADR-010), `summarizer` + `llm` (the copied AI router,
|
||||||
Primary→Fallback per ADR-004), `store` (Postgres, golang-migrate migrations 001–006),
|
Primary→Fallback per ADR-004), `store` (Postgres, golang-migrate migrations 001–015),
|
||||||
`secrets` (file-backed `SecretStore`). The brain HTTP sink (ADR-005) is the remaining
|
`secrets` (file-backed `SecretStore`). The brain HTTP sink (ADR-005) is the remaining
|
||||||
optional sink.
|
optional sink.
|
||||||
- Stage 1 is open (ADR-012): multi-user with **DB-enforced** isolation — Postgres RLS `FORCE`d
|
- Stage 1 is open (ADR-012): multi-user with **DB-enforced** isolation — Postgres RLS `FORCE`d
|
||||||
on all user-owned tables (migration 003), two-user isolation test in
|
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
|
`internal/adapters/store/rls_test.go`. Registration gate, per-user YouTube web connect, and
|
||||||
account management (disconnect / delete, ADR-013) all shipped.
|
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
|
- `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
|
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.
|
transport over the unchanged engine/ports — ADR-003). `tapir env` prints config.
|
||||||
|
|||||||
+61
-1
@@ -743,6 +743,66 @@ 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 1–5 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.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
## Rejected alternatives
|
## Rejected alternatives
|
||||||
|
|
||||||
Approaches considered during the 2026-06-02 planning + grill session and **deliberately not
|
Approaches considered during the 2026-06-02 planning + grill session and **deliberately not
|
||||||
@@ -757,7 +817,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 |
|
| 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 |
|
| 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 |
|
| 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 1–5 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 1–5 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 |
|
| 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 |
|
| 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) |
|
| 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) |
|
||||||
|
|||||||
@@ -63,7 +63,12 @@ func buildProcessor(cfg config.Config, st *store.Store) (*usecase.Engine, error)
|
|||||||
}
|
}
|
||||||
sum := summarizer.New(primary, nil)
|
sum := summarizer.New(primary, nil)
|
||||||
|
|
||||||
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
|
// engineProcessor adapts the engine (which works in terms of a domain.Video) to
|
||||||
|
|||||||
+21
-13
@@ -10,12 +10,15 @@ only opaque references to them; the secret material lives in ESO/1Password (ADR-
|
|||||||
|
|
||||||
## Design decisions baked into this model
|
## Design decisions baked into this model
|
||||||
|
|
||||||
- **Per-user isolation, not a shared global video table.** The earlier draft proposed a
|
- **Per-user isolation for everything except transcripts.** The earlier draft proposed a
|
||||||
global `videos`/`transcripts` table deduped across tenants. Rejected for Future B: it
|
global `videos`/`transcripts` table deduped across tenants. **Videos** stay per-user and
|
||||||
reintroduces exactly the cross-domain coupling the homelab architecture review is
|
RLS-scoped — a shared video table reintroduces exactly the cross-domain coupling the homelab
|
||||||
removing, and at 1–5 users the cost of occasionally re-summarizing the same video is
|
architecture review is removing. **Transcripts**, however, are now shared (ADR-021): keyed by
|
||||||
trivial compared to the isolation it would cost. Each user's data is self-contained.
|
`(provider, provider_video_id)`, no `user_id`, **not** RLS-scoped. The cost avoided there is
|
||||||
(Revisit only if Future C makes GPU/transcription cost dominate — a new ADR, not a default.)
|
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
|
- **Secrets by reference only.** Tables hold a `secret_ref` (opaque string/UUID resolved via
|
||||||
the `SecretStore` port), never tokens or keys.
|
the `SecretStore` port), never tokens or keys.
|
||||||
- **The brain sink is just a delivery target.** No brain-specific tables. Whether a summary
|
- **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)"
|
USER ||--o{ AI_CREDENTIAL : "has (planned)"
|
||||||
VIDEO_CONNECTION ||--o{ SUBSCRIPTION : "exposes (planned)"
|
VIDEO_CONNECTION ||--o{ SUBSCRIPTION : "exposes (planned)"
|
||||||
SUBSCRIPTION ||--o{ VIDEO : "produces (per user)"
|
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"
|
VIDEO ||--o| SUMMARY : "has at most one"
|
||||||
SUMMARY ||--o{ SINK_DELIVERY : "delivered via"
|
SUMMARY ||--o{ SINK_DELIVERY : "delivered via"
|
||||||
USER ||--o{ CHANNEL_ERROR : "reports unavailable channels"
|
USER ||--o{ CHANNEL_ERROR : "reports unavailable channels"
|
||||||
@@ -92,12 +95,12 @@ erDiagram
|
|||||||
timestamptz rate_limited_at "backoff clock for 429 retries (migration 007)"
|
timestamptz rate_limited_at "backoff clock for 429 retries (migration 007)"
|
||||||
}
|
}
|
||||||
TRANSCRIPT {
|
TRANSCRIPT {
|
||||||
uuid video_id PK_FK
|
text provider PK "part of shared key (ADR-021)"
|
||||||
uuid user_id FK
|
text provider_video_id PK "part of shared key — the cross-user dedup key"
|
||||||
text source "captions | none"
|
text source "captions | none"
|
||||||
text language
|
text language
|
||||||
text content "null when source = none"
|
text content "null when source = none"
|
||||||
timestamptz resolved_at
|
timestamptz fetched_at
|
||||||
}
|
}
|
||||||
SUMMARY {
|
SUMMARY {
|
||||||
uuid id PK
|
uuid id PK
|
||||||
@@ -173,8 +176,12 @@ mechanism.
|
|||||||
`transcript_status` and `rate_limited_at` (migration 007) track caption-fetch outcomes for
|
`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
|
rate-limit backoff: `NULL` = not attempted; `rate_limited` = 429 seen, skip until
|
||||||
`NOW() - rate_limited_at > TAPIR_FETCH_BACKOFF`; `fetched` = resolved; `none` = no transcript.
|
`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** — shared public caption content, one row per `(provider, provider_video_id)`,
|
||||||
transcript" so the watcher doesn't reprocess (ADR-007). `content` null in that case.
|
**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
|
- **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`/
|
"is local good enough?" question queryable (the Stage 0 quality signal). `highlights`/
|
||||||
`takeaways` as jsonb to stay schema-flexible while the output format settles.
|
`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)
|
## 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.
|
- Sharding / per-tenant physical databases.
|
||||||
- Soft-delete + full audit trail on connections/credentials (a Stage 2 hardening item; add
|
- Soft-delete + full audit trail on connections/credentials (a Stage 2 hardening item; add
|
||||||
via ADR when Stage 2 work starts).
|
via ADR when Stage 2 work starts).
|
||||||
|
|||||||
@@ -31,5 +31,16 @@ Feature: Summarize new videos from subscribed channels
|
|||||||
When the watcher sees "Designing for Attention" again
|
When the watcher sees "Designing for Attention" again
|
||||||
Then Tapir does not produce a second summary for it
|
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
|
# Captions-first is the core path (ADR-007). Audio-download + speech-to-text is
|
||||||
# deferred and intentionally has no scenario here yet.
|
# 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.
|
||||||
|
|||||||
@@ -10,10 +10,13 @@ import (
|
|||||||
// DeleteUser permanently removes a user and all of their data. It runs through
|
// 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.
|
// withUser so RLS confines every statement to the calling user's own rows.
|
||||||
//
|
//
|
||||||
// Deleting the users row cascades (ON DELETE CASCADE) to videos, transcripts,
|
// Deleting the users row cascades (ON DELETE CASCADE) to videos, summaries
|
||||||
// summaries (→ sink_deliveries), video_connections, and the user_identities map
|
// (→ sink_deliveries), video_connections, and the user_identities map —
|
||||||
// — referential-integrity cascades bypass RLS, so a user's child rows are removed
|
// 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
|
// 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
|
// 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
|
// (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
|
// explicitly in the same scoped transaction. Deleting an absent user is a no-op
|
||||||
|
|||||||
@@ -53,7 +53,11 @@ func TestMigration010LoginEventsUpDown(t *testing.T) {
|
|||||||
require.True(t, loginEventsExists(t), "login_events must exist at latest migration")
|
require.True(t, loginEventsExists(t), "login_events must exist at latest migration")
|
||||||
|
|
||||||
m := fileMigrator(t)
|
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.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.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")
|
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.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.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")
|
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")
|
require.Equal(t, "true", autoSummarizeDefault(t), "011 sets the default to TRUE")
|
||||||
|
|
||||||
m := fileMigrator(t)
|
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 013 drops channel_errors")
|
||||||
require.NoError(t, m.Steps(-1), "down 012 is a no-op")
|
require.NoError(t, m.Steps(-1), "down 012 is a no-op")
|
||||||
require.NoError(t, m.Steps(-1), "down 011 reverts the column default")
|
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.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 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 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
|
// 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).';
|
||||||
@@ -29,6 +29,7 @@ type SummaryRow struct {
|
|||||||
ProviderVideoID string // videos.provider_video_id; empty when no videos row
|
ProviderVideoID string // videos.provider_video_id; empty when no videos row
|
||||||
Title string // videos.title; 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
|
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
|
URL string // videos.url; empty when no videos row
|
||||||
PublishedAt time.Time // videos.published_at; zero when absent
|
PublishedAt time.Time // videos.published_at; zero when absent
|
||||||
Summary string
|
Summary string
|
||||||
@@ -137,7 +138,8 @@ const selectVideo = `
|
|||||||
COALESCE(s.created_at, v.seen_at),
|
COALESCE(s.created_at, v.seen_at),
|
||||||
(s.id IS NOT NULL) AS summarized,
|
(s.id IS NOT NULL) AS summarized,
|
||||||
v.summarize_requested,
|
v.summarize_requested,
|
||||||
COALESCE(v.transcript_status, '')
|
COALESCE(v.transcript_status, ''),
|
||||||
|
COALESCE(v.channel_title, '')
|
||||||
FROM videos v
|
FROM videos v
|
||||||
LEFT JOIN summaries s ON s.video_id = v.id AND s.user_id = v.user_id`
|
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.Summarized,
|
||||||
&row.SummarizeRequested,
|
&row.SummarizeRequested,
|
||||||
&row.TranscriptStatus,
|
&row.TranscriptStatus,
|
||||||
|
&row.ChannelTitle,
|
||||||
); err != nil {
|
); err != nil {
|
||||||
return SummaryRow{}, fmt.Errorf("store: scan video: %w", err)
|
return SummaryRow{}, fmt.Errorf("store: scan video: %w", err)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -22,9 +22,11 @@ import (
|
|||||||
// (no GUC set → zero rows) proves the enforcement path is live, not bypassed.
|
// (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
|
// 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{
|
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
|
// allIsolatedTables adds sink_deliveries, whose ownership is derived from its
|
||||||
@@ -38,9 +40,10 @@ type seeded struct {
|
|||||||
summaryID string
|
summaryID string
|
||||||
}
|
}
|
||||||
|
|
||||||
// seedUser inserts one full chain (user → video → transcript → summary →
|
// seedUser inserts one full chain (user → video → summary → action → delivery)
|
||||||
// action → delivery) as the superuser pool, which bypasses RLS so both users'
|
// as the superuser pool, which bypasses RLS so both users' data lands regardless
|
||||||
// data lands regardless of the GUC.
|
// 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 {
|
func seedUser(t *testing.T, p *pgxpool.Pool, userID string) seeded {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
ctx := context.Background()
|
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`,
|
VALUES ($1, 'youtube', $2, 'title') RETURNING id`,
|
||||||
userID, "vid-"+userID).Scan(&videoID))
|
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
|
var summaryID string
|
||||||
require.NoError(t, p.QueryRow(ctx,
|
require.NoError(t, p.QueryRow(ctx,
|
||||||
`INSERT INTO summaries (user_id, video_id, summary) VALUES ($1, $2, 'sum')
|
`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()
|
t.Helper()
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
|
|
||||||
// Idempotent across test runs (schema/role persist for the TestMain PG).
|
// Idempotent across tests AND runs: the role persists for the TestMain PG and
|
||||||
_, _ = super.Exec(ctx, `DROP ROLE IF EXISTS app`)
|
// owns granted privileges, so a plain DROP ROLE fails once any GRANT exists
|
||||||
_, err := super.Exec(ctx, `CREATE ROLE app LOGIN PASSWORD 'app'`)
|
// (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)
|
require.NoError(t, err)
|
||||||
_, err = super.Exec(ctx, `GRANT USAGE ON SCHEMA public TO app`)
|
_, err = super.Exec(ctx, `GRANT USAGE ON SCHEMA public TO app`)
|
||||||
require.NoError(t, err)
|
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 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},
|
{"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},
|
{"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 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 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},
|
{"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
|
_ = 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
|
||||||
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
@@ -46,14 +46,15 @@ func (s *Store) UpsertVideo(ctx context.Context, v domain.Video) (string, error)
|
|||||||
}
|
}
|
||||||
|
|
||||||
if err := tx.QueryRow(ctx,
|
if err := tx.QueryRow(ctx,
|
||||||
`INSERT INTO videos (user_id, provider, provider_video_id, title, url, published_at)
|
`INSERT INTO videos (user_id, provider, provider_video_id, title, url, published_at, channel_title)
|
||||||
VALUES ($1, $2, $3, $4, $5, $6)
|
VALUES ($1, $2, $3, $4, $5, $6, $7)
|
||||||
ON CONFLICT (user_id, provider, provider_video_id) DO UPDATE SET
|
ON CONFLICT (user_id, provider, provider_video_id) DO UPDATE SET
|
||||||
title = EXCLUDED.title,
|
title = EXCLUDED.title,
|
||||||
url = EXCLUDED.url,
|
url = EXCLUDED.url,
|
||||||
published_at = EXCLUDED.published_at
|
published_at = EXCLUDED.published_at,
|
||||||
|
channel_title = COALESCE(NULLIF(EXCLUDED.channel_title, ''), videos.channel_title)
|
||||||
RETURNING id`,
|
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 {
|
).Scan(&id); err != nil {
|
||||||
return fmt.Errorf("store: upsert video: %w", err)
|
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
|
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
|
||||||
|
}
|
||||||
|
|||||||
@@ -113,3 +113,29 @@ func TestNewestUnsummarizedVideoIDs(t *testing.T) {
|
|||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
require.Empty(t, none, "limit 0 returns nothing")
|
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)")
|
||||||
|
}
|
||||||
|
|||||||
@@ -21,7 +21,7 @@ func TestVideoByID(t *testing.T) {
|
|||||||
if got := r.URL.Query().Get("part"); got != "snippet" {
|
if got := r.URL.Query().Get("part"); got != "snippet" {
|
||||||
t.Errorf("expected part=snippet, got %q", got)
|
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)
|
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" {
|
if v.ProviderVideoID != id || v.Title != "Never Gonna Give You Up" {
|
||||||
t.Errorf("unexpected video: %+v", v)
|
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 {
|
if v.Provider != domain.ProviderYouTube || v.URL != "https://www.youtube.com/watch?v="+id {
|
||||||
t.Errorf("video not wired correctly: %+v", v)
|
t.Errorf("video not wired correctly: %+v", v)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -234,6 +234,7 @@ func (a *Adapter) NewVideos(ctx context.Context, sub domain.Subscription) ([]dom
|
|||||||
Provider: domain.ProviderYouTube,
|
Provider: domain.ProviderYouTube,
|
||||||
ProviderVideoID: vid,
|
ProviderVideoID: vid,
|
||||||
Title: item.Snippet.Title,
|
Title: item.Snippet.Title,
|
||||||
|
ChannelTitle: sub.ChannelTitle,
|
||||||
URL: "https://www.youtube.com/watch?v=" + vid,
|
URL: "https://www.youtube.com/watch?v=" + vid,
|
||||||
PublishedAt: item.Snippet.PublishedAt,
|
PublishedAt: item.Snippet.PublishedAt,
|
||||||
})
|
})
|
||||||
@@ -270,6 +271,7 @@ func (a *Adapter) VideoByID(ctx context.Context, userID, videoID string) (domain
|
|||||||
Provider: domain.ProviderYouTube,
|
Provider: domain.ProviderYouTube,
|
||||||
ProviderVideoID: videoID,
|
ProviderVideoID: videoID,
|
||||||
Title: it.Snippet.Title,
|
Title: it.Snippet.Title,
|
||||||
|
ChannelTitle: it.Snippet.ChannelTitle,
|
||||||
URL: "https://www.youtube.com/watch?v=" + videoID,
|
URL: "https://www.youtube.com/watch?v=" + videoID,
|
||||||
PublishedAt: it.Snippet.PublishedAt,
|
PublishedAt: it.Snippet.PublishedAt,
|
||||||
}, nil
|
}, nil
|
||||||
@@ -376,8 +378,9 @@ type playlistItemListResponse struct {
|
|||||||
type videoListResponse struct {
|
type videoListResponse struct {
|
||||||
Items []struct {
|
Items []struct {
|
||||||
Snippet struct {
|
Snippet struct {
|
||||||
Title string `json:"title"`
|
Title string `json:"title"`
|
||||||
PublishedAt time.Time `json:"publishedAt"`
|
ChannelTitle string `json:"channelTitle"`
|
||||||
|
PublishedAt time.Time `json:"publishedAt"`
|
||||||
} `json:"snippet"`
|
} `json:"snippet"`
|
||||||
} `json:"items"`
|
} `json:"items"`
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
vids, err := a.NewVideos(context.Background(), sub)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("NewVideos: %v", err)
|
t.Fatalf("NewVideos: %v", err)
|
||||||
@@ -143,6 +143,9 @@ func TestNewVideos(t *testing.T) {
|
|||||||
if v.ProviderVideoID != "vid1" || v.Title != "Designing for Attention" {
|
if v.ProviderVideoID != "vid1" || v.Title != "Designing for Attention" {
|
||||||
t.Errorf("unexpected video: %+v", v)
|
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" {
|
if v.Provider != domain.ProviderYouTube || v.URL != "https://www.youtube.com/watch?v=vid1" {
|
||||||
t.Errorf("video not wired correctly: %+v", v)
|
t.Errorf("video not wired correctly: %+v", v)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -73,6 +73,7 @@ type Video struct {
|
|||||||
Provider Provider
|
Provider Provider
|
||||||
ProviderVideoID string
|
ProviderVideoID string
|
||||||
Title string
|
Title string
|
||||||
|
ChannelTitle string
|
||||||
URL string
|
URL string
|
||||||
PublishedAt time.Time
|
PublishedAt time.Time
|
||||||
SeenAt time.Time
|
SeenAt time.Time
|
||||||
|
|||||||
@@ -27,6 +27,26 @@ type Summarizer interface {
|
|||||||
Summarize(ctx context.Context, v domain.Video, t domain.Transcript) (domain.Summary, error)
|
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, ...).
|
// Sink delivers a summary to a destination (user store, brain, ...).
|
||||||
// Implementations fail independently of one another.
|
// Implementations fail independently of one another.
|
||||||
type Sink interface {
|
type Sink interface {
|
||||||
|
|||||||
@@ -27,6 +27,13 @@ type Engine struct {
|
|||||||
AI ports.Summarizer
|
AI ports.Summarizer
|
||||||
Sinks []ports.Sink
|
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
|
// processed dedups videos within this engine's lifetime so a video is not
|
||||||
// summarized twice when the watcher sees it again. Durable cross-restart
|
// summarized twice when the watcher sees it again. Durable cross-restart
|
||||||
// dedup is the store's concern (a resolved TRANSCRIPT / existing SUMMARY,
|
// dedup is the store's concern (a resolved TRANSCRIPT / existing SUMMARY,
|
||||||
@@ -57,9 +64,9 @@ type ProcessResult struct {
|
|||||||
// resolve transcript -> (summarize -> deliver) | skip.
|
// resolve transcript -> (summarize -> deliver) | skip.
|
||||||
// See docs/use-cases/summarize_new_video.feature.
|
// See docs/use-cases/summarize_new_video.feature.
|
||||||
func (e *Engine) ProcessNewVideo(ctx context.Context, v domain.Video) (ProcessResult, error) {
|
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 {
|
if err != nil {
|
||||||
return ProcessResult{Video: v}, fmt.Errorf("fetch transcript: %w", err)
|
return ProcessResult{Video: v}, err
|
||||||
}
|
}
|
||||||
if !t.HasText() {
|
if !t.HasText() {
|
||||||
// No usable transcript: record the skip, produce no summary, deliver nothing
|
// 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...)
|
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
|
// ProcessNewVideos walks a user's subscriptions and processes each newly seen
|
||||||
// video. Only videos surfaced via the user's subscriptions are considered, so a
|
// 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
|
// channel the user is not subscribed to is never processed. A video already
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
+22
-14
@@ -21,6 +21,9 @@ import (
|
|||||||
// fake without a database.
|
// fake without a database.
|
||||||
type Store interface {
|
type Store interface {
|
||||||
ListVideos(ctx context.Context, userID string, limit int) ([]store.SummaryRow, error)
|
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)
|
GetSummaryByVideo(ctx context.Context, userID, videoID string) (*store.SummaryRow, error)
|
||||||
GetVideoRow(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)
|
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()
|
q := r.URL.Query()
|
||||||
f := Filter{
|
f := Filter{
|
||||||
Channel: q.Get("channel"),
|
Channels: nonEmptyStrings(q["channel"]),
|
||||||
From: q.Get("from"),
|
From: q.Get("from"),
|
||||||
To: q.Get("to"),
|
To: q.Get("to"),
|
||||||
OnlySummarized: q.Get("summarized") == "1",
|
OnlySummarized: q.Get("summarized") == "1",
|
||||||
@@ -211,25 +214,30 @@ func (a *App) handleList(w http.ResponseWriter, r *http.Request) {
|
|||||||
rows := f.apply(allRows)
|
rows := f.apply(allRows)
|
||||||
buckets := bucketRows(rows, a.recencyCutoff())
|
buckets := bucketRows(rows, a.recencyCutoff())
|
||||||
|
|
||||||
// hasConnected drives the empty state: a fresh account with a connection but
|
// Channel options for the multi-select filter (the user's source channels).
|
||||||
// no discovery pass yet has zero rows, and we want it to read "connected,
|
channels, err := a.Store.DistinctChannels(r.Context(), userID)
|
||||||
// summaries land gradually" rather than "nothing here". Only needed when the
|
if err != nil {
|
||||||
// list is empty.
|
a.serverError(w, r, "distinct channels", err)
|
||||||
hasConnected := false
|
return
|
||||||
if buckets.empty() {
|
|
||||||
conns, err := a.Store.ConnectionsForUser(r.Context(), userID)
|
|
||||||
if err != nil {
|
|
||||||
a.serverError(w, r, "connections for user", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
hasConnected = len(conns) > 0
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 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").
|
||||||
|
// Computed every render — not only when empty — so a user with videos still
|
||||||
|
// gets the paste box.
|
||||||
|
conns, err := a.Store.ConnectionsForUser(r.Context(), userID)
|
||||||
|
if err != nil {
|
||||||
|
a.serverError(w, r, "connections for user", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
hasConnected := len(conns) > 0
|
||||||
|
|
||||||
if isHTMX(r) {
|
if isHTMX(r) {
|
||||||
a.render(w, r, summaryList(buckets, hasConnected))
|
a.render(w, r, summaryList(buckets, hasConnected))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
a.render(w, r, ListPage(buckets, f, stats, takeFlash(w, r), hasConnected))
|
a.render(w, r, ListPage(buckets, f, stats, takeFlash(w, r), hasConnected, channels))
|
||||||
}
|
}
|
||||||
|
|
||||||
// handleDetail renders one summary in full (highlights, takeaways, action group).
|
// handleDetail renders one summary in full (highlights, takeaways, action group).
|
||||||
|
|||||||
@@ -207,13 +207,15 @@ func TestListChannelFilter(t *testing.T) {
|
|||||||
resetDB(t, p)
|
resetDB(t, p)
|
||||||
require.NoError(t, deliver(ctx, app, videoX, "body x"))
|
require.NoError(t, deliver(ctx, app, videoX, "body x"))
|
||||||
seedVideo(t, p, videoX, "X Title", "https://x", time.Time{})
|
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.
|
// Selecting a different channel hides the row; selecting its channel shows it.
|
||||||
rec := do(t, app, httptest.NewRequest(http.MethodGet, "/?channel=vimeo", nil))
|
rec := do(t, app, httptest.NewRequest(http.MethodGet, "/?channel=Other+Channel", nil))
|
||||||
require.Equal(t, http.StatusOK, rec.Code)
|
require.Equal(t, http.StatusOK, rec.Code)
|
||||||
require.NotContains(t, body(t, rec), "X Title")
|
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")
|
require.Contains(t, body(t, rec), "X Title")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ import (
|
|||||||
"net/url"
|
"net/url"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
"github.com/stretchr/testify/require"
|
"github.com/stretchr/testify/require"
|
||||||
|
|
||||||
@@ -111,3 +112,23 @@ func TestPasteDedupNoDuplicate(t *testing.T) {
|
|||||||
userID).Scan(&count))
|
userID).Scan(&count))
|
||||||
require.Equal(t, 1, count, "pasting the same video twice must not duplicate the row")
|
require.Equal(t, 1, count, "pasting the same video twice must not duplicate the row")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestListShowsPasteFormForConnectedUserWithVideos(t *testing.T) {
|
||||||
|
app := newApp(t)
|
||||||
|
resetDB(t, rawPool(t))
|
||||||
|
app.Fetcher = &fakeFetcher{title: "x"}
|
||||||
|
p := rawPool(t)
|
||||||
|
|
||||||
|
// Connected user with a non-empty feed (the case the bug missed: hasConnected
|
||||||
|
// was only computed for an empty feed).
|
||||||
|
_, err := p.Exec(context.Background(),
|
||||||
|
`INSERT INTO video_connections (user_id, provider, token_ref, status)
|
||||||
|
VALUES ($1, 'youtube', 'youtube/x/refresh_token', 'active')`, userID)
|
||||||
|
require.NoError(t, err)
|
||||||
|
seedVideo(t, p, "11111111-1111-1111-1111-111111111111", "A talk", "https://youtu.be/aaaaaaaaaaa", time.Now())
|
||||||
|
|
||||||
|
rec := do(t, app, httptest.NewRequest(http.MethodGet, "/", nil))
|
||||||
|
require.Equal(t, http.StatusOK, rec.Code)
|
||||||
|
require.Contains(t, body(t, rec), `action="/paste"`,
|
||||||
|
"a connected user must see the paste box even when the feed has videos")
|
||||||
|
}
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package web
|
|||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
|
"gitea.d-ma.be/mathias/tapir/internal/adapters/store"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
)
|
)
|
||||||
@@ -64,7 +65,7 @@ func TestParseYouTubeVideoID(t *testing.T) {
|
|||||||
func TestListPageShowsPasteFormOnlyWhenConnected(t *testing.T) {
|
func TestListPageShowsPasteFormOnlyWhenConnected(t *testing.T) {
|
||||||
render := func(connected bool) string {
|
render := func(connected bool) string {
|
||||||
var buf bytes.Buffer
|
var buf bytes.Buffer
|
||||||
if err := ListPage(listBuckets{}, Filter{}, PipelineStats{}, "", connected).Render(context.Background(), &buf); err != nil {
|
if err := ListPage(listBuckets{}, Filter{}, PipelineStats{}, "", connected, nil).Render(context.Background(), &buf); err != nil {
|
||||||
t.Fatalf("render: %v", err)
|
t.Fatalf("render: %v", err)
|
||||||
}
|
}
|
||||||
return buf.String()
|
return buf.String()
|
||||||
@@ -78,3 +79,19 @@ func TestListPageShowsPasteFormOnlyWhenConnected(t *testing.T) {
|
|||||||
t.Errorf("disconnected feed must not show the paste form")
|
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
@@ -2,6 +2,7 @@ package web
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"regexp"
|
"regexp"
|
||||||
|
"slices"
|
||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
"unicode/utf8"
|
"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
|
// Dates are kept as the raw YYYY-MM-DD strings so the form re-renders the user's
|
||||||
// input verbatim; parsing happens in matchFilter.
|
// input verbatim; parsing happens in matchFilter.
|
||||||
type Filter struct {
|
type Filter struct {
|
||||||
Channel string
|
Channels []string // selected channel titles; empty = all channels
|
||||||
From string
|
From string
|
||||||
To string
|
To string
|
||||||
OnlySummarized bool // show only videos that have a summary
|
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) 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.
|
// filter that happens to match nothing still shows the bar so it can be cleared.
|
||||||
func (f Filter) active() bool {
|
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;
|
// 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 {
|
if f.OnlySummarized && !r.Summarized {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
if f.Channel != "" && r.Channel != f.Channel {
|
if len(f.Channels) > 0 && !slices.Contains(f.Channels, r.ChannelTitle) {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
if from, ok := parseDate(f.From); ok {
|
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.
|
// apply returns the subset of rows matching the filter, preserving order.
|
||||||
func (f Filter) apply(rows []store.SummaryRow) []store.SummaryRow {
|
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
|
return rows
|
||||||
}
|
}
|
||||||
out := rows[:0:0]
|
out := rows[:0:0]
|
||||||
|
|||||||
@@ -104,14 +104,14 @@ templ flashBanner(code string) {
|
|||||||
// #summary-list region; a non-HTMX request renders the whole page. flash carries
|
// #summary-list region; a non-HTMX request renders the whole page. flash carries
|
||||||
// a one-shot notification (e.g. "connected", "registered") surfaced on arrival
|
// a one-shot notification (e.g. "connected", "registered") surfaced on arrival
|
||||||
// after a POST→redirect.
|
// after a POST→redirect.
|
||||||
templ ListPage(b listBuckets, f Filter, stats PipelineStats, flash string, hasConnected bool) {
|
templ ListPage(b listBuckets, f Filter, stats PipelineStats, flash string, hasConnected bool, channels []string) {
|
||||||
@Layout("Tapir — Summaries") {
|
@Layout("Tapir — Summaries") {
|
||||||
@flashBanner(flash)
|
@flashBanner(flash)
|
||||||
if hasConnected {
|
if hasConnected {
|
||||||
@pasteForm()
|
@pasteForm()
|
||||||
}
|
}
|
||||||
if !b.empty() || f.active() {
|
if !b.empty() || f.active() {
|
||||||
@filterForm(f)
|
@filterForm(f, channels)
|
||||||
}
|
}
|
||||||
if stats.RateLimited > 0 || stats.Pending > 0 || stats.NoText > 0 {
|
if stats.RateLimited > 0 || stats.Pending > 0 || stats.NoText > 0 {
|
||||||
@pipelineBar(stats)
|
@pipelineBar(stats)
|
||||||
@@ -168,7 +168,7 @@ templ pasteForm() {
|
|||||||
<div id="paste-result"></div>
|
<div id="paste-result"></div>
|
||||||
}
|
}
|
||||||
|
|
||||||
templ filterForm(f Filter) {
|
templ filterForm(f Filter, channels []string) {
|
||||||
<form
|
<form
|
||||||
class="filters"
|
class="filters"
|
||||||
method="get"
|
method="get"
|
||||||
@@ -178,7 +178,16 @@ templ filterForm(f Filter) {
|
|||||||
hx-swap="innerHTML"
|
hx-swap="innerHTML"
|
||||||
hx-indicator="#filter-indicator"
|
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">
|
<label class="filter-check">
|
||||||
<input type="checkbox" name="summarized" value="1" if f.OnlySummarized { checked }/>
|
<input type="checkbox" name="summarized" value="1" if f.OnlySummarized { checked }/>
|
||||||
Summarized only
|
Summarized only
|
||||||
|
|||||||
+597
-554
File diff suppressed because it is too large
Load Diff
@@ -67,6 +67,7 @@ var scenarioCoverage = map[string]string{
|
|||||||
"A subscribed channel posts a video with no usable transcript": "TestVideoWithNoTranscriptIsSkipped",
|
"A subscribed channel posts a video with no usable transcript": "TestVideoWithNoTranscriptIsSkipped",
|
||||||
"A channel I am not subscribed to posts a video": "TestUnsubscribedChannelVideoIsNotProcessed",
|
"A channel I am not subscribed to posts a video": "TestUnsubscribedChannelVideoIsNotProcessed",
|
||||||
"The same video is not summarized twice": "TestAlreadySummarizedVideoIsNotReprocessed",
|
"The same video is not summarized twice": "TestAlreadySummarizedVideoIsNotReprocessed",
|
||||||
|
"Re-analyzing a stored video does not re-fetch its transcript": "TestProcessNewVideo_SecondSummarizeDoesNotRefetch",
|
||||||
}
|
}
|
||||||
|
|
||||||
var (
|
var (
|
||||||
|
|||||||
Reference in New Issue
Block a user