diff --git a/internal/scheduler/archive.go b/internal/scheduler/archive.go new file mode 100644 index 0000000..586cd0f --- /dev/null +++ b/internal/scheduler/archive.go @@ -0,0 +1,105 @@ +package scheduler + +import ( + "context" + "log/slog" + "time" + + "tumble/internal/archive" + "tumble/internal/data" +) + +const ( + archiveRecheckNotFound = 30 * 24 * time.Hour + archiveRecheckError = 24 * time.Hour +) + +func runArchiveBatch(ctx context.Context, store data.Store, client *archive.Client) { + // 1. Get unchecked dead link URLs + unchecked, err := store.GetUncheckedDeadLinkURLs(ctx) + if err != nil { + slog.Error("Archive batch: failed to get unchecked URLs", "error", err) + return + } + + // 2. Get stale not_found/error URLs needing recheck + staleNotFound, err := store.GetStaleArchiveLookups(ctx, archiveRecheckNotFound) + if err != nil { + slog.Error("Archive batch: failed to get stale not_found URLs", "error", err) + } + staleError, err := store.GetStaleArchiveLookups(ctx, archiveRecheckError) + if err != nil { + slog.Error("Archive batch: failed to get stale error URLs", "error", err) + } + + // Combine and deduplicate + seen := make(map[string]bool) + var urls []string + for _, u := range unchecked { + if !seen[u] { + seen[u] = true + urls = append(urls, u) + } + } + for _, u := range staleNotFound { + if !seen[u] { + seen[u] = true + urls = append(urls, u) + } + } + for _, u := range staleError { + if !seen[u] { + seen[u] = true + urls = append(urls, u) + } + } + + if len(urls) == 0 { + slog.Info("Archive batch: no URLs to check") + return + } + + slog.Info("Archive batch: starting", "total", len(urls)) + + for i, u := range urls { + if ctx.Err() != nil { + slog.Info("Archive batch: context cancelled, stopping", "checked", i) + return + } + + result, err := client.Check(ctx, u) + if err != nil { + slog.Warn("Archive batch: API error", "url", u, "error", err) + store.UpsertArchiveLookup(ctx, &data.ArchiveLookup{ + URL: u, + Status: "error", + CheckedAt: time.Now(), + }) + continue + } + + lookup := &data.ArchiveLookup{ + URL: u, + CheckedAt: time.Now(), + } + if result.Found { + lookup.Status = "found" + lookup.ArchiveURL = &result.ArchiveURL + if !result.SnapshotAt.IsZero() { + lookup.SnapshotAt = &result.SnapshotAt + } + } else { + lookup.Status = "not_found" + } + + if err := store.UpsertArchiveLookup(ctx, lookup); err != nil { + slog.Warn("Archive batch: failed to store result", "url", u, "error", err) + } + + if (i+1)%50 == 0 { + slog.Info("Archive batch: progress", "checked", i+1, "total", len(urls)) + } + } + + slog.Info("Archive batch: complete", "total", len(urls)) +} diff --git a/internal/scheduler/scheduler.go b/internal/scheduler/scheduler.go index 58c7883..cb0ce35 100644 --- a/internal/scheduler/scheduler.go +++ b/internal/scheduler/scheduler.go @@ -7,6 +7,7 @@ import ( "time" "github.com/robfig/cron/v3" + "tumble/internal/archive" "tumble/internal/data" ) @@ -17,11 +18,12 @@ const ( // Scheduler manages scheduled tasks for the application. type Scheduler struct { - cron *cron.Cron - store data.Store - retryCount int - retryMu sync.Mutex - stopRetry chan struct{} + cron *cron.Cron + store data.Store + archiveClient *archive.Client + retryCount int + retryMu sync.Mutex + stopRetry chan struct{} } // New creates a new Scheduler with the given store. @@ -34,9 +36,10 @@ func New(store data.Store) *Scheduler { } return &Scheduler{ - cron: cron.New(cron.WithLocation(loc)), - store: store, - stopRetry: make(chan struct{}), + cron: cron.New(cron.WithLocation(loc)), + store: store, + archiveClient: archive.NewClient(5), + stopRetry: make(chan struct{}), } } @@ -50,12 +53,23 @@ func (s *Scheduler) Start(ctx context.Context) error { return err } + // Schedule archive batch at 3 AM Central + _, err = s.cron.AddFunc("0 3 * * *", func() { + runArchiveBatch(ctx, s.store, s.archiveClient) + }) + if err != nil { + return err + } + s.cron.Start() slog.Info("Scheduler started", "nextRun", s.cron.Entries()[0].Next) // Check if we need to fetch today's kitten on startup go s.checkStartupKitten(ctx) + // Run archive batch on startup (in background) + go runArchiveBatch(ctx, s.store, s.archiveClient) + return nil }