Files
tapir/cmd/tapir/processor.go
T
mathiasandClaude Opus 4.8 71df696448
CI / Lint / Test / Vet (push) Successful in 10s
CI / Build & Import (push) Successful in 10s
feat(web): per-video chat over the stored transcript (ADR-027)
A deeper-dive chat entered from the summary view: ask questions about a video
against its already-stored transcript (ADR-021), no caption fetch, ever.

Safety by construction — the load-bearing property. The chat handlers reach the
chat.Service only after reading the SHARED stored transcript via Store.GetTranscript
(a pure DB read); the service holds no VideoSource. So an enabled chat cannot
trigger a caption fetch, touch the rate gate, or reach YouTube. A video with no
stored transcript gets an honest "not available" — no fetch, no model call. The
key web test wires the summarize/fetch collaborators as tripwires that fail the
test if chat ever routes into them, and asserts the model answered from the
stored text.

Model defaults to the summary's own model and is switchable among the ADR-022
chain (phi4-mini → gemma4-26b → mistral-small); switching re-runs against the
same transcript — deliberate model-comparison instrumentation. The cloud model
is absent from the switcher when TAPIR_CLOUD_FALLBACK_MODEL="" (the local-first /
NDA lever), honoured the same way the summarizer honours it. Reuses the existing
LiteLLM gateway client (a chat is a different call, not a new integration) and
the TAPIR_MAX_TRANSCRIPT_CHARS truncation, surfacing an honest bounded-context
note when a long transcript is cut.

Ephemeral v1: the multi-turn conversation rides in hidden request fields; no
table, no migration, nothing persisted. Entry is RLS-scoped through
GetSummaryByVideo, so chat is reachable only from the user's own summary view.
Show-source verification and on-demand fetch are deferred (ADR-027).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-11 09:26:58 +02:00

212 lines
8.3 KiB
Go

package main
import (
"context"
"fmt"
"strings"
"gitea.d-ma.be/mathias/tapir/internal/adapters/chat"
"gitea.d-ma.be/mathias/tapir/internal/adapters/llm"
"gitea.d-ma.be/mathias/tapir/internal/adapters/secrets"
"gitea.d-ma.be/mathias/tapir/internal/adapters/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/domain"
"gitea.d-ma.be/mathias/tapir/internal/ports"
"gitea.d-ma.be/mathias/tapir/internal/usecase"
"gitea.d-ma.be/mathias/tapir/internal/web"
)
// videoFetcher adapts the YouTube adapter to web.VideoFetcher for the paste flow
// (Feature 2). It builds a per-user adapter bound to that user's token ref and
// resolves a single video's metadata via the Data API — ungated; only the later
// transcript fetch goes through globalFetchGate.
type videoFetcher struct {
cfg config.Config
secrets ports.SecretStore
}
func (f videoFetcher) FetchVideo(ctx context.Context, userID, videoID string) (domain.Video, error) {
a := youtube.New(youtube.Config{
ClientID: f.cfg.YTClientID,
ClientSecret: f.cfg.YTClientSecret,
TokenSecretRef: web.YouTubeTokenRef(userID),
}, f.secrets)
return a.VideoByID(ctx, userID, videoID)
}
// buildProcessor wires the summarization engine — YouTube source (captions-first),
// AI-router summarizer, store sink — shared by `tapir run` and the web
// "Summarize now" path so the wiring lives in one place. It returns (nil, nil) —
// not an error — when the config cannot support live summarization (no gateway
// URL, no YouTube client credentials, or no secrets file). That nil is the
// 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)
}
// chatModels is the ordered, local-first set of models offered in the chat
// switcher (ADR-027), reusing the ADR-022 chain: primary → local fallback →
// cloud. Empty entries are dropped and duplicates collapsed, so a client/NDA
// deployment that sets the cloud fallback empty simply has no external model in
// the switcher — the same local-first lever the summarizer honours.
func chatModels(cfg config.Config) []string {
var models []string
add := func(m string) {
if m == "" {
return
}
for _, e := range models {
if e == m {
return
}
}
models = append(models, m)
}
add(cfg.SummarizerModel)
add(cfg.FallbackModel)
add(cfg.CloudFallbackModel)
return models
}
// buildChat wires the per-video chat service (ADR-027): a Completer factory over
// the SAME LiteLLM gateway the summarizer uses (a different alias per model, not a
// second client config) and the same transcript-truncation budget. It returns nil
// when no gateway is configured — chat is simply not mounted, the read path is
// unaffected. It deliberately takes NO YouTube source: chat is stored-only.
func buildChat(cfg config.Config) *chat.Service {
if cfg.GatewayURL == "" {
return nil
}
models := chatModels(cfg)
if len(models) == 0 {
return nil
}
newClient := func(model string) chat.Completer {
return llm.New(cfg.GatewayURL, cfg.GatewayKey, model, cfg.SummarizerTimeout, llm.WithMaxTokens(cfg.SummaryMaxTokens))
}
return chat.New(newClient, models, 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
}
secretStore := secrets.NewFileStore(cfg.SecretsFile)
src := youtube.New(youtube.Config{
ClientID: cfg.YTClientID,
ClientSecret: cfg.YTClientSecret,
TokenSecretRef: cfg.YTTokenRef,
PreferredLanguages: []string{"en"},
MinVideoSeconds: cfg.MinVideoSeconds,
}, secretStore)
sum := buildSummarizer(cfg)
// The store is both the summary sink and the shared transcript cache (ADR-021):
// the engine reads stored transcripts before any caption fetch and writes
// 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
// the web.Processor port (which works in terms of a stored video id): it loads the
// video row, runs the engine, and — on a produced summary — clears the manual
// queue flag, mirroring the runner so the video is not re-summarized on the next
// `tapir run` and the UI drops the "Queued" chip. A skip (no transcript) leaves
// the flag set so a later run can retry.
type engineProcessor struct {
engine *usecase.Engine
store *store.Store
}
func (p *engineProcessor) ProcessVideo(ctx context.Context, userID, videoID string) error {
// This is the user-initiated (foreground) path — a click on "Summarize",
// "Try now", or a pasted URL. Mark the context so the caption gate gives it
// priority over the background sweep (ADR-026, Pillar A).
ctx = youtube.ForegroundContext(ctx)
row, err := p.store.GetVideoRow(ctx, userID, videoID)
if err != nil {
return fmt.Errorf("load video %q: %w", videoID, err)
}
v := domain.Video{
ID: row.VideoID,
UserID: userID,
Provider: domain.Provider(row.Channel),
ProviderVideoID: row.ProviderVideoID,
Title: row.Title,
URL: row.URL,
PublishedAt: row.PublishedAt,
}
res, err := p.engine.ProcessNewVideo(ctx, v)
if err != nil {
return fmt.Errorf("process video %q: %w", videoID, err)
}
// Record the outcome so the status endpoint can show honest state (ADR-025):
// a 429'd or caption-less click used to leave transcript_status unset, so the
// poll silently reverted to the "Summarize" button. Mirror the runner: stamp
// rate_limited / none / fetched. A rate-limited video keeps its requested flag
// so the background sweep retries it; none and fetched are terminal here.
switch {
case res.Skipped && res.TranscriptSource == string(domain.SourceRateLimited):
if err := p.store.SetTranscriptStatus(ctx, userID, videoID, "rate_limited"); err != nil {
return fmt.Errorf("set rate_limited status %q: %w", videoID, err)
}
case res.Skipped:
if err := p.store.SetTranscriptStatus(ctx, userID, videoID, "none"); err != nil {
return fmt.Errorf("set none status %q: %w", videoID, err)
}
if err := p.store.ClearSummarizeRequested(ctx, userID, videoID); err != nil {
return fmt.Errorf("clear summarize flag %q: %w", videoID, err)
}
case res.Summary != nil:
if err := p.store.SetTranscriptStatus(ctx, userID, videoID, "fetched"); err != nil {
return fmt.Errorf("set fetched status %q: %w", videoID, err)
}
if err := p.store.ClearSummarizeRequested(ctx, userID, videoID); err != nil {
return fmt.Errorf("clear summarize flag %q: %w", videoID, err)
}
}
return nil
}