From 937674f808ab13d48252ed5b187ca16141e969fc Mon Sep 17 00:00:00 2001 From: AmanTahiliani Date: Mon, 25 May 2026 13:52:43 -0400 Subject: [PATCH] feat(backend): implement session coverage, resumable ingestion, deep year backfill, backoff rate-limit, and dynamic cache TTLs --- cmd/main.go | 138 ++++++++++- internal/api/cache.go | 22 +- internal/api/openf1.go | 19 ++ internal/ingest/ingest.go | 269 +++++++++++++++++---- internal/ingest/ingest_test.go | 62 +++++ internal/store/coverage.go | 117 +++++++++ internal/store/coverage_test.go | 134 ++++++++++ internal/store/migrations/004_coverage.sql | 9 + internal/store/store_test.go | 11 +- 9 files changed, 731 insertions(+), 50 deletions(-) create mode 100644 internal/store/coverage.go create mode 100644 internal/store/coverage_test.go create mode 100644 internal/store/migrations/004_coverage.sql diff --git a/cmd/main.go b/cmd/main.go index 5c07232..90dd733 100644 --- a/cmd/main.go +++ b/cmd/main.go @@ -7,6 +7,7 @@ import ( "log" "net/http" "os" + "strings" "time" "github.com/AmanTahiliani/box-box/internal/api" @@ -22,13 +23,24 @@ func main() { webMode := flag.Bool("web", false, "Start web companion server instead of TUI") port := flag.Int("port", 8080, "Port for web server (used with --web)") ingestYear := flag.Int("ingest-year", 0, "Ingest OpenF1 meetings for a season year") + backfillSeason := flag.Int("backfill-season", 0, "Trigger full-season backfill/deep-ingestion for the given year") ingestMeeting := flag.Int("ingest-meeting", 0, "Ingest meeting metadata and Race Hub datasets for all sessions") ingestSession := flag.Int("ingest-session", 0, "Ingest Race Hub datasets for a session key") ingestNews := flag.Bool("ingest-news", false, "Refresh RSS/Atom paddock briefing feeds") dryRun := flag.Bool("dry-run", false, "Preview ingestion without writing domain rows") + force := flag.Bool("force", false, "Re-ingest datasets even if already tracked in the session_coverage table as completed") + coverageYear := flag.Int("coverage", 0, "Show season coverage report for the given year") dbPath := flag.String("db", "", "Domain database path (default: ~/.local/share/box-box/boxbox.db)") flag.Parse() + if *coverageYear != 0 { + if err := runCoverageReport(*coverageYear, *dbPath); err != nil { + fmt.Fprintf(os.Stderr, "coverage report error: %v\n", err) + os.Exit(1) + } + return + } + var client *api.OpenF1Client if apiKey := os.Getenv("OPENF1_API_KEY"); apiKey != "" { client = api.NewOpenF1ClientWithKey("https://api.openf1.org", 15*time.Second, apiKey) @@ -44,6 +56,9 @@ func main() { if *ingestYear != 0 { ingestFlags++ } + if *backfillSeason != 0 { + ingestFlags++ + } if *ingestMeeting != 0 { ingestFlags++ } @@ -55,7 +70,7 @@ func main() { } if ingestFlags > 0 { if ingestFlags > 1 { - fmt.Fprintln(os.Stderr, "box-box: only one of --ingest-year, --ingest-meeting, --ingest-session, or --ingest-news may be set") + fmt.Fprintln(os.Stderr, "box-box: only one of --ingest-year, --backfill-season, --ingest-meeting, --ingest-session, or --ingest-news may be set") os.Exit(1) } if *ingestNews { @@ -65,7 +80,13 @@ func main() { } return } - if err := runIngestion(client, *ingestYear, *ingestMeeting, *ingestSession, *dryRun, *dbPath); err != nil { + + yearVal := *ingestYear + if *backfillSeason != 0 { + yearVal = *backfillSeason + } + + if err := runIngestion(client, yearVal, *ingestMeeting, *ingestSession, *force, *dryRun, *dbPath); err != nil { fmt.Fprintf(os.Stderr, "box-box ingest error: %v\n", err) os.Exit(1) } @@ -113,7 +134,7 @@ func main() { } } -func runIngestion(client *api.OpenF1Client, year, meetingKey, sessionKey int, dryRun bool, dbPath string) error { +func runIngestion(client *api.OpenF1Client, year, meetingKey, sessionKey int, force, dryRun bool, dbPath string) error { log.SetOutput(os.Stderr) path := dbPath @@ -129,6 +150,7 @@ func runIngestion(client *api.OpenF1Client, year, meetingKey, sessionKey int, dr opts := ingest.DefaultOptions() opts.DryRun = dryRun + opts.Force = force opts.Progress = ingest.NewProgress(os.Stderr) svc := ingest.NewService(st, ingest.NewOpenF1Source(client), opts) @@ -185,3 +207,113 @@ func runNewsIngestion(dryRun bool, dbPath string) error { ) return err } + +func runCoverageReport(year int, dbPath string) error { + path := dbPath + if path == "" { + path = store.DefaultDBPath() + } + + st, err := store.Open(path) + if err != nil { + return fmt.Errorf("open domain database: %w", err) + } + defer st.Close() + + rows, err := st.GetSeasonCoverage(year) + if err != nil { + return fmt.Errorf("get season coverage: %w", err) + } + + if len(rows) == 0 { + fmt.Printf("No session coverage data found for year %d.\n", year) + return nil + } + + type datasetStatus struct { + Status string + Count int + } + + type sessionInfo struct { + MeetingName string + SessionName string + SessionKey int + Datasets map[string]datasetStatus + } + + var sessions []sessionInfo + sessionMap := make(map[int]int) + + for _, row := range rows { + idx, exists := sessionMap[row.SessionKey] + if !exists { + idx = len(sessions) + sessions = append(sessions, sessionInfo{ + MeetingName: row.MeetingName, + SessionName: row.SessionName, + SessionKey: row.SessionKey, + Datasets: make(map[string]datasetStatus), + }) + sessionMap[row.SessionKey] = idx + } + if row.Dataset != "" { + sessions[idx].Datasets[row.Dataset] = datasetStatus{ + Status: row.Status, + Count: row.RowCount, + } + } + } + + fmt.Printf("\n--- Season %d Coverage Report ---\n\n", year) + fmt.Printf("%-35s | %-5s | %-2s | %-2s | %-2s | %-2s | %-2s | %-2s | %-2s | %-2s | %-2s\n", + "Meeting / Session (Key)", "ID", "DR", "SR", "SG", "ST", "PS", "PO", "RC", "WE", "LA") + fmt.Println(strings.Repeat("-", 82)) + + for _, sess := range sessions { + statusChar := func(ds string) string { + dsStatus, ok := sess.Datasets[ds] + if !ok { + return "." + } + switch dsStatus.Status { + case "complete": + return "✓" + case "failed": + return "✗" + default: + return "." + } + } + + nameCol := fmt.Sprintf("%s - %s (%d)", sess.MeetingName, sess.SessionName, sess.SessionKey) + if len(nameCol) > 35 { + nameCol = nameCol[:32] + "..." + } + + fmt.Printf("%-35s | %-5d | %s | %s | %s | %s | %s | %s | %s | %s | %s\n", + nameCol, + sess.SessionKey, + statusChar("drivers"), + statusChar("session_result"), + statusChar("starting_grid"), + statusChar("stints"), + statusChar("pit_stops"), + statusChar("positions"), + statusChar("race_control"), + statusChar("weather"), + statusChar("laps"), + ) + } + + fmt.Println(strings.Repeat("-", 82)) + fmt.Println("\nLegend:") + fmt.Println(" [✓] Complete [✗] Failed [.] Pending/Unattempted") + fmt.Println("Datasets:") + fmt.Println(" DR: drivers SR: session_result SG: starting_grid") + fmt.Println(" ST: stints PS: pit_stops PO: positions") + fmt.Println(" RC: race_control WE: weather LA: laps") + fmt.Println() + + return nil +} diff --git a/internal/api/cache.go b/internal/api/cache.go index 3b2826a..81f63f6 100644 --- a/internal/api/cache.go +++ b/internal/api/cache.go @@ -5,6 +5,7 @@ import ( "encoding/json" "os" "path/filepath" + "strconv" "strings" "sync/atomic" "time" @@ -96,11 +97,30 @@ func cacheDBPath() string { // ttlForURL determines the appropriate TTL based on the URL pattern. // Returns 0 (CacheTTLForever) for historical data that will never change. func ttlForURL(url string) time.Duration { + var year int + if idx := strings.Index(url, "year="); idx != -1 && len(url) >= idx+9 { + yearStr := url[idx+5 : idx+9] + if y, err := strconv.Atoi(yearStr); err == nil { + year = y + } + } + + currentYear := time.Now().Year() + // Historical data — completed past seasons never change. - if strings.Contains(url, "year=2023") || strings.Contains(url, "year=2024") { + if year > 0 && year < currentYear { return CacheTTLForever } + // For current year or unspecified year (e.g. meeting/session list endpoint that includes a session key query): + if year == currentYear || year == 0 { + // Cache current year meetings and sessions metadata for 24h + if (strings.Contains(url, "/meetings") || strings.Contains(url, "/sessions")) && + !strings.Contains(url, "/session_result") { + return CacheTTLLong + } + } + // Live telemetry endpoints — change every few seconds during a session. if strings.Contains(url, "/position") || strings.Contains(url, "/intervals") || diff --git a/internal/api/openf1.go b/internal/api/openf1.go index 3c4d700..6bd6183 100644 --- a/internal/api/openf1.go +++ b/internal/api/openf1.go @@ -20,6 +20,15 @@ import ( // from ~30 min before a session starts until ~30 min after it ends. var ErrLiveSessionLocked = errors.New("live F1 session in progress — API access is restricted to authenticated users until the session ends") +// RateLimitError is returned when the OpenF1 API rate limit is reached (HTTP 429). +type RateLimitError struct { + RetryAfter time.Duration +} + +func (e *RateLimitError) Error() string { + return fmt.Sprintf("openf1 API rate limit hit: retry after %v", e.RetryAfter) +} + // IsLiveSessionError reports whether err (or any error in its chain) is the // live-session lockout error from the OpenF1 API. func IsLiveSessionError(err error) bool { @@ -111,6 +120,16 @@ func (c *OpenF1Client) FetchStrict(url string) ([]byte, error) { } defer resp.Body.Close() + if resp.StatusCode == http.StatusTooManyRequests { + retryAfterDur := 500 * time.Millisecond + if retryAfterHeader := resp.Header.Get("Retry-After"); retryAfterHeader != "" { + if seconds, err := strconv.Atoi(retryAfterHeader); err == nil { + retryAfterDur = time.Duration(seconds) * time.Second + } + } + return nil, &RateLimitError{RetryAfter: retryAfterDur} + } + data, err := io.ReadAll(resp.Body) if err != nil { return nil, err diff --git a/internal/ingest/ingest.go b/internal/ingest/ingest.go index d2e7e1c..c93f39f 100644 --- a/internal/ingest/ingest.go +++ b/internal/ingest/ingest.go @@ -4,6 +4,7 @@ import ( "encoding/json" "errors" "fmt" + "math/rand" "strings" "time" @@ -15,6 +16,7 @@ import ( // Options configures ingestion behavior. type Options struct { DryRun bool + Force bool // Re-fetch even if session_coverage says 'complete' RequestDelay time.Duration MaxRetries int RetryBackoff time.Duration @@ -25,7 +27,7 @@ type Options struct { func DefaultOptions() Options { return Options{ RequestDelay: 300 * time.Millisecond, - MaxRetries: 3, + MaxRetries: 5, RetryBackoff: 500 * time.Millisecond, Progress: NewProgress(nil), } @@ -71,7 +73,7 @@ type Service struct { // NewService creates an ingestion service. func NewService(st *store.Store, source Source, opts Options) *Service { if opts.MaxRetries <= 0 { - opts.MaxRetries = 3 + opts.MaxRetries = 5 } if opts.RetryBackoff <= 0 { opts.RetryBackoff = 500 * time.Millisecond @@ -132,7 +134,75 @@ func (s *Service) IngestYear(year int) (Summary, error) { summary.Meetings++ } - summary.Status = statusForDryRun(s.opts.DryRun) + sessionFailures := 0 + partialSessions := 0 + var totalSessionsCount int + + for _, m := range meetings { + s.opts.Progress.Step("fetching sessions for meeting %d (%s)", m.MeetingKey, m.MeetingName) + sessionFetch, sessions, err := fetchWithRetry(s, func() (FetchResult, []models.Session, error) { + return s.source.FetchSessionsForMeeting(int(m.MeetingKey)) + }) + if err != nil { + summary.Errors = append(summary.Errors, fmt.Sprintf("meeting %d (%s) sessions: %v", m.MeetingKey, m.MeetingName, err)) + continue + } + summary.RawPayloads++ + if !s.opts.DryRun { + mkVal := int(m.MeetingKey) + inserted, err := s.storeRaw(sessionFetch, &mkVal, nil) + if err != nil { + return s.finishFailed(runID, summary, err) + } + if inserted { + summary.RawInserted++ + } + } + s.delay() + + for _, sess := range sessions { + if s.opts.DryRun { + summary.Sessions++ + totalSessionsCount++ + continue + } + if err := s.store.UpsertSession(sessionToStore(sess)); err != nil { + return s.finishFailed(runID, summary, err) + } + summary.Sessions++ + totalSessionsCount++ + } + + for _, sess := range sessions { + s.opts.Progress.Step("ingesting Race Hub datasets for session %d (%s)", sess.SessionKey, sess.SessionName) + sessSummary, err := s.ingestSessionDatasets(sess) + ss := SessionSummary{ + SessionKey: sess.SessionKey, + SessionName: sess.SessionName, + Summary: sessSummary, + } + if err != nil { + sessionFailures++ + ss.Summary.Status = "failed" + ss.Summary.Errors = append(ss.Summary.Errors, err.Error()) + summary.Errors = append(summary.Errors, fmt.Sprintf( + "session %d (%s): %v", sess.SessionKey, sess.SessionName, err, + )) + } + if err == nil && sessSummary.Status == "partial" { + partialSessions++ + for _, partialErr := range sessSummary.Errors { + summary.Errors = append(summary.Errors, fmt.Sprintf( + "session %d (%s): %s", sess.SessionKey, sess.SessionName, partialErr, + )) + } + } + summary.SessionSummaries = append(summary.SessionSummaries, ss) + summary.mergeCounts(sessSummary) + } + } + + summary.Status = meetingStatus(sessionFailures, partialSessions, totalSessionsCount, s.opts.DryRun) s.finishRun(runID, summary) s.opts.Progress.Summary(summary) return summary, nil @@ -355,63 +425,94 @@ func (s *Service) ingestSessionDatasets(sess models.Session) (Summary, error) { } sk := sessionKey - s.opts.Progress.Step("fetching drivers for session %d", sessionKey) - driverFetch, drivers, err := fetchWithRetry(s, func() (FetchResult, []models.Driver, error) { - return s.source.FetchDriversForSession(sessionKey) - }) + coverage, err := s.store.GetSessionCoverage(sessionKey) if err != nil { - return summary, err + coverage = make(map[string]store.CoverageEntry) } - summary.RawPayloads++ - if !s.opts.DryRun { - inserted, err := s.storeRaw(driverFetch, &meetingKey, &sk) + + // 1. Ingest drivers + if cov, ok := coverage["drivers"]; ok && cov.Status == "complete" && !s.opts.Force { + s.opts.Progress.Step("drivers already complete for session %d, skipping", sessionKey) + summary.Drivers = cov.RowCount + } else { + s.opts.Progress.Step("fetching drivers for session %d", sessionKey) + driverFetch, drivers, err := fetchWithRetry(s, func() (FetchResult, []models.Driver, error) { + return s.source.FetchDriversForSession(sessionKey) + }) if err != nil { + if !s.opts.DryRun { + _ = s.store.UpsertCoverage(sessionKey, "drivers", "failed", 0, err.Error()) + } return summary, err } - if inserted { - summary.RawInserted++ - } - for _, d := range drivers { - if err := s.store.UpsertDriver(driverToStore(d)); err != nil { + summary.RawPayloads++ + if !s.opts.DryRun { + inserted, err := s.storeRaw(driverFetch, &meetingKey, &sk) + if err != nil { + _ = s.store.UpsertCoverage(sessionKey, "drivers", "failed", 0, err.Error()) return summary, err } - if err := s.store.UpsertSessionDriver(sessionDriverToStore(d)); err != nil { - return summary, err + if inserted { + summary.RawInserted++ } - summary.Drivers++ + for _, d := range drivers { + if err := s.store.UpsertDriver(driverToStore(d)); err != nil { + _ = s.store.UpsertCoverage(sessionKey, "drivers", "failed", 0, err.Error()) + return summary, err + } + if err := s.store.UpsertSessionDriver(sessionDriverToStore(d)); err != nil { + _ = s.store.UpsertCoverage(sessionKey, "drivers", "failed", 0, err.Error()) + return summary, err + } + summary.Drivers++ + } + _ = s.store.UpsertCoverage(sessionKey, "drivers", "complete", len(drivers), "") + } else { + summary.Drivers = len(drivers) } - } else { - summary.Drivers = len(drivers) + s.delay() } - s.delay() - s.opts.Progress.Step("fetching session results for session %d", sessionKey) - resultFetch, results, err := fetchWithRetry(s, func() (FetchResult, []models.SessionResult, error) { - return s.source.FetchSessionResult(sessionKey) - }) - if err != nil { - return summary, err - } - summary.RawPayloads++ - if !s.opts.DryRun { - inserted, err := s.storeRaw(resultFetch, &meetingKey, &sk) + // 2. Ingest session_result + if cov, ok := coverage["session_result"]; ok && cov.Status == "complete" && !s.opts.Force { + s.opts.Progress.Step("session_result already complete for session %d, skipping", sessionKey) + summary.SessionResults = cov.RowCount + } else { + s.opts.Progress.Step("fetching session results for session %d", sessionKey) + resultFetch, results, err := fetchWithRetry(s, func() (FetchResult, []models.SessionResult, error) { + return s.source.FetchSessionResult(sessionKey) + }) if err != nil { + if !s.opts.DryRun { + _ = s.store.UpsertCoverage(sessionKey, "session_result", "failed", 0, err.Error()) + } return summary, err } - if inserted { - summary.RawInserted++ - } - for _, r := range results { - if err := s.store.UpsertSessionResult(sessionResultToStore(r)); err != nil { + summary.RawPayloads++ + if !s.opts.DryRun { + inserted, err := s.storeRaw(resultFetch, &meetingKey, &sk) + if err != nil { + _ = s.store.UpsertCoverage(sessionKey, "session_result", "failed", 0, err.Error()) return summary, err } - summary.SessionResults++ + if inserted { + summary.RawInserted++ + } + for _, r := range results { + if err := s.store.UpsertSessionResult(sessionResultToStore(r)); err != nil { + _ = s.store.UpsertCoverage(sessionKey, "session_result", "failed", 0, err.Error()) + return summary, err + } + summary.SessionResults++ + } + _ = s.store.UpsertCoverage(sessionKey, "session_result", "complete", len(results), "") + } else { + summary.SessionResults = len(results) } - } else { - summary.SessionResults = len(results) + s.delay() } - s.delay() + // 3. Optional datasets optionalIngests := []struct { name string run func(*Summary, int, int) error @@ -426,9 +527,74 @@ func (s *Service) ingestSessionDatasets(sess models.Session) (Summary, error) { {name: "weather", run: s.ingestWeather}, {name: "laps", run: s.ingestLaps}, } + for _, optional := range optionalIngests { - if err := optional.run(&summary, meetingKey, sk); err != nil { + if cov, ok := coverage[optional.name]; ok && cov.Status == "complete" && !s.opts.Force { + s.opts.Progress.Step("%s already complete for session %d, skipping", optional.name, sessionKey) + switch optional.name { + case "starting_grid": + summary.StartingGrid = cov.RowCount + case "stints": + summary.Stints = cov.RowCount + case "pit_stops": + summary.PitStops = cov.RowCount + case "positions": + summary.Positions = cov.RowCount + case "race_control": + summary.RaceControl = cov.RowCount + case "weather": + summary.Weather = cov.RowCount + case "laps": + summary.Laps = cov.RowCount + } + continue + } + + var prevCount int + switch optional.name { + case "starting_grid": + prevCount = summary.StartingGrid + case "stints": + prevCount = summary.Stints + case "pit_stops": + prevCount = summary.PitStops + case "positions": + prevCount = summary.Positions + case "race_control": + prevCount = summary.RaceControl + case "weather": + prevCount = summary.Weather + case "laps": + prevCount = summary.Laps + } + + err := optional.run(&summary, meetingKey, sk) + if err != nil { summary.Errors = append(summary.Errors, fmt.Sprintf("%s: %v", optional.name, err)) + if !s.opts.DryRun { + _ = s.store.UpsertCoverage(sessionKey, optional.name, "failed", 0, err.Error()) + } + } else { + if !s.opts.DryRun { + var newCount int + switch optional.name { + case "starting_grid": + newCount = summary.StartingGrid - prevCount + case "stints": + newCount = summary.Stints - prevCount + case "pit_stops": + newCount = summary.PitStops - prevCount + case "positions": + newCount = summary.Positions - prevCount + case "race_control": + newCount = summary.RaceControl - prevCount + case "weather": + newCount = summary.Weather - prevCount + case "laps": + newCount = summary.Laps - prevCount + } + _ = s.store.UpsertCoverage(sessionKey, optional.name, "complete", newCount, "") + } } } @@ -761,7 +927,17 @@ func fetchWithRetry[T any](s *Service, fn fetchFunc[T]) (FetchResult, T, error) for attempt := 0; attempt < s.opts.MaxRetries; attempt++ { if attempt > 0 { - time.Sleep(s.opts.RetryBackoff * time.Duration(attempt)) + backoff := s.opts.RetryBackoff * time.Duration(1< wait { + wait = rle.RetryAfter + } + } + time.Sleep(wait) } fetch, data, err := fn() @@ -785,12 +961,17 @@ func isRetryable(err error) bool { if err == nil { return false } + var rle *api.RateLimitError + if errors.As(err, &rle) { + return true + } msg := strings.ToLower(err.Error()) if strings.Contains(msg, "status 429") || strings.Contains(msg, "status 5") || strings.Contains(msg, "timeout") || strings.Contains(msg, "connection reset") || - strings.Contains(msg, "temporary") { + strings.Contains(msg, "temporary") || + strings.Contains(msg, "rate limit") { return true } var netErr interface{ Timeout() bool } diff --git a/internal/ingest/ingest_test.go b/internal/ingest/ingest_test.go index 490c8b2..f7b6549 100644 --- a/internal/ingest/ingest_test.go +++ b/internal/ingest/ingest_test.go @@ -760,3 +760,65 @@ func TestIngestMeetingPartialFailurePreservesSuccessfulSessions(t *testing.T) { t.Fatalf("meetings after partial failure = %+v, err = %v, want 1 meeting preserved", meetings, err) } } + +func TestResumableIngestion(t *testing.T) { + _, sessionKey, src := testSessionFixtures() + st := openTestStore(t) + + opts := DefaultOptions() + opts.RequestDelay = 0 + svc := NewService(st, src, opts) + + // First run: ingest session. This should populate session_coverage. + summary, err := svc.IngestSession(sessionKey) + if err != nil { + t.Fatalf("first IngestSession() error = %v", err) + } + if summary.Status != "completed" { + t.Fatalf("first run summary.Status = %q, want completed", summary.Status) + } + + // Verify coverage was populated in the database. + cov, err := st.GetSessionCoverage(sessionKey) + if err != nil { + t.Fatalf("GetSessionCoverage() error = %v", err) + } + if len(cov) != 9 { + t.Fatalf("cov length = %d, want 9 datasets", len(cov)) + } + for ds, entry := range cov { + if entry.Status != "complete" { + t.Errorf("dataset %s status = %q, want complete", ds, entry.Status) + } + } + + // Modify the fake source so it returns empty datasets or fails. + // If resumability is working, it should skip all fetches because they are already complete! + src.failOn = "drivers" // If it fetches, it will fail! + + // Second run: without Force, it should skip and succeed. + summary2, err := svc.IngestSession(sessionKey) + if err != nil { + t.Fatalf("second IngestSession() error = %v", err) + } + if summary2.Status != "completed" { + t.Fatalf("second run summary.Status = %q, want completed", summary2.Status) + } + + // Third run: with Force, it should try to fetch and fail as expected! + opts.Force = true + svc2 := NewService(st, src, opts) + _, err = svc2.IngestSession(sessionKey) + if err == nil { + t.Fatal("third IngestSession() expected error due to failOn, got nil") + } + + // Check if drivers dataset is now marked as failed in DB + cov2, err := st.GetSessionCoverage(sessionKey) + if err != nil { + t.Fatalf("GetSessionCoverage() error = %v", err) + } + if cov2["drivers"].Status != "failed" { + t.Fatalf("drivers status = %q, want failed", cov2["drivers"].Status) + } +} diff --git a/internal/store/coverage.go b/internal/store/coverage.go new file mode 100644 index 0000000..50647a8 --- /dev/null +++ b/internal/store/coverage.go @@ -0,0 +1,117 @@ +package store + +import ( + "database/sql" + "fmt" +) + +// CoverageEntry represents the coverage status for a single dataset of a session. +type CoverageEntry struct { + Status string + ErrorMsg string + RowCount int + UpdatedAt string +} + +// SessionCoverageRow represents a joined session coverage record for reporting. +type SessionCoverageRow struct { + MeetingKey int + MeetingName string + SessionKey int + SessionName string + Dataset string + Status string + ErrorMsg string + RowCount int + UpdatedAt string +} + +// UpsertCoverage inserts or updates a session coverage record. +func (s *Store) UpsertCoverage(sessionKey int, dataset string, status string, rowCount int, errMsg string) error { + _, err := s.db.Exec(` + INSERT INTO session_coverage (session_key, dataset, status, row_count, error_msg, updated_at) + VALUES (?, ?, ?, ?, ?, datetime('now')) + ON CONFLICT(session_key, dataset) DO UPDATE SET + status = excluded.status, + row_count = excluded.row_count, + error_msg = excluded.error_msg, + updated_at = excluded.updated_at + `, sessionKey, dataset, status, rowCount, nullString(errMsg)) + if err != nil { + return fmt.Errorf("upsert coverage: %w", err) + } + return nil +} + +// GetSessionCoverage fetches the coverage statuses for all datasets of a given session. +func (s *Store) GetSessionCoverage(sessionKey int) (map[string]CoverageEntry, error) { + rows, err := s.db.Query(` + SELECT dataset, status, row_count, error_msg, updated_at + FROM session_coverage + WHERE session_key = ? + `, sessionKey) + if err != nil { + return nil, fmt.Errorf("get session coverage: %w", err) + } + defer rows.Close() + + coverage := make(map[string]CoverageEntry) + for rows.Next() { + var dataset string + var entry CoverageEntry + var errMsg sql.NullString + if err := rows.Scan(&dataset, &entry.Status, &entry.RowCount, &errMsg, &entry.UpdatedAt); err != nil { + return nil, fmt.Errorf("scan session coverage: %w", err) + } + entry.ErrorMsg = errMsg.String + coverage[dataset] = entry + } + return coverage, rows.Err() +} + +// GetSeasonCoverage returns coverage records for all sessions of a given year. +func (s *Store) GetSeasonCoverage(year int) ([]SessionCoverageRow, error) { + rows, err := s.db.Query(` + SELECT + m.meeting_key, + m.meeting_name, + s.session_key, + s.session_name, + COALESCE(c.dataset, '') as dataset, + COALESCE(c.status, 'pending') as status, + COALESCE(c.error_msg, '') as error_msg, + COALESCE(c.row_count, 0) as row_count, + COALESCE(c.updated_at, '') as updated_at + FROM sessions s + JOIN meetings m ON s.meeting_key = m.meeting_key + LEFT JOIN session_coverage c ON s.session_key = c.session_key + WHERE m.year = ? + ORDER BY m.date_start ASC, m.meeting_key ASC, s.date_start ASC, s.session_key ASC + `, year) + if err != nil { + return nil, fmt.Errorf("get season coverage query: %w", err) + } + defer rows.Close() + + var coverageRows []SessionCoverageRow + for rows.Next() { + var r SessionCoverageRow + var errMsg sql.NullString + if err := rows.Scan( + &r.MeetingKey, + &r.MeetingName, + &r.SessionKey, + &r.SessionName, + &r.Dataset, + &r.Status, + &r.ErrorMsg, + &r.RowCount, + &r.UpdatedAt, + ); err != nil { + return nil, fmt.Errorf("scan season coverage row: %w", err) + } + r.ErrorMsg = errMsg.String + coverageRows = append(coverageRows, r) + } + return coverageRows, rows.Err() +} diff --git a/internal/store/coverage_test.go b/internal/store/coverage_test.go new file mode 100644 index 0000000..e7767e3 --- /dev/null +++ b/internal/store/coverage_test.go @@ -0,0 +1,134 @@ +package store + +import ( + "testing" +) + +func TestCoverageCRUD(t *testing.T) { + s := openTestStore(t) + + // Verify schema migration version is 4 (since we added 004_coverage.sql) + version, err := s.SchemaVersion() + if err != nil { + t.Fatalf("SchemaVersion() error = %v", err) + } + if version != 4 { + t.Fatalf("SchemaVersion() = %d, want 4", version) + } + + // Verify session_coverage table exists + var tableName string + err = s.db.QueryRow( + `SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'session_coverage'`, + ).Scan(&tableName) + if err != nil { + t.Fatalf("session_coverage table is missing: %v", err) + } + + // Insert meeting and session to test GetSeasonCoverage joins + meetingKey := 1234 + sessionKey := 5678 + year := 2024 + + if err := s.UpsertMeeting(Meeting{ + MeetingKey: meetingKey, + MeetingName: "Test Grand Prix", + Year: year, + DateStart: "2024-03-01T12:00:00Z", + }); err != nil { + t.Fatalf("UpsertMeeting() error = %v", err) + } + + if err := s.UpsertSession(Session{ + SessionKey: sessionKey, + MeetingKey: meetingKey, + SessionName: "Qualifying", + SessionType: "Qualifying", + DateStart: "2024-03-02T14:00:00Z", + }); err != nil { + t.Fatalf("UpsertSession() error = %v", err) + } + + // 1. Test empty session coverage + cov, err := s.GetSessionCoverage(sessionKey) + if err != nil { + t.Fatalf("GetSessionCoverage() error = %v", err) + } + if len(cov) != 0 { + t.Fatalf("Expected empty coverage, got: %v", cov) + } + + // 2. Upsert multiple datasets + err = s.UpsertCoverage(sessionKey, "drivers", "complete", 20, "") + if err != nil { + t.Fatalf("UpsertCoverage() drivers error = %v", err) + } + err = s.UpsertCoverage(sessionKey, "laps", "failed", 0, "429 Rate Limit") + if err != nil { + t.Fatalf("UpsertCoverage() laps error = %v", err) + } + + // 3. Retrieve session coverage + cov, err = s.GetSessionCoverage(sessionKey) + if err != nil { + t.Fatalf("GetSessionCoverage() error = %v", err) + } + if len(cov) != 2 { + t.Fatalf("Expected coverage size 2, got: %d", len(cov)) + } + + drv, ok := cov["drivers"] + if !ok { + t.Fatalf("Expected 'drivers' key to exist in coverage map") + } + if drv.Status != "complete" || drv.RowCount != 20 || drv.ErrorMsg != "" { + t.Fatalf("Unexpected drivers status: %+v", drv) + } + + laps, ok := cov["laps"] + if !ok { + t.Fatalf("Expected 'laps' key to exist in coverage map") + } + if laps.Status != "failed" || laps.RowCount != 0 || laps.ErrorMsg != "429 Rate Limit" { + t.Fatalf("Unexpected laps status: %+v", laps) + } + + // 4. Test updates (idempotency/upsert) + err = s.UpsertCoverage(sessionKey, "laps", "complete", 150, "") + if err != nil { + t.Fatalf("UpsertCoverage() second laps error = %v", err) + } + cov, err = s.GetSessionCoverage(sessionKey) + if err != nil { + t.Fatalf("GetSessionCoverage() error = %v", err) + } + laps = cov["laps"] + if laps.Status != "complete" || laps.RowCount != 150 || laps.ErrorMsg != "" { + t.Fatalf("Expected updated laps status to be complete with 150 rows, got: %+v", laps) + } + + // 5. Test GetSeasonCoverage joins + seasonRows, err := s.GetSeasonCoverage(year) + if err != nil { + t.Fatalf("GetSeasonCoverage() error = %v", err) + } + if len(seasonRows) != 2 { + t.Fatalf("Expected 2 season coverage rows, got: %d", len(seasonRows)) + } + + // The two rows should correspond to drivers and laps datasets + for _, row := range seasonRows { + if row.MeetingKey != meetingKey || row.MeetingName != "Test Grand Prix" { + t.Errorf("Unexpected meeting info: %+v", row) + } + if row.SessionKey != sessionKey || row.SessionName != "Qualifying" { + t.Errorf("Unexpected session info: %+v", row) + } + if row.Dataset != "drivers" && row.Dataset != "laps" { + t.Errorf("Unexpected dataset: %q", row.Dataset) + } + if row.Status != "complete" { + t.Errorf("Expected status to be 'complete', got %q", row.Status) + } + } +} diff --git a/internal/store/migrations/004_coverage.sql b/internal/store/migrations/004_coverage.sql new file mode 100644 index 0000000..169ef0e --- /dev/null +++ b/internal/store/migrations/004_coverage.sql @@ -0,0 +1,9 @@ +CREATE TABLE IF NOT EXISTS session_coverage ( + session_key INTEGER NOT NULL, + dataset TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'pending', + row_count INTEGER NOT NULL DEFAULT 0, + updated_at TEXT NOT NULL DEFAULT (datetime('now')), + error_msg TEXT, + PRIMARY KEY (session_key, dataset) +); diff --git a/internal/store/store_test.go b/internal/store/store_test.go index 4626185..9f338b7 100644 --- a/internal/store/store_test.go +++ b/internal/store/store_test.go @@ -28,8 +28,8 @@ func TestOpenAppliesMigrations(t *testing.T) { if err != nil { t.Fatalf("SchemaVersion() error = %v", err) } - if version != 3 { - t.Fatalf("SchemaVersion() = %d, want 3", version) + if version != 4 { + t.Fatalf("SchemaVersion() = %d, want 4", version) } tables := []string{ @@ -50,6 +50,7 @@ func TestOpenAppliesMigrations(t *testing.T) { "laps", "news_sources", "news_items", + "session_coverage", } for _, table := range tables { var name string @@ -92,6 +93,12 @@ func TestMigrationsAreIdempotent(t *testing.T) { if count != 1 { t.Fatalf("schema_migrations v3 count = %d, want 1", count) } + if err := s.db.QueryRow(`SELECT COUNT(*) FROM schema_migrations WHERE version = 4`).Scan(&count); err != nil { + t.Fatalf("count schema_migrations v4: %v", err) + } + if count != 1 { + t.Fatalf("schema_migrations v4 count = %d, want 1", count) + } } func TestRawPayloadInsertAndRead(t *testing.T) {