After ProcessNewVideo the runner records transcript_status per outcome: rate_limited (stamps the backoff clock), none, or fetched. Before fetching, a video inside the TAPIR_FETCH_BACKOFF window is skipped (SkippedRateLimited) so a just-429'd caption endpoint is not re-hit; once the window expires it retries. Backoff/clock injected via variadic Options (WithBackoff, WithClock) so existing New call sites and the fake-driven loop tests stay valid. Backoff 0 = always retry. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
299 lines
11 KiB
Go
299 lines
11 KiB
Go
package runner_test
|
|
|
|
import (
|
|
"context"
|
|
"io"
|
|
"log/slog"
|
|
"testing"
|
|
"time"
|
|
|
|
"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.
|
|
// auto controls the summarization mode; requested is the manual-mode queue keyed
|
|
// by store id; cleared records the ids whose queue flag the runner reset.
|
|
type fakeStore struct {
|
|
seen map[string]bool
|
|
upserted []domain.Video
|
|
auto bool
|
|
requested map[string]bool
|
|
cleared []string
|
|
rateLimited map[string]time.Time // id -> when 429'd (seeds the backoff window)
|
|
statuses map[string]string // id -> last SetTranscriptStatus value
|
|
}
|
|
|
|
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
|
|
}
|
|
|
|
func (f *fakeStore) GetAutoSummarize(_ context.Context, _ string) (bool, error) {
|
|
return f.auto, nil
|
|
}
|
|
|
|
func (f *fakeStore) RequestedVideoIDs(_ context.Context, _ string) (map[string]bool, error) {
|
|
cp := make(map[string]bool, len(f.requested))
|
|
for k, v := range f.requested {
|
|
cp[k] = v
|
|
}
|
|
return cp, nil
|
|
}
|
|
|
|
func (f *fakeStore) ClearSummarizeRequested(_ context.Context, _, videoID string) error {
|
|
f.cleared = append(f.cleared, videoID)
|
|
return nil
|
|
}
|
|
|
|
func (f *fakeStore) RateLimitedVideoIDs(_ context.Context, _ string) (map[string]time.Time, error) {
|
|
cp := make(map[string]time.Time, len(f.rateLimited))
|
|
for k, v := range f.rateLimited {
|
|
cp[k] = v
|
|
}
|
|
return cp, nil
|
|
}
|
|
|
|
func (f *fakeStore) SetTranscriptStatus(_ context.Context, _, videoID, status string) error {
|
|
if f.statuses == nil {
|
|
f.statuses = map[string]string{}
|
|
}
|
|
f.statuses[videoID] = status
|
|
return 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{}, auto: 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, 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}, auto: 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{}, auto: 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.SkippedNoText)
|
|
require.Equal(t, 0, stats.Summarized)
|
|
require.Empty(t, sink.delivered, "no summary delivered when there is no transcript")
|
|
}
|
|
|
|
func TestRunOnce_ManualMode_SkipsUnrequested(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")}},
|
|
}
|
|
// Manual mode, nothing queued: discover (upsert) but summarize nothing.
|
|
st := &fakeStore{seen: map[string]bool{}, auto: false, requested: 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.SkippedManual, "manual mode skips unqueued videos")
|
|
require.Equal(t, 0, stats.Summarized)
|
|
require.Empty(t, sink.delivered, "no summary in manual mode without a request")
|
|
require.Len(t, st.upserted, 2, "discovery still persists every candidate")
|
|
}
|
|
|
|
func TestRunOnce_ManualMode_ProcessesRequested(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")}},
|
|
}
|
|
// Manual mode, v1 queued (by store id). Only v1 is summarized; its flag clears.
|
|
st := &fakeStore{seen: map[string]bool{}, auto: false, requested: 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.Summarized, "only the queued video is summarized")
|
|
require.Equal(t, 1, stats.SkippedManual, "the unqueued video is skipped")
|
|
require.Len(t, sink.delivered, 1)
|
|
require.Equal(t, "id-v1", sink.delivered[0].VideoID)
|
|
require.Equal(t, []string{"id-v1"}, st.cleared, "the queue flag is cleared after summarizing")
|
|
}
|
|
|
|
// noFetchSource fails the test if a transcript fetch happens — used to prove the
|
|
// runner skips a rate-limited video before touching the caption endpoint.
|
|
type noFetchSource struct{ *fakeSource }
|
|
|
|
func (noFetchSource) FetchTranscript(context.Context, domain.Video) (domain.Transcript, error) {
|
|
panic("FetchTranscript must not be called for a rate-limited video within the backoff window")
|
|
}
|
|
|
|
func TestRunOnce_SkipsRateLimitedWithinBackoff(t *testing.T) {
|
|
base := time.Date(2026, 6, 3, 12, 0, 0, 0, time.UTC)
|
|
src := &fakeSource{
|
|
subs: []domain.Subscription{sub("chan1", "Channel One")},
|
|
videos: map[string][]domain.Video{"chan1": {vid("v1", "Video 1")}},
|
|
}
|
|
// v1 was rate-limited 5m ago; backoff is 1h, so it is still inside the window.
|
|
st := &fakeStore{
|
|
seen: map[string]bool{},
|
|
auto: true,
|
|
rateLimited: map[string]time.Time{"id-v1": base.Add(-5 * time.Minute)},
|
|
}
|
|
eng := usecase.NewEngine(noFetchSource{src}, fakeSummarizer{}, &recordingSink{})
|
|
r := runner.New(noFetchSource{src}, st, eng, testUser, quietLogger(),
|
|
runner.WithBackoff(time.Hour), runner.WithClock(func() time.Time { return base }))
|
|
|
|
stats, err := r.RunOnce(context.Background())
|
|
require.NoError(t, err)
|
|
require.Equal(t, 1, stats.SkippedRateLimited, "still throttled -> skipped")
|
|
require.Equal(t, 0, stats.Summarized)
|
|
require.Empty(t, st.statuses, "no status write: the engine was never invoked")
|
|
}
|
|
|
|
func TestRunOnce_RetriesRateLimitedAfterBackoff(t *testing.T) {
|
|
base := time.Date(2026, 6, 3, 12, 0, 0, 0, time.UTC)
|
|
src := &fakeSource{
|
|
subs: []domain.Subscription{sub("chan1", "Channel One")},
|
|
videos: map[string][]domain.Video{"chan1": {vid("v1", "Video 1")}},
|
|
}
|
|
// v1 was rate-limited 2h ago; backoff is 1h, so the window has expired.
|
|
st := &fakeStore{
|
|
seen: map[string]bool{},
|
|
auto: true,
|
|
rateLimited: map[string]time.Time{"id-v1": base.Add(-2 * time.Hour)},
|
|
}
|
|
sink := &recordingSink{}
|
|
eng := usecase.NewEngine(src, fakeSummarizer{}, sink)
|
|
r := runner.New(src, st, eng, testUser, quietLogger(),
|
|
runner.WithBackoff(time.Hour), runner.WithClock(func() time.Time { return base }))
|
|
|
|
stats, err := r.RunOnce(context.Background())
|
|
require.NoError(t, err)
|
|
require.Equal(t, 0, stats.SkippedRateLimited, "window expired -> not skipped")
|
|
require.Equal(t, 1, stats.Summarized, "the video is retried and summarized")
|
|
require.Len(t, sink.delivered, 1)
|
|
require.Equal(t, "fetched", st.statuses["id-v1"], "status advances to fetched on success")
|
|
}
|
|
|
|
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}, auto: 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")
|
|
}
|