diff --git a/internal/runner/runner.go b/internal/runner/runner.go new file mode 100644 index 0000000..7df7d31 --- /dev/null +++ b/internal/runner/runner.go @@ -0,0 +1,159 @@ +// Package runner wires the engine to the durable store for the `tapir run` +// command. It owns the cross-restart dedup the engine core deliberately does +// not: the engine's in-memory processed map is process-lifetime only, so this +// loads the store's SeenVideoIDs and skips videos already summarized in a prior +// run. It also assigns each video its durable store id (UpsertVideo) before +// processing, so the summary's video_id equals the dedup key. +// +// It depends on small local interfaces (VideoStore, Processor), not concrete +// types, so the loop is tested with fakes — no live YouTube, gateway, or PG. +package runner + +import ( + "context" + "errors" + "fmt" + "log/slog" + "time" + + "gitea.d-ma.be/mathias/tapir/internal/domain" + "gitea.d-ma.be/mathias/tapir/internal/ports" + "gitea.d-ma.be/mathias/tapir/internal/usecase" +) + +// VideoStore is the durable persistence the run loop needs: assign a stable id + +// metadata, and read the already-summarized set. *store.Store satisfies it. +type VideoStore interface { + UpsertVideo(ctx context.Context, v domain.Video) (string, error) + SeenVideoIDs(ctx context.Context, userID string) (map[string]bool, error) +} + +// Processor runs the core use case for a single video. *usecase.Engine +// satisfies it. +type Processor interface { + ProcessNewVideo(ctx context.Context, v domain.Video) (usecase.ProcessResult, error) +} + +// Runner walks a user's subscriptions, persists each candidate video, skips the +// ones already summarized (durably), and processes the rest through the engine. +type Runner struct { + src ports.VideoSource + store VideoStore + engine Processor + userID string + log *slog.Logger +} + +// New builds a Runner. A nil logger falls back to slog.Default. +func New(src ports.VideoSource, store VideoStore, engine Processor, userID string, log *slog.Logger) *Runner { + if log == nil { + log = slog.Default() + } + return &Runner{src: src, store: store, engine: engine, userID: userID, log: log} +} + +// Stats summarizes one RunOnce pass. +type Stats struct { + Candidates int + Summarized int + SkippedSeen int + SkippedNoText int + Errors int +} + +// RunOnce performs a single pass over the user's subscriptions. Per-item errors +// are logged and collected (one bad video or channel does not abort the pass) +// and returned joined alongside the Stats gathered. +func (r *Runner) RunOnce(ctx context.Context) (Stats, error) { + var ( + stats Stats + errs []error + ) + + seen, err := r.store.SeenVideoIDs(ctx, r.userID) + if err != nil { + return stats, fmt.Errorf("runner: load seen videos: %w", err) + } + + subs, err := r.src.ListSubscriptions(ctx, r.userID) + if err != nil { + return stats, fmt.Errorf("runner: list subscriptions: %w", err) + } + + for _, sub := range subs { + vids, err := r.src.NewVideos(ctx, sub) + if err != nil { + errs = append(errs, fmt.Errorf("new videos for %q: %w", sub.ChannelTitle, err)) + stats.Errors++ + continue + } + for _, v := range vids { + stats.Candidates++ + v.UserID = r.userID // keep the dedup/FK key consistent with config + + id, err := r.store.UpsertVideo(ctx, v) + if err != nil { + errs = append(errs, fmt.Errorf("upsert video %q: %w", v.ProviderVideoID, err)) + stats.Errors++ + continue + } + v.ID = id + + if seen[id] { + stats.SkippedSeen++ + continue + } + seen[id] = true // also guard against the same video within this pass + + res, err := r.engine.ProcessNewVideo(ctx, v) + if err != nil { + errs = append(errs, fmt.Errorf("process %q: %w", v.ProviderVideoID, err)) + stats.Errors++ + continue + } + switch { + case res.Skipped: + stats.SkippedNoText++ + r.log.Info("skipped video (no transcript)", "video", v.ProviderVideoID, "title", v.Title) + case res.Summary != nil: + stats.Summarized++ + r.log.Info("summarized video", "video", v.ProviderVideoID, "title", v.Title, + "provider", res.Summary.AIProvider, "model", res.Summary.AIModel) + } + } + } + + return stats, errors.Join(errs...) +} + +// Loop runs RunOnce immediately, then on every interval tick until ctx is +// cancelled. A zero or negative interval means a single pass (no loop). Per-pass +// errors are logged, not fatal, so a transient failure doesn't kill the watcher. +func (r *Runner) Loop(ctx context.Context, interval time.Duration) error { + runPass := func() { + stats, err := r.RunOnce(ctx) + r.log.Info("run pass complete", + "candidates", stats.Candidates, "summarized", stats.Summarized, + "skipped_seen", stats.SkippedSeen, "skipped_no_text", stats.SkippedNoText, + "errors", stats.Errors) + if err != nil { + r.log.Warn("run pass had errors", "err", err) + } + } + + runPass() + if interval <= 0 { + return nil + } + + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return ctx.Err() + case <-ticker.C: + runPass() + } + } +} diff --git a/internal/runner/runner_test.go b/internal/runner/runner_test.go new file mode 100644 index 0000000..5bacdce --- /dev/null +++ b/internal/runner/runner_test.go @@ -0,0 +1,161 @@ +package runner_test + +import ( + "context" + "io" + "log/slog" + "testing" + + "github.com/stretchr/testify/require" + + "gitea.d-ma.be/mathias/tapir/internal/domain" + "gitea.d-ma.be/mathias/tapir/internal/runner" + "gitea.d-ma.be/mathias/tapir/internal/usecase" +) + +const testUser = "11111111-1111-1111-1111-111111111111" + +// --- fakes ----------------------------------------------------------------- + +type fakeSource struct { + subs []domain.Subscription + videos map[string][]domain.Video // keyed by channel id + transcripts map[string]domain.Transcript +} + +func (f *fakeSource) ListSubscriptions(_ context.Context, _ string) ([]domain.Subscription, error) { + return f.subs, nil +} + +func (f *fakeSource) NewVideos(_ context.Context, sub domain.Subscription) ([]domain.Video, error) { + return f.videos[sub.ChannelID], nil +} + +func (f *fakeSource) FetchTranscript(_ context.Context, v domain.Video) (domain.Transcript, error) { + if t, ok := f.transcripts[v.ProviderVideoID]; ok { + return t, nil + } + return domain.Transcript{VideoID: v.ID, UserID: v.UserID, Source: domain.SourceCaptions, Content: "default transcript text"}, nil +} + +// fakeStore assigns deterministic ids ("id-"+provider video id) so a pre-seeded +// seen set lines up with UpsertVideo output, modelling cross-restart dedup. +type fakeStore struct { + seen map[string]bool + upserted []domain.Video +} + +func (f *fakeStore) UpsertVideo(_ context.Context, v domain.Video) (string, error) { + f.upserted = append(f.upserted, v) + return "id-" + v.ProviderVideoID, nil +} + +func (f *fakeStore) SeenVideoIDs(_ context.Context, _ string) (map[string]bool, error) { + cp := make(map[string]bool, len(f.seen)) + for k, v := range f.seen { + cp[k] = v + } + return cp, nil +} + +type fakeSummarizer struct{} + +func (fakeSummarizer) Summarize(_ context.Context, v domain.Video, _ domain.Transcript) (domain.Summary, error) { + return domain.Summary{UserID: v.UserID, VideoID: v.ID, Summary: "s", AIProvider: "local", AIModel: "koala/phi4-mini"}, nil +} + +type recordingSink struct{ delivered []domain.Summary } + +func (s *recordingSink) Name() string { return "store" } +func (s *recordingSink) Deliver(_ context.Context, sum domain.Summary) error { + s.delivered = append(s.delivered, sum) + return nil +} + +func sub(channelID, title string) domain.Subscription { + return domain.Subscription{UserID: testUser, ChannelID: channelID, ChannelTitle: title, Active: true} +} + +func vid(provID, title string) domain.Video { + return domain.Video{UserID: testUser, Provider: domain.ProviderYouTube, ProviderVideoID: provID, Title: title} +} + +func quietLogger() *slog.Logger { + return slog.New(slog.NewTextHandler(io.Discard, nil)) +} + +// --- tests ----------------------------------------------------------------- + +func TestRunOnce_SummarizesNewVideos(t *testing.T) { + src := &fakeSource{ + subs: []domain.Subscription{sub("chan1", "Channel One")}, + videos: map[string][]domain.Video{"chan1": {vid("v1", "Video 1"), vid("v2", "Video 2")}}, + } + st := &fakeStore{seen: map[string]bool{}} + sink := &recordingSink{} + eng := usecase.NewEngine(src, fakeSummarizer{}, sink) + r := runner.New(src, st, eng, testUser, quietLogger()) + + stats, err := r.RunOnce(context.Background()) + require.NoError(t, err) + require.Equal(t, 2, stats.Candidates) + require.Equal(t, 2, stats.Summarized) + require.Equal(t, 0, stats.SkippedSeen) + require.Len(t, sink.delivered, 2) + + // Each delivered summary must carry the durable store id as its video id. + require.Equal(t, "id-v1", sink.delivered[0].VideoID) + require.Equal(t, "id-v2", sink.delivered[1].VideoID) +} + +func TestRunOnce_SkipsAlreadySummarized(t *testing.T) { + src := &fakeSource{ + subs: []domain.Subscription{sub("chan1", "Channel One")}, + videos: map[string][]domain.Video{"chan1": {vid("v1", "Video 1"), vid("v2", "Video 2")}}, + } + // v1 was summarized in a prior run (durable seen set). + st := &fakeStore{seen: map[string]bool{"id-v1": true}} + sink := &recordingSink{} + eng := usecase.NewEngine(src, fakeSummarizer{}, sink) + r := runner.New(src, st, eng, testUser, quietLogger()) + + stats, err := r.RunOnce(context.Background()) + require.NoError(t, err) + require.Equal(t, 1, stats.SkippedSeen) + require.Equal(t, 1, stats.Summarized) + require.Len(t, sink.delivered, 1) + require.Equal(t, "id-v2", sink.delivered[0].VideoID, "only the unseen video is summarized") +} + +func TestRunOnce_SkipsVideosWithoutTranscript(t *testing.T) { + src := &fakeSource{ + subs: []domain.Subscription{sub("chan1", "Channel One")}, + videos: map[string][]domain.Video{"chan1": {vid("v1", "Video 1")}}, + transcripts: map[string]domain.Transcript{"v1": {Source: domain.SourceNone}}, + } + st := &fakeStore{seen: map[string]bool{}} + sink := &recordingSink{} + eng := usecase.NewEngine(src, fakeSummarizer{}, sink) + r := runner.New(src, st, eng, testUser, quietLogger()) + + stats, err := r.RunOnce(context.Background()) + require.NoError(t, err) + require.Equal(t, 1, stats.SkippedNoText) + require.Equal(t, 0, stats.Summarized) + require.Empty(t, sink.delivered, "no summary delivered when there is no transcript") +} + +func TestRunOnce_UpsertsEveryCandidate(t *testing.T) { + src := &fakeSource{ + subs: []domain.Subscription{sub("chan1", "Channel One")}, + videos: map[string][]domain.Video{"chan1": {vid("v1", "Video 1"), vid("v2", "Video 2")}}, + } + // Even an already-seen video gets upserted so its metadata stays fresh. + st := &fakeStore{seen: map[string]bool{"id-v1": true}} + eng := usecase.NewEngine(src, fakeSummarizer{}, &recordingSink{}) + r := runner.New(src, st, eng, testUser, quietLogger()) + + _, err := r.RunOnce(context.Background()) + require.NoError(t, err) + require.Len(t, st.upserted, 2, "every candidate is upserted, including seen ones") +}