The first friendly-pilot live run produced zero summaries: koala/phi4-mini hit three silent failure modes — 8k context overflow on long transcripts (HTTP 400), intermittent malformed JSON (highlights as a bare string), and no fallback wired at all (summarizer.New(primary, nil)). Keep phi4-mini as the fast primary and add resilience around it: - Ordered endpoint chain (summarizer.NewChain): phi4-mini → koala/phi4-14b (local) → berget/mistral-small (worst-case external). All reached through the one LiteLLM gateway by alias. - A parse failure now advances the chain like a transport error — the old Primary→Fallback shape returned the parse error without trying anyone else. - Tolerant parse: highlights/takeaways coerce string→[]string, absorbing the common small-model quirk without spending a fallback round-trip. - Transcript truncation (TAPIR_MAX_TRANSCRIPT_CHARS=18000) prevents the overflow rather than recovering from it; validated to fit phi4-mini's 8k window. - Bounded completion budget (TAPIR_SUMMARY_MAX_TOKENS=1500) — the old 8192 budget itself contributed to the overflow. Local-first guarantee preserved by ordering: external endpoint is tried only after every local one fails. TAPIR_CLOUD_FALLBACK_MODEL="" disables it entirely for client/NDA deployments. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
141 lines
5.4 KiB
Go
141 lines
5.4 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"strings"
|
|
|
|
"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)
|
|
}
|
|
|
|
// 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"},
|
|
}, 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 {
|
|
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)
|
|
}
|
|
if res.Summary != nil {
|
|
if err := p.store.ClearSummarizeRequested(ctx, userID, videoID); err != nil {
|
|
return fmt.Errorf("clear summarize flag %q: %w", videoID, err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|