diff --git a/internal/api/cache.go b/internal/api/cache.go index 7057e0e..5ae2d5a 100644 --- a/internal/api/cache.go +++ b/internal/api/cache.go @@ -164,8 +164,8 @@ func (c *Cache) Get(key string) ([]byte, bool) { if ttl > 0 { age := time.Since(time.Unix(createdAt, 0)) if age > ttl { - // Expired — delete and return miss. - _, _ = c.db.Exec(`DELETE FROM cache WHERE key = ?`, key) + // Expired entries remain stored so get() can use them as a stale + // fallback if the live request fails. Prune() owns physical cleanup. atomic.AddInt64(&c.stats.Misses, 1) return nil, false } diff --git a/internal/api/client.go b/internal/api/client.go index fe04625..49d93b9 100644 --- a/internal/api/client.go +++ b/internal/api/client.go @@ -1,6 +1,7 @@ package api import ( + "context" "net/http" "sync" "sync/atomic" @@ -26,8 +27,12 @@ type requestPacer struct { // wait blocks until this caller's reserved slot arrives. func (p *requestPacer) wait() { + _ = p.waitContext(context.Background()) +} + +func (p *requestPacer) waitContext(ctx context.Context) error { if p == nil || p.interval <= 0 { - return + return nil } p.mu.Lock() now := time.Now() @@ -38,8 +43,18 @@ func (p *requestPacer) wait() { p.next = p.next.Add(p.interval) p.mu.Unlock() if sleep > 0 { - time.Sleep(sleep) + timer := time.NewTimer(sleep) + defer timer.Stop() + select { + case <-timer.C: + case <-ctx.Done(): + // Keep the unused reservation in the schedule. Blindly reclaiming an + // interval can collide with later callers that already reserved their + // wake times, releasing two requests simultaneously. + return ctx.Err() + } } + return nil } type OpenF1Client struct { @@ -56,6 +71,26 @@ type OpenF1Client struct { staleFlag int32 } +// Scoped returns a lightweight request-scoped view of the client. Network, +// pacing and cache resources are shared, while the stale fallback indicator is +// deliberately not shared. Web handlers use this view so a stale fallback in +// one concurrent HTTP request can never mark an unrelated response as stale. +// +// The legacy client-wide stale flag remains available for the TUI, whose loads +// are intentionally aggregated into one navigation-level notice. +func (c *OpenF1Client) Scoped() *OpenF1Client { + if c == nil { + return nil + } + return &OpenF1Client{ + url: c.url, + apiKey: c.apiKey, + httpClient: c.httpClient, + cache: c.cache, + pacer: c.pacer, + } +} + func NewOpenF1Client(url string, timeout time.Duration) *OpenF1Client { return &OpenF1Client{ url: url, diff --git a/internal/api/freshness_test.go b/internal/api/freshness_test.go new file mode 100644 index 0000000..16c4fa5 --- /dev/null +++ b/internal/api/freshness_test.go @@ -0,0 +1,123 @@ +package api + +import ( + "fmt" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + "time" +) + +func expireCacheEntry(t *testing.T, c *OpenF1Client, key string) { + t.Helper() + if _, err := c.cache.db.Exec(`UPDATE cache SET created_at = ? WHERE key = ?`, time.Now().Add(-48*time.Hour).Unix(), key); err != nil { + t.Fatal(err) + } +} + +func TestScopedClientReportsStaleFallbackWithoutMutatingParent(t *testing.T) { + year := time.Now().Year() + var fail atomic.Bool + upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if fail.Load() { + http.Error(w, "unavailable", http.StatusServiceUnavailable) + return + } + w.Header().Set("Content-Type", "application/json") + _, _ = fmt.Fprintf(w, `[{ + "meeting_key": 1, + "meeting_name": "British Grand Prix", + "year": %d + }]`, year) + })) + defer upstream.Close() + + client := NewOpenF1Client(upstream.URL, time.Second) + defer client.Close() + client.pacer.interval = 0 + key := fmt.Sprintf("%s/v1/meetings?year=%d", upstream.URL, year) + _, _ = client.cache.db.Exec(`DELETE FROM cache WHERE key = ?`, key) + defer func() { _, _ = client.cache.db.Exec(`DELETE FROM cache WHERE key = ?`, key) }() + + if _, err := client.GetMeetingsForYear(year); err != nil { + t.Fatalf("prime cache: %v", err) + } + expireCacheEntry(t, client, key) + fail.Store(true) + + scoped := client.Scoped() + meetings, err := scoped.GetMeetingsForYear(year) + if err != nil || len(meetings) != 1 { + t.Fatalf("stale fallback = (%+v, %v)", meetings, err) + } + if !scoped.LastResponseWasStale() { + t.Fatal("scoped request did not report its stale fallback") + } + if client.LastResponseWasStale() { + t.Fatal("request-scoped fallback leaked into the parent client") + } +} + +func TestScopedClientsDoNotLeakFreshnessAcrossConcurrentRequests(t *testing.T) { + year := time.Now().Year() + staleStarted := make(chan struct{}) + releaseStale := make(chan struct{}) + var failMeetings atomic.Bool + upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/v1/meetings": + if failMeetings.Load() { + close(staleStarted) + <-releaseStale + http.Error(w, "unavailable", http.StatusServiceUnavailable) + return + } + _, _ = fmt.Fprintf(w, `[{"meeting_key":1,"meeting_name":"British Grand Prix","year":%d}]`, year) + case "/v1/sessions": + _, _ = w.Write([]byte(`[{"session_key":11,"meeting_key":2,"session_name":"Race"}]`)) + default: + http.NotFound(w, r) + } + })) + defer upstream.Close() + + client := NewOpenF1Client(upstream.URL, 2*time.Second) + defer client.Close() + client.pacer.interval = 0 + meetingKey := fmt.Sprintf("%s/v1/meetings?year=%d", upstream.URL, year) + sessionKey := upstream.URL + "/v1/sessions?meeting_key=2" + _, _ = client.cache.db.Exec(`DELETE FROM cache WHERE key IN (?, ?)`, meetingKey, sessionKey) + defer func() { _, _ = client.cache.db.Exec(`DELETE FROM cache WHERE key IN (?, ?)`, meetingKey, sessionKey) }() + if _, err := client.GetMeetingsForYear(year); err != nil { + t.Fatal(err) + } + expireCacheEntry(t, client, meetingKey) + failMeetings.Store(true) + + staleClient := client.Scoped() + staleDone := make(chan error, 1) + go func() { + _, err := staleClient.GetMeetingsForYear(year) + staleDone <- err + }() + <-staleStarted + + freshClient := client.Scoped() + if _, err := freshClient.GetSessionsForMeeting(2); err != nil { + t.Fatalf("fresh concurrent request: %v", err) + } + if freshClient.LastResponseWasStale() { + t.Fatal("fresh request inherited concurrent request's stale state") + } + close(releaseStale) + if err := <-staleDone; err != nil { + t.Fatalf("stale request: %v", err) + } + if !staleClient.LastResponseWasStale() { + t.Fatal("stale request lost its own freshness state") + } + if freshClient.LastResponseWasStale() || client.LastResponseWasStale() { + t.Fatal("stale state leaked after concurrent requests completed") + } +} diff --git a/internal/api/openf1.go b/internal/api/openf1.go index d61874e..b0be5d4 100644 --- a/internal/api/openf1.go +++ b/internal/api/openf1.go @@ -2,6 +2,7 @@ package api import ( "bytes" + "context" "encoding/json" "errors" "fmt" @@ -59,8 +60,14 @@ func retryAfter429(resp *http.Response) time.Duration { // Without this, concurrent fan-outs (championship hub, track prefetch) burst // past the free-tier limit and callers silently treat 429s as missing data. func (c *OpenF1Client) doPaced(req *http.Request) (*http.Response, error) { + return c.doPacedContext(req.Context(), req) +} + +func (c *OpenF1Client) doPacedContext(ctx context.Context, req *http.Request) (*http.Response, error) { for attempt := 0; ; attempt++ { - c.pacer.wait() + if err := c.pacer.waitContext(ctx); err != nil { + return nil, err + } resp, err := c.httpClient.Do(req) if err != nil { return nil, err @@ -70,7 +77,14 @@ func (c *OpenF1Client) doPaced(req *http.Request) (*http.Response, error) { } delay := retryAfter429(resp) resp.Body.Close() - time.Sleep(delay) + timer := time.NewTimer(delay) + select { + case <-timer.C: + case <-ctx.Done(): + timer.Stop() + return nil, ctx.Err() + } + timer.Stop() } } @@ -83,13 +97,23 @@ func (c *OpenF1Client) doPaced(req *http.Request) (*http.Response, error) { // entry for this URL, that stale entry is returned instead of propagating the // error. The client's staleFlag is set so the UI can show a disclaimer. func (c *OpenF1Client) get(url string) (io.ReadCloser, error) { + return c.getContext(context.Background(), url) +} + +// getContext is the cancellable form used by bounded optional web enrichment. +// A caller cancellation never falls back to stale data: the work is no longer +// relevant to that response and must stop instead of continuing in background. +func (c *OpenF1Client) getContext(ctx context.Context, url string) (io.ReadCloser, error) { + if err := ctx.Err(); err != nil { + return nil, err + } // 1. Check the cache for a fresh (non-expired) entry. if cachedData, ok := c.cache.Get(url); ok { return io.NopCloser(bytes.NewReader(cachedData)), nil } // 2. Attempt a live network request. - req, err := http.NewRequest("GET", url, nil) + req, err := http.NewRequestWithContext(ctx, "GET", url, nil) if err != nil { // Even a request-construction failure warrants a stale fallback. return c.tryStale(url, err) @@ -98,8 +122,11 @@ func (c *OpenF1Client) get(url string) (io.ReadCloser, error) { req.Header.Set("Authorization", "Bearer "+c.apiKey) } - resp, err := c.doPaced(req) + resp, err := c.doPacedContext(ctx, req) if err != nil { + if ctx.Err() != nil { + return nil, ctx.Err() + } return c.tryStale(url, err) } defer resp.Body.Close() @@ -128,6 +155,9 @@ func (c *OpenF1Client) get(url string) (io.ReadCloser, error) { // 3. Success — read the body, store in cache, return. data, err := io.ReadAll(resp.Body) if err != nil { + if ctx.Err() != nil { + return nil, ctx.Err() + } return c.tryStale(url, err) } @@ -245,7 +275,13 @@ func (c *OpenF1Client) GetDriversForSession(sessionKey int) ([]models.Driver, er } func (c *OpenF1Client) GetDriver(sessionKey, driverNumber int) (*models.Driver, error) { - body, err := c.get(fmt.Sprintf("%s/v1/drivers?session_key=%d&driver_number=%d", c.url, sessionKey, driverNumber)) + return c.GetDriverContext(context.Background(), sessionKey, driverNumber) +} + +// GetDriverContext is a cancellable single-driver lookup for optional bounded +// enrichment. Other public methods retain their existing background semantics. +func (c *OpenF1Client) GetDriverContext(ctx context.Context, sessionKey, driverNumber int) (*models.Driver, error) { + body, err := c.getContext(ctx, fmt.Sprintf("%s/v1/drivers?session_key=%d&driver_number=%d", c.url, sessionKey, driverNumber)) if err != nil { return nil, err } diff --git a/internal/api/pacing_test.go b/internal/api/pacing_test.go index 284f819..689ea13 100644 --- a/internal/api/pacing_test.go +++ b/internal/api/pacing_test.go @@ -1,7 +1,9 @@ package api import ( + "context" "encoding/json" + "errors" "io" "net/http" "net/http/httptest" @@ -49,6 +51,70 @@ func TestRequestPacerNilSafe(t *testing.T) { p.wait() // must not panic } +func TestRequestPacerCancellationDoesNotCollideReservedWaiters(t *testing.T) { + const interval = 80 * time.Millisecond + p := &requestPacer{interval: interval} + if err := p.waitContext(context.Background()); err != nil { + t.Fatal(err) + } + p.mu.Lock() + initialNext := p.next + p.mu.Unlock() + + waitForReservation := func(want time.Time) { + t.Helper() + deadline := time.Now().Add(250 * time.Millisecond) + for time.Now().Before(deadline) { + p.mu.Lock() + got := p.next + p.mu.Unlock() + if got.Equal(want) { + return + } + time.Sleep(time.Millisecond) + } + t.Fatalf("reservation did not reach %v", want) + } + + ctxB, cancelB := context.WithCancel(context.Background()) + bDone := make(chan error, 1) + go func() { bDone <- p.waitContext(ctxB) }() + waitForReservation(initialNext.Add(interval)) + + cDone := make(chan time.Time, 1) + go func() { + _ = p.waitContext(context.Background()) + cDone <- time.Now() + }() + waitForReservation(initialNext.Add(2 * interval)) + + cancelStarted := time.Now() + cancelB() + select { + case err := <-bDone: + if !errors.Is(err, context.Canceled) { + t.Fatalf("B error = %v, want context.Canceled", err) + } + if elapsed := time.Since(cancelStarted); elapsed > 30*time.Millisecond { + t.Fatalf("B cancellation took %v", elapsed) + } + case <-time.After(50 * time.Millisecond): + t.Fatal("B did not return promptly after cancellation") + } + + dDone := make(chan time.Time, 1) + go func() { + _ = p.waitContext(context.Background()) + dDone <- time.Now() + }() + waitForReservation(initialNext.Add(3 * interval)) + + cAt, dAt := <-cDone, <-dDone + if separation := dAt.Sub(cAt); separation < interval/2 { + t.Fatalf("C and D collided: wake separation %v, want at least %v", separation, interval/2) + } +} + func TestGetRetriesOn429(t *testing.T) { var calls atomic.Int32 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { diff --git a/internal/query/context.go b/internal/query/context.go index 7c61e51..b876f61 100644 --- a/internal/query/context.go +++ b/internal/query/context.go @@ -40,6 +40,7 @@ type LiveEvidence struct { // ContextAvailability is structured source state for a referenced session. type ContextAvailability struct { + Source string `json:"source"` Schedule string `json:"schedule"` LiveTransport string `json:"live_transport"` LiveSession string `json:"live_session"` @@ -298,21 +299,28 @@ func meetingModelByKey(meetings []store.Meeting, key int) *models.Meeting { func sessionRef(c contextCandidate, evidence LiveEvidence, now time.Time) *ContextSession { session := sessionToModel(c.session) meeting := meetingToModel(c.meeting) - availability := ContextAvailability{Schedule: "available", LiveSession: "inactive", Archive: "unavailable", Freshness: "fresh", Limitations: []string{}} + // Schedule and analysis are domain-store facts. Without an ingestion + // timestamp the resolver cannot honestly call them network-fresh, so local + // is the baseline freshness vocabulary exposed to clients. + availability := ContextAvailability{Source: "local", Schedule: "available", LiveSession: "inactive", Archive: "unavailable", Freshness: "local", Limitations: []string{}} availability.LiveTransport = "unknown" if c.session.SessionKey == 0 { availability.Schedule = "unavailable" availability.Limitations = append(availability.Limitations, "schedule_identity_unmatched") } if evidence.Active && liveMatches(evidence, c.meeting, c.session) { + availability.Source = "mixed" availability.LiveTransport = "connected" availability.LiveSession = "active" + availability.Freshness = "live" if !evidence.ObservedAt.IsZero() { availability.ObservedAt = evidence.ObservedAt.Format(time.RFC3339) } } if c.archived { + availability.Source = "mixed" availability.Archive = "available" + availability.Freshness = "archive" if !evidence.ObservedAt.IsZero() { availability.ObservedAt = evidence.ObservedAt.Format(time.RFC3339) } @@ -326,6 +334,12 @@ func sessionRef(c contextCandidate, evidence LiveEvidence, now time.Time) *Conte } else { availability.LocalAnalysis = "pending" } + if availability.LocalAnalysis == "partial" && availability.Freshness == "local" { + availability.Freshness = "partial" + } + if c.session.SessionKey == 0 && availability.LiveSession == "active" { + availability.Source = "fia" + } return &ContextSession{Session: session, Meeting: &meeting, Availability: availability} } diff --git a/internal/query/context_test.go b/internal/query/context_test.go index 321aad3..7c4bac7 100644 --- a/internal/query/context_test.go +++ b/internal/query/context_test.go @@ -160,6 +160,54 @@ func TestResolveWeekendContextPassedTimeDoesNotCompleteSession(t *testing.T) { } } +func TestResolveWeekendContextAvailabilityUsesTruthfulSources(t *testing.T) { + now, _ := time.Parse(time.RFC3339, "2026-07-05T18:00:00Z") + svc := contextService(t, now) + addContextMeeting(t, svc, 1, "British Grand Prix", "2026-07-03T00:00:00Z", "2026-07-05T16:00:00Z", false) + addContextSession(t, svc, 11, 1, "Race", "2026-07-05T14:00:00Z", "2026-07-05T16:00:00Z", false) + // Results without laps/stints/positions are meaningful but incomplete local + // analysis, so they must not be labelled universally fresh. + completeContextSession(t, svc, 11, 1) + addContextMeeting(t, svc, 2, "Belgian Grand Prix", "2026-07-17T00:00:00Z", "2026-07-19T16:00:00Z", false) + addContextSession(t, svc, 21, 2, "Practice 1", "2026-07-17T09:00:00Z", "2026-07-17T10:00:00Z", false) + addContextSession(t, svc, 22, 2, "Race", "2026-07-19T14:00:00Z", "2026-07-19T16:00:00Z", false) + + local, err := svc.ResolveWeekendContext(LiveEvidence{}) + if err != nil { + t.Fatal(err) + } + if got := local.PreviousCompletedSession.Availability; got.Source != "local" || got.Freshness != "partial" || got.LocalAnalysis != "partial" { + t.Fatalf("partial local availability = %+v", got) + } + if got := local.NextSession.Availability; got.Source != "local" || got.Freshness != "local" { + t.Fatalf("future local availability = %+v", got) + } + + liveContext, err := svc.ResolveWeekendContext(LiveEvidence{Active: true, MeetingName: "Belgian Grand Prix", CircuitName: "Belgian Grand Prix", SessionName: "Practice 1", SessionType: "Practice 1", ObservedAt: now}) + if err != nil { + t.Fatal(err) + } + if got := liveContext.ActiveSession.Availability; got.Source != "mixed" || got.Freshness != "live" || got.LiveSession != "active" { + t.Fatalf("FIA + local availability = %+v", got) + } + + archiveContext, err := svc.ResolveWeekendContext(LiveEvidence{Final: true, MeetingName: "British Grand Prix", CircuitName: "British Grand Prix", SessionName: "Race", SessionType: "Race", ObservedAt: now}) + if err != nil { + t.Fatal(err) + } + if got := archiveContext.PreviousCompletedSession.Availability; got.Source != "mixed" || got.Freshness != "archive" || got.Archive != "available" { + t.Fatalf("FIA archive + local availability = %+v", got) + } + + synthetic, err := svc.ResolveWeekendContext(LiveEvidence{Active: true, MeetingName: "Unscheduled Grand Prix", SessionName: "Race", SessionType: "Race", ObservedAt: now}) + if err != nil { + t.Fatal(err) + } + if got := synthetic.ActiveSession.Availability; got.Source != "fia" || got.Freshness != "live" || got.Schedule != "unavailable" { + t.Fatalf("synthetic FIA availability = %+v", got) + } +} + func TestResolveWeekendContextNeverUsesFutureAnalysis(t *testing.T) { now, _ := time.Parse(time.RFC3339, "2026-07-01T12:00:00Z") svc := contextService(t, now) diff --git a/internal/web/api.go b/internal/web/api.go index fca561e..de21b7a 100644 --- a/internal/web/api.go +++ b/internal/web/api.go @@ -16,6 +16,7 @@ import ( readability "codeberg.org/readeck/go-readability/v2" + "github.com/AmanTahiliani/box-box/internal/api" "github.com/AmanTahiliani/box-box/internal/models" "github.com/AmanTahiliani/box-box/internal/query" ) @@ -43,6 +44,7 @@ func (s *Server) handleMeetings(w http.ResponseWriter, r *http.Request) { switch parseSourceMode(r) { case sourceLocal: + markLocalResponse(w, false) if !s.hasLocalQuery() { writeJSON(w, []models.Meeting{}) return @@ -62,17 +64,20 @@ func (s *Server) handleMeetings(w http.ResponseWriter, r *http.Request) { return } if len(meetings) > 0 { + markLocalResponse(w, false) writeJSON(w, meetings) return } } } - meetings, err := s.client.GetMeetingsForYear(year) + client := s.client.Scoped() + meetings, err := client.GetMeetingsForYear(year) if err != nil { - writeError(w, err, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return } + markOpenF1Response(w, client) writeJSON(w, meetings) } @@ -87,6 +92,7 @@ func (s *Server) handleSessions(w http.ResponseWriter, r *http.Request) { switch parseSourceMode(r) { case sourceLocal: + markLocalResponse(w, false) if !s.hasLocalQuery() { writeJSON(w, []models.Session{}) return @@ -106,23 +112,27 @@ func (s *Server) handleSessions(w http.ResponseWriter, r *http.Request) { return } if len(sessions) > 0 { + markLocalResponse(w, false) writeJSON(w, sessions) return } } } - sessions, err := s.client.GetSessionsForMeeting(meetingKey) + client := s.client.Scoped() + sessions, err := client.GetSessionsForMeeting(meetingKey) if err != nil { - writeError(w, err, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return } + markOpenF1Response(w, client) writeJSON(w, sessions) } // --- /api/v1/news --- func (s *Server) handleNews(w http.ResponseWriter, r *http.Request) { + markLocalResponse(w, false) if !s.hasLocalQuery() { writeJSON(w, []query.NewsItem{}) return @@ -223,6 +233,7 @@ func (s *Server) handleDrivers(w http.ResponseWriter, r *http.Request) { switch parseSourceMode(r) { case sourceLocal: + markLocalResponse(w, false) if !s.hasLocalQuery() { writeJSON(w, []models.Driver{}) return @@ -242,6 +253,7 @@ func (s *Server) handleDrivers(w http.ResponseWriter, r *http.Request) { if s.hasLocalQuery() { drivers, err := s.query.ListDrivers(sessionKey) if err == nil && len(drivers) > 0 { + markLocalResponse(w, false) writeJSON(w, drivers) return } @@ -252,11 +264,13 @@ func (s *Server) handleDrivers(w http.ResponseWriter, r *http.Request) { } } - drivers, err := s.client.GetDriversForSession(sessionKey) + client := s.client.Scoped() + drivers, err := client.GetDriversForSession(sessionKey) if err != nil { - writeError(w, err, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return } + markOpenF1Response(w, client) writeJSON(w, drivers) } @@ -279,6 +293,7 @@ func (s *Server) handleResults(w http.ResponseWriter, r *http.Request) { switch parseSourceMode(r) { case sourceLocal: + markLocalResponse(w, false) if !s.hasLocalQuery() { writeJSON(w, []resultWithDriver{}) return @@ -294,6 +309,7 @@ func (s *Server) handleResults(w http.ResponseWriter, r *http.Request) { if s.hasLocalQuery() { results, err := s.query.ListResults(sessionKey) if err == nil && len(results) > 0 { + markLocalResponse(w, false) writeJSON(w, enrichedResultsToAPI(results)) return } @@ -308,30 +324,42 @@ func (s *Server) handleResults(w http.ResponseWriter, r *http.Request) { results []models.SessionResult drivers []models.Driver resultsErr error + driversErr error wg sync.WaitGroup ) + client := s.client.Scoped() wg.Add(2) - go func() { defer wg.Done(); results, resultsErr = s.client.GetSessionResult(sessionKey) }() - go func() { defer wg.Done(); drivers, _ = s.client.GetDriversForSession(sessionKey) }() + go func() { defer wg.Done(); results, resultsErr = client.GetSessionResult(sessionKey) }() + go func() { defer wg.Done(); drivers, driversErr = client.GetDriversForSession(sessionKey) }() wg.Wait() if resultsErr != nil { - writeError(w, resultsErr, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, resultsErr, http.StatusInternalServerError, client.LastResponseWasStale()) return } driverMap := buildDriverMap(drivers) + incomplete := driversErr != nil enriched := make([]resultWithDriver, 0, len(results)) for _, res := range results { e := resultWithDriver{SessionResult: res} - if d, ok := driverMap[res.DriverNumber]; ok { + if d, ok := driverMap[res.DriverNumber]; ok && hasDriverPresentation(d) { e.NameAcronym = d.NameAcronym e.FullName = d.FullName e.TeamName = d.TeamName e.TeamColour = d.TeamColour + } else { + incomplete = true } enriched = append(enriched, e) } + resultsFreshness := "fresh" + if len(results) == 0 { + resultsFreshness = "limited" + } else if incomplete { + resultsFreshness = "partial" + } + markOpenF1Availability(w, client, resultsFreshness) writeJSON(w, enriched) } @@ -354,6 +382,7 @@ func (s *Server) handleGrid(w http.ResponseWriter, r *http.Request) { switch parseSourceMode(r) { case sourceLocal: + markLocalResponse(w, false) if !s.hasLocalQuery() { writeJSON(w, []gridWithDriver{}) return @@ -369,6 +398,7 @@ func (s *Server) handleGrid(w http.ResponseWriter, r *http.Request) { if s.hasLocalQuery() { grid, err := s.query.ListStartingGrid(sessionKey) if err == nil && len(grid) > 0 { + markLocalResponse(w, false) writeJSON(w, enrichedGridToAPI(grid)) return } @@ -380,33 +410,45 @@ func (s *Server) handleGrid(w http.ResponseWriter, r *http.Request) { } var ( - grid []models.StartingGrid - drivers []models.Driver - gridErr error - wg sync.WaitGroup + grid []models.StartingGrid + drivers []models.Driver + gridErr error + driversErr error + wg sync.WaitGroup ) + client := s.client.Scoped() wg.Add(2) - go func() { defer wg.Done(); grid, gridErr = s.client.GetStartingGrid(sessionKey) }() - go func() { defer wg.Done(); drivers, _ = s.client.GetDriversForSession(sessionKey) }() + go func() { defer wg.Done(); grid, gridErr = client.GetStartingGrid(sessionKey) }() + go func() { defer wg.Done(); drivers, driversErr = client.GetDriversForSession(sessionKey) }() wg.Wait() if gridErr != nil { - writeError(w, gridErr, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, gridErr, http.StatusInternalServerError, client.LastResponseWasStale()) return } driverMap := buildDriverMap(drivers) + incomplete := driversErr != nil enriched := make([]gridWithDriver, 0, len(grid)) for _, g := range grid { e := gridWithDriver{StartingGrid: g} - if d, ok := driverMap[g.DriverNumber]; ok { + if d, ok := driverMap[g.DriverNumber]; ok && hasDriverPresentation(d) { e.NameAcronym = d.NameAcronym e.FullName = d.FullName e.TeamName = d.TeamName e.TeamColour = d.TeamColour + } else { + incomplete = true } enriched = append(enriched, e) } + gridFreshness := "fresh" + if len(grid) == 0 { + gridFreshness = "limited" + } else if incomplete { + gridFreshness = "partial" + } + markOpenF1Availability(w, client, gridFreshness) writeJSON(w, enriched) } @@ -419,26 +461,29 @@ func (s *Server) handleLaps(w http.ResponseWriter, r *http.Request) { return } + client := s.client.Scoped() if dnStr := r.URL.Query().Get("driver_number"); dnStr != "" { driverNumber, err := strconv.Atoi(dnStr) if err != nil || driverNumber == 0 { http.Error(w, "invalid driver_number", http.StatusBadRequest) return } - laps, err := s.client.GetLapsForDriver(sessionKey, driverNumber) + laps, err := client.GetLapsForDriver(sessionKey, driverNumber) if err != nil { - writeError(w, err, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return } + markOpenF1Response(w, client) writeJSON(w, laps) return } - laps, err := s.client.GetLapsForSession(sessionKey) + laps, err := client.GetLapsForSession(sessionKey) if err != nil { - writeError(w, err, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return } + markOpenF1Response(w, client) writeJSON(w, laps) } @@ -450,11 +495,13 @@ func (s *Server) handleWeather(w http.ResponseWriter, r *http.Request) { http.Error(w, "session_key required", http.StatusBadRequest) return } - weather, err := s.client.GetWeather(sessionKey) + client := s.client.Scoped() + weather, err := client.GetWeather(sessionKey) if err != nil { - writeError(w, err, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return } + markOpenF1Response(w, client) writeJSON(w, weather) } @@ -466,11 +513,13 @@ func (s *Server) handleRaceControl(w http.ResponseWriter, r *http.Request) { http.Error(w, "session_key required", http.StatusBadRequest) return } - rc, err := s.client.GetRaceControl(sessionKey) + client := s.client.Scoped() + rc, err := client.GetRaceControl(sessionKey) if err != nil { - writeError(w, err, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return } + markOpenF1Response(w, client) writeJSON(w, rc) } @@ -487,11 +536,13 @@ func (s *Server) handleTelemetry(w http.ResponseWriter, r *http.Request) { http.Error(w, "driver_number required", http.StatusBadRequest) return } - carData, err := s.client.GetCarData(sessionKey, driverNumber) + client := s.client.Scoped() + carData, err := client.GetCarData(sessionKey, driverNumber) if err != nil { - writeError(w, err, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return } + markOpenF1Response(w, client) writeJSON(w, carData) } @@ -503,11 +554,13 @@ func (s *Server) handleOvertakes(w http.ResponseWriter, r *http.Request) { http.Error(w, "session_key required", http.StatusBadRequest) return } - overtakes, err := s.client.GetOvertakesForSession(sessionKey) + client := s.client.Scoped() + overtakes, err := client.GetOvertakesForSession(sessionKey) if err != nil { - writeError(w, err, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return } + markOpenF1Response(w, client) writeJSON(w, overtakes) } @@ -524,11 +577,13 @@ func (s *Server) handleTeamRadio(w http.ResponseWriter, r *http.Request) { http.Error(w, "driver_number required", http.StatusBadRequest) return } - radios, err := s.client.GetTeamRadio(sessionKey, driverNumber) + client := s.client.Scoped() + radios, err := client.GetTeamRadio(sessionKey, driverNumber) if err != nil { - writeError(w, err, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return } + markOpenF1Response(w, client) writeJSON(w, radios) } @@ -547,36 +602,42 @@ func (s *Server) handleChampionshipDrivers(w http.ResponseWriter, r *http.Reques if year == 0 { year = time.Now().Year() } - champ, err := s.client.GetDriverChampionshipForYear(year) + client := s.client.Scoped() + champ, err := client.GetDriverChampionshipForYear(year) if err != nil { - writeError(w, err, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return } if len(champ) == 0 { + markOpenF1Availability(w, client, "limited") writeJSON(w, []any{}) return } - drivers, _ := s.client.GetDriversForSession(champ[0].SessionKey) + drivers, driversErr := client.GetDriversForSession(champ[0].SessionKey) driverMap := buildDriverMapFirst(drivers) + incomplete := driversErr != nil enriched := make([]champDriverWithInfo, 0, len(champ)) for _, c := range champ { e := champDriverWithInfo{ChampionshipDriver: c} - d, ok := s.championshipDriverInfo(c.SessionKey, c.DriverNumber, driverMap) - if ok { + d, ok := championshipDriverInfo(client, c.SessionKey, c.DriverNumber, driverMap) + if ok && hasDriverPresentation(d) { e.NameAcronym = d.NameAcronym e.FullName = d.FullName e.TeamName = d.TeamName e.TeamColour = d.TeamColour + } else { + incomplete = true } enriched = append(enriched, e) } + markOpenF1AggregateResponse(w, client, incomplete) writeJSON(w, enriched) } -func (s *Server) championshipDriverInfo(sessionKey, driverNumber int, fallback map[int]models.Driver) (models.Driver, bool) { - if d, err := s.client.GetDriver(sessionKey, driverNumber); err == nil && d != nil { +func championshipDriverInfo(client *api.OpenF1Client, sessionKey, driverNumber int, fallback map[int]models.Driver) (models.Driver, bool) { + if d, err := client.GetDriver(sessionKey, driverNumber); err == nil && d != nil { return *d, true } d, ok := fallback[driverNumber] @@ -590,11 +651,17 @@ func (s *Server) handleChampionshipTeams(w http.ResponseWriter, r *http.Request) if year == 0 { year = time.Now().Year() } - teams, err := s.client.GetTeamChampionshipForYear(year) + client := s.client.Scoped() + teams, err := client.GetTeamChampionshipForYear(year) if err != nil { - writeError(w, err, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return } + if len(teams) == 0 { + markOpenF1Availability(w, client, "limited") + } else { + markOpenF1Response(w, client) + } writeJSON(w, teams) } @@ -671,8 +738,10 @@ func champHubTTL(year int, now time.Time) time.Duration { } type champHubEntry struct { - resp champHubResponse - expires time.Time + resp champHubResponse + source string + freshness string + expires time.Time } // champHubCache is an in-memory cache of aggregated hub responses keyed by @@ -693,12 +762,26 @@ func (c *champHubCache) get(year int, now time.Time) (champHubResponse, bool) { } func (c *champHubCache) put(year int, resp champHubResponse, now time.Time, ttl time.Duration) { + c.putWithMetadata(year, resp, "local", "local", now, ttl) +} + +func (c *champHubCache) getWithMetadata(year int, now time.Time) (champHubResponse, string, string, bool) { + c.mu.Lock() + defer c.mu.Unlock() + e, ok := c.entries[year] + if !ok || now.After(e.expires) { + return champHubResponse{}, "", "", false + } + return e.resp, e.source, e.freshness, true +} + +func (c *champHubCache) putWithMetadata(year int, resp champHubResponse, source, freshness string, now time.Time, ttl time.Duration) { c.mu.Lock() defer c.mu.Unlock() if c.entries == nil { c.entries = map[int]champHubEntry{} } - c.entries[year] = champHubEntry{resp: resp, expires: now.Add(ttl)} + c.entries[year] = champHubEntry{resp: resp, source: source, freshness: freshness, expires: now.Add(ttl)} } // fetchMeetingRaces fans fetch out across meetings with bounded concurrency. @@ -740,13 +823,6 @@ func (s *Server) handleChampionshipHub(w http.ResponseWriter, r *http.Request) { } mode := parseSourceMode(r) - if mode != sourceLocal { - if resp, ok := s.hubCache.get(year, time.Now()); ok { - writeJSON(w, resp) - return - } - } - if mode == sourceLocal || mode == sourceAuto { resp, ok, err := s.localChampionshipHub(year) if err != nil { @@ -754,21 +830,33 @@ func (s *Server) handleChampionshipHub(w http.ResponseWriter, r *http.Request) { return } if ok { - s.hubCache.put(year, resp, time.Now(), champHubTTL(year, time.Now())) + s.hubCache.putWithMetadata(year, resp, "local", "local", time.Now(), champHubTTL(year, time.Now())) + markLocalResponse(w, false) writeJSON(w, resp) return } if mode == sourceLocal { + markDataResponse(w, "none", "limited") writeJSON(w, resp) return } } - resp, err := s.openF1ChampionshipHub(year) - if err != nil { - writeError(w, err, http.StatusInternalServerError, s.client.LastResponseWasStale()) + // At this point an auto request has no usable domain aggregate and an + // explicit OpenF1 request must not be satisfied by a local cache entry. + if resp, source, freshness, ok := s.hubCache.getWithMetadata(year, time.Now()); ok && source == "openf1" { + markDataResponse(w, source, freshness) + writeJSON(w, resp) return } + + client := s.client.Scoped() + resp, incomplete, err := s.openF1ChampionshipHub(client, year) + if err != nil { + writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) + return + } + markOpenF1AggregateResponse(w, client, incomplete) writeJSON(w, resp) } @@ -797,25 +885,35 @@ func (s *Server) localChampionshipHub(year int) (champHubResponse, bool, error) return aggregateChampionshipHub(year, races, inputs.Champ, inputs.Teams, inputs.DriverMap), true, nil } -func (s *Server) openF1ChampionshipHub(year int) (champHubResponse, error) { - champ, err := s.client.GetDriverChampionshipForYear(year) +func (s *Server) openF1ChampionshipHub(client *api.OpenF1Client, year int) (champHubResponse, bool, error) { + champ, err := client.GetDriverChampionshipForYear(year) if err != nil { - return champHubResponse{}, err + return champHubResponse{}, false, err } if len(champ) == 0 { - return champHubResponse{Season: year, RoundLabels: []string{}, Drivers: []champHubDriver{}, Teams: []champHubTeam{}}, nil + return champHubResponse{Season: year, RoundLabels: []string{}, Drivers: []champHubDriver{}, Teams: []champHubTeam{}}, true, nil } - teams, _ := s.client.GetTeamChampionshipForYear(year) + teams, teamsErr := client.GetTeamChampionshipForYear(year) driverInfo := map[int]models.Driver{} - if ds, derr := s.client.GetDriversForSession(champ[0].SessionKey); derr == nil { + driversIncomplete := false + if ds, derr := client.GetDriversForSession(champ[0].SessionKey); derr == nil { driverInfo = buildDriverMapFirst(ds) + } else { + driversIncomplete = true + } + for _, standing := range champ { + if !hasDriverPresentation(driverInfo[standing.DriverNumber]) { + driversIncomplete = true + break + } } - races, incomplete, err := s.fetchSeasonRaces(year) + races, incomplete, err := fetchSeasonRaces(client, year) if err != nil { - return champHubResponse{}, err + return champHubResponse{}, false, err } + incomplete = incomplete || teamsErr != nil || driversIncomplete resp := aggregateChampionshipHub(year, races, champ, teams, driverInfo) ttl := champHubTTL(year, time.Now()) @@ -825,16 +923,22 @@ func (s *Server) openF1ChampionshipHub(year int) (champHubResponse, error) { // so a partial view of the season doesn't stick around for the full TTL. ttl = champHubIncompleteTTL } - s.hubCache.put(year, resp, time.Now(), ttl) - return resp, nil + freshness := "fresh" + if client.LastResponseWasStale() { + freshness = "stale" + } else if incomplete { + freshness = "partial" + } + s.hubCache.putWithMetadata(year, resp, "openf1", freshness, time.Now(), ttl) + return resp, incomplete, nil } // fetchSeasonRaces returns a season's GP meetings in date order, each bundled // with its race results and starting grid fetched from OpenF1. incomplete // reports whether any per-meeting fetch failed, so callers can avoid caching a // partial view of the season for long. -func (s *Server) fetchSeasonRaces(year int) (races []meetingRace, incomplete bool, err error) { - meetings, err := s.client.GetMeetingsForYear(year) +func fetchSeasonRaces(client *api.OpenF1Client, year int) (races []meetingRace, incomplete bool, err error) { + meetings, err := client.GetMeetingsForYear(year) if err != nil { return nil, false, err } @@ -842,7 +946,7 @@ func (s *Server) fetchSeasonRaces(year int) (races []meetingRace, incomplete boo var failed atomic.Bool races = fetchMeetingRaces(meetings, champHubWorkers, func(m models.Meeting) (meetingRace, bool) { - sessions, serr := s.client.GetSessionsForMeeting(int(m.MeetingKey)) + sessions, serr := client.GetSessionsForMeeting(int(m.MeetingKey)) if serr != nil { failed.Store(true) return meetingRace{}, false @@ -855,10 +959,17 @@ func (s *Server) fetchSeasonRaces(year int) (races []meetingRace, incomplete boo } } if raceKey == 0 { - return meetingRace{}, false // not a GP meeting (e.g. pre-season testing) + if isKnownNonChampionshipMeeting(m, sessions) { + return meetingRace{}, false + } + failed.Store(true) + // The meeting list does not identify non-championship events. Skipping + // a meeting without a Race may be expected (testing), but the aggregate + // is not proven complete and must be labelled partial. + return meetingRace{}, false } - results, rerr := s.client.GetSessionResult(raceKey) - grid, gerr := s.client.GetStartingGrid(raceKey) + results, rerr := client.GetSessionResult(raceKey) + grid, gerr := client.GetStartingGrid(raceKey) if rerr != nil || gerr != nil { failed.Store(true) } @@ -867,6 +978,29 @@ func (s *Server) fetchSeasonRaces(year int) (races []meetingRace, incomplete boo return races, failed.Load(), nil } +func isKnownNonChampionshipMeeting(meeting models.Meeting, sessions []models.Session) bool { + if hasTestingToken(meeting.MeetingName + " " + meeting.MeetingOfficialName) { + return true + } + for _, session := range sessions { + if hasTestingToken(session.SessionName + " " + session.SessionType) { + return true + } + } + return false +} + +func hasTestingToken(value string) bool { + for _, token := range strings.FieldsFunc(strings.ToLower(value), func(r rune) bool { + return (r < 'a' || r > 'z') && (r < '0' || r > '9') + }) { + if token == "test" || token == "tests" || token == "testing" { + return true + } + } + return false +} + // aggregateChampionshipHub is the pure aggregation core (no network) so it can be // unit-tested with synthetic data. races must be ordered ascending by date and // contain only GP meetings (those with a Race session). @@ -1077,6 +1211,7 @@ type trackOutlineResponse struct { } func (s *Server) handleTrackOutline(w http.ResponseWriter, r *http.Request) { + markLocalResponse(w, false) year, _ := strconv.Atoi(r.URL.Query().Get("year")) if year == 0 { year = time.Now().Year() @@ -1273,22 +1408,25 @@ func (s *Server) handleStrategy(w http.ResponseWriter, r *http.Request) { } var ( - stints []models.Stint - pits []models.Pit - results []models.SessionResult - drivers []models.Driver - rc []models.RaceControl - stintsErr error - pitsErr error - resErr error - wg sync.WaitGroup + stints []models.Stint + pits []models.Pit + results []models.SessionResult + drivers []models.Driver + rc []models.RaceControl + stintsErr error + pitsErr error + resErr error + driversErr error + rcErr error + wg sync.WaitGroup ) + client := s.client.Scoped() wg.Add(5) - go func() { defer wg.Done(); stints, stintsErr = s.client.GetStintsForSession(sessionKey) }() - go func() { defer wg.Done(); pits, pitsErr = s.client.GetPitStopsForSession(sessionKey) }() - go func() { defer wg.Done(); results, resErr = s.client.GetSessionResult(sessionKey) }() - go func() { defer wg.Done(); drivers, _ = s.client.GetDriversForSession(sessionKey) }() - go func() { defer wg.Done(); rc, _ = s.client.GetRaceControl(sessionKey) }() + go func() { defer wg.Done(); stints, stintsErr = client.GetStintsForSession(sessionKey) }() + go func() { defer wg.Done(); pits, pitsErr = client.GetPitStopsForSession(sessionKey) }() + go func() { defer wg.Done(); results, resErr = client.GetSessionResult(sessionKey) }() + go func() { defer wg.Done(); drivers, driversErr = client.GetDriversForSession(sessionKey) }() + go func() { defer wg.Done(); rc, rcErr = client.GetRaceControl(sessionKey) }() wg.Wait() if stintsErr != nil || pitsErr != nil || resErr != nil { @@ -1299,17 +1437,21 @@ func (s *Server) handleStrategy(w http.ResponseWriter, r *http.Request) { if e == nil { e = resErr } - writeError(w, e, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, e, http.StatusInternalServerError, client.LastResponseWasStale()) return } - // Non-race sessions have no stints. + // Empty strategy data may mean a non-race session or a race still settling. if len(stints) == 0 { + // Without session-type evidence, an empty primary strategy dataset is + // not enough to prove "not applicable" (it may still be settling). + markOpenF1Availability(w, client, "limited") writeJSON(w, map[string]any{"note": "Not applicable", "drivers": []any{}}) return } driverMap := buildDriverMap(drivers) + incomplete := driversErr != nil || rcErr != nil resultMap := make(map[int]models.SessionResult, len(results)) totalLaps := 0 @@ -1342,6 +1484,9 @@ func (s *Server) handleStrategy(w http.ResponseWriter, r *http.Request) { stratDrivers := make([]strategyDriver, 0, len(seenDrivers)) for dn := range seenDrivers { d := driverMap[dn] + if !hasDriverPresentation(d) { + incomplete = true + } res := resultMap[dn] sd := strategyDriver{ @@ -1396,6 +1541,7 @@ func (s *Server) handleStrategy(w http.ResponseWriter, r *http.Request) { return pi < pj }) + markOpenF1AggregateResponse(w, client, incomplete) writeJSON(w, strategyResponse{ SessionKey: sessionKey, TotalLaps: totalLaps, @@ -1491,21 +1637,31 @@ func (s *Server) handleLapsComparison(w http.ResponseWriter, r *http.Request) { } var ( - allLaps []models.Lap - stints []models.Stint - pits []models.Pit - rc []models.RaceControl - wg sync.WaitGroup + allLaps []models.Lap + stints []models.Stint + pits []models.Pit + rc []models.RaceControl + lapsErr error + stintsErr error + pitsErr error + rcErr error + wg sync.WaitGroup ) + client := s.client.Scoped() wg.Add(4) - go func() { defer wg.Done(); allLaps, _ = s.client.GetLapsForSession(sessionKey) }() - go func() { defer wg.Done(); stints, _ = s.client.GetStintsForSession(sessionKey) }() - go func() { defer wg.Done(); pits, _ = s.client.GetPitStopsForSession(sessionKey) }() - go func() { defer wg.Done(); rc, _ = s.client.GetRaceControl(sessionKey) }() + go func() { defer wg.Done(); allLaps, lapsErr = client.GetLapsForSession(sessionKey) }() + go func() { defer wg.Done(); stints, stintsErr = client.GetStintsForSession(sessionKey) }() + go func() { defer wg.Done(); pits, pitsErr = client.GetPitStopsForSession(sessionKey) }() + go func() { defer wg.Done(); rc, rcErr = client.GetRaceControl(sessionKey) }() wg.Wait() + if lapsErr != nil { + writeError(w, lapsErr, http.StatusInternalServerError, client.LastResponseWasStale()) + return + } - allDrivers, _ := s.client.GetDriversForSession(sessionKey) + allDrivers, driversErr := client.GetDriversForSession(sessionKey) driverMap := buildDriverMap(allDrivers) + incomplete := stintsErr != nil || pitsErr != nil || rcErr != nil || driversErr != nil // If no filter, default to first 3 unique driver numbers from lap data. if len(requestedDrivers) == 0 { @@ -1543,6 +1699,9 @@ func (s *Server) handleLapsComparison(w http.ResponseWriter, r *http.Request) { compDrivers := make([]comparisonDriver, 0, len(requestedDrivers)) for _, dn := range requestedDrivers { d := driverMap[dn] + if !hasDriverPresentation(d) { + incomplete = true + } cd := comparisonDriver{ DriverNumber: dn, NameAcronym: d.NameAcronym, @@ -1558,6 +1717,13 @@ func (s *Server) handleLapsComparison(w http.ResponseWriter, r *http.Request) { compDrivers = append(compDrivers, cd) } + freshness := "fresh" + if len(allLaps) == 0 { + freshness = "limited" + } else if incomplete { + freshness = "partial" + } + markOpenF1Availability(w, client, freshness) writeJSON(w, lapsComparisonResponse{ SessionKey: sessionKey, SCPeriods: extractSCPeriods(rc), diff --git a/internal/web/championship_hub_test.go b/internal/web/championship_hub_test.go index dc005b0..05d63fd 100644 --- a/internal/web/championship_hub_test.go +++ b/internal/web/championship_hub_test.go @@ -1,13 +1,136 @@ package web import ( + "fmt" + "net/http" + "net/http/httptest" "sync/atomic" "testing" "time" + "github.com/AmanTahiliani/box-box/internal/api" "github.com/AmanTahiliani/box-box/internal/models" ) +func championshipTestUpstream(t *testing.T, driversOK, meetingHasRace bool) *httptest.Server { + t.Helper() + completed := time.Now().Add(-time.Hour).UTC().Format(time.RFC3339) + return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + switch r.URL.Path { + case "/v1/sessions": + if r.URL.Query().Get("session_name") == "Race" { + _, _ = fmt.Fprintf(w, `[{"session_key":99,"session_name":"Race","date_end":%q}]`, completed) + return + } + if meetingHasRace { + _, _ = w.Write([]byte(`[{"session_key":101,"meeting_key":1,"session_name":"Race"}]`)) + } else { + _, _ = w.Write([]byte(`[{"session_key":100,"meeting_key":1,"session_name":"Practice 1"}]`)) + } + case "/v1/championship_drivers": + _, _ = w.Write([]byte(`[{"driver_number":1,"session_key":99,"position_current":1,"points_current":25}]`)) + case "/v1/championship_teams": + _, _ = w.Write([]byte(`[{"team_name":"Red Bull","position_current":1,"points_current":25}]`)) + case "/v1/drivers": + if !driversOK { + http.Error(w, "identity unavailable", http.StatusBadGateway) + return + } + _, _ = w.Write([]byte(`[{"driver_number":1,"name_acronym":"VER","full_name":"Max Verstappen","team_name":"Red Bull","team_colour":"3671c6"}]`)) + case "/v1/meetings": + _, _ = w.Write([]byte(`[{"meeting_key":1,"meeting_name":"Mystery Grand Prix"}]`)) + case "/v1/session_result": + _, _ = w.Write([]byte(`[{"driver_number":1,"position":1,"points":25}]`)) + case "/v1/starting_grid": + _, _ = w.Write([]byte(`[{"driver_number":1,"position":1}]`)) + default: + http.NotFound(w, r) + } + })) +} + +func TestOpenF1ChampionshipHubIdentityFailureIsPartialAndCached(t *testing.T) { + upstream := championshipTestUpstream(t, false, true) + defer upstream.Close() + client := api.NewOpenF1Client(upstream.URL, 2*time.Second) + defer client.Close() + server := NewServer(client, 0, nil) + year := time.Now().Year() + + _, incomplete, err := server.openF1ChampionshipHub(client.Scoped(), year) + if err != nil { + t.Fatal(err) + } + if !incomplete { + t.Fatal("missing championship driver identity was labelled complete") + } + _, source, freshness, ok := server.hubCache.getWithMetadata(year, time.Now()) + if !ok || source != "openf1" || freshness != "partial" { + t.Fatalf("cached metadata = hit %v, %q/%q", ok, source, freshness) + } +} + +func TestFetchSeasonRacesMeetingWithoutRaceIsIncomplete(t *testing.T) { + upstream := championshipTestUpstream(t, true, false) + defer upstream.Close() + client := api.NewOpenF1Client(upstream.URL, 2*time.Second) + defer client.Close() + + races, incomplete, err := fetchSeasonRaces(client.Scoped(), time.Now().Year()) + if err != nil { + t.Fatal(err) + } + if !incomplete || len(races) != 0 { + t.Fatalf("no-Race meeting = races %d, incomplete %v", len(races), incomplete) + } +} + +func TestFetchSeasonRacesRecognizedTestingMeetingIsNotIncomplete(t *testing.T) { + upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/v1/meetings": + _, _ = w.Write([]byte(`[{"meeting_key":1253,"meeting_name":"Pre-Season Testing"}]`)) + case "/v1/sessions": + _, _ = w.Write([]byte(`[{"session_key":1,"meeting_key":1253,"session_name":"Day 1","session_type":"Testing"}]`)) + default: + http.NotFound(w, r) + } + })) + defer upstream.Close() + client := api.NewOpenF1Client(upstream.URL, 2*time.Second) + defer client.Close() + + races, incomplete, err := fetchSeasonRaces(client.Scoped(), time.Now().Year()) + if err != nil { + t.Fatal(err) + } + if incomplete || len(races) != 0 { + t.Fatalf("recognized testing meeting = races %d, incomplete %v", len(races), incomplete) + } +} + +func TestKnownNonChampionshipMeetingRequiresTestingToken(t *testing.T) { + if !isKnownNonChampionshipMeeting(models.Meeting{MeetingName: "Pre-Season Testing"}, nil) { + t.Fatal("pre-season testing was not recognized") + } + if isKnownNonChampionshipMeeting(models.Meeting{MeetingName: "Fastest Grand Prix"}, nil) { + t.Fatal("substring inside a normal word was treated as testing") + } + if !isKnownNonChampionshipMeeting(models.Meeting{MeetingName: "Winter Event"}, []models.Session{{SessionType: "Test"}}) { + t.Fatal("explicit Test session was not recognized") + } +} + +func TestHandleChampionshipHubSourceLocalWithoutAggregateIsLimited(t *testing.T) { + server := NewServer(nil, 0, nil) + recorder := httptest.NewRecorder() + server.handleChampionshipHub(recorder, httptest.NewRequest(http.MethodGet, "/api/v1/championship/hub?year=2026&source=local", nil)) + if recorder.Code != http.StatusOK || recorder.Header().Get(dataSourceHeader) != "none" || recorder.Header().Get(dataFreshnessHeader) != "limited" { + t.Fatalf("empty local championship = %d %q/%q body=%s", recorder.Code, recorder.Header().Get(dataSourceHeader), recorder.Header().Get(dataFreshnessHeader), recorder.Body.String()) + } +} + func raceResult(num, pos int, pts float64) models.SessionResult { return models.SessionResult{DriverNumber: num, Position: pos, Points: pts} } diff --git a/internal/web/component_freshness_test.go b/internal/web/component_freshness_test.go new file mode 100644 index 0000000..f6818b3 --- /dev/null +++ b/internal/web/component_freshness_test.go @@ -0,0 +1,120 @@ +package web + +import ( + "net/http" + "net/http/httptest" + "testing" + "time" + + "github.com/AmanTahiliani/box-box/internal/api" +) + +func componentTestServer(t *testing.T, responses map[string]string, failures map[string]bool) *Server { + t.Helper() + upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if failures[r.URL.Path] { + http.Error(w, "component unavailable", http.StatusBadGateway) + return + } + body, ok := responses[r.URL.Path] + if !ok { + http.NotFound(w, r) + return + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(body)) + })) + t.Cleanup(upstream.Close) + client := api.NewOpenF1Client(upstream.URL, 2*time.Second) + t.Cleanup(func() { _ = client.Close() }) + return NewServer(client, 0, nil) +} + +func assertAvailabilityHeaders(t *testing.T, recorder *httptest.ResponseRecorder, source, freshness string) { + t.Helper() + if recorder.Code != http.StatusOK || recorder.Header().Get(dataSourceHeader) != source || recorder.Header().Get(dataFreshnessHeader) != freshness { + t.Fatalf("response = status %d, metadata %q/%q, body=%s", recorder.Code, recorder.Header().Get(dataSourceHeader), recorder.Header().Get(dataFreshnessHeader), recorder.Body.String()) + } +} + +func TestResultsAndGridIdentityFailuresReportPartial(t *testing.T) { + tests := []struct { + name string + path string + body string + run func(*Server, http.ResponseWriter, *http.Request) + }{ + {name: "results", path: "/v1/session_result", body: `[{"driver_number":1,"position":1}]`, run: func(s *Server, w http.ResponseWriter, r *http.Request) { s.handleResults(w, r) }}, + {name: "grid", path: "/v1/starting_grid", body: `[{"driver_number":1,"position":1}]`, run: func(s *Server, w http.ResponseWriter, r *http.Request) { s.handleGrid(w, r) }}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + server := componentTestServer(t, map[string]string{tt.path: tt.body}, map[string]bool{"/v1/drivers": true}) + recorder := httptest.NewRecorder() + tt.run(server, recorder, httptest.NewRequest(http.MethodGet, "/api/v1/"+tt.name+"?session_key=99&source=openf1", nil)) + assertAvailabilityHeaders(t, recorder, "openf1", "partial") + }) + } +} + +func TestStrategyOptionalComponentFailureReportsPartial(t *testing.T) { + server := componentTestServer(t, map[string]string{ + "/v1/stints": `[{"driver_number":1,"stint_number":1,"lap_start":1,"lap_end":10,"compound":"MEDIUM"}]`, + "/v1/pit": `[]`, + "/v1/session_result": `[{"driver_number":1,"position":1,"number_of_laps":10}]`, + "/v1/race_control": `[]`, + }, map[string]bool{"/v1/drivers": true}) + recorder := httptest.NewRecorder() + server.handleStrategy(recorder, httptest.NewRequest(http.MethodGet, "/api/v1/strategy?session_key=99", nil)) + assertAvailabilityHeaders(t, recorder, "openf1", "partial") +} + +func TestStrategyEmptyPrimaryDataReportsLimited(t *testing.T) { + server := componentTestServer(t, map[string]string{ + "/v1/stints": `[]`, + "/v1/pit": `[]`, + "/v1/session_result": `[{"driver_number":1,"position":1,"number_of_laps":10}]`, + "/v1/drivers": `[{"driver_number":1,"full_name":"Max Verstappen","team_name":"Red Bull","team_colour":"3671c6"}]`, + "/v1/race_control": `[]`, + }, nil) + recorder := httptest.NewRecorder() + server.handleStrategy(recorder, httptest.NewRequest(http.MethodGet, "/api/v1/strategy?session_key=99", nil)) + assertAvailabilityHeaders(t, recorder, "openf1", "limited") +} + +func TestLapsComparisonDoesNotLabelMissingComponentsFresh(t *testing.T) { + tests := []struct { + name string + laps string + freshness string + }{ + {name: "empty primary data", laps: `[]`, freshness: "limited"}, + {name: "optional components failed", laps: `[{"driver_number":1,"lap_number":1}]`, freshness: "partial"}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + server := componentTestServer(t, map[string]string{"/v1/laps": tt.laps}, map[string]bool{ + "/v1/stints": true, + "/v1/pit": true, + "/v1/race_control": true, + "/v1/drivers": true, + }) + recorder := httptest.NewRecorder() + server.handleLapsComparison(recorder, httptest.NewRequest(http.MethodGet, "/api/v1/laps/comparison?session_key=99", nil)) + assertAvailabilityHeaders(t, recorder, "openf1", tt.freshness) + }) + } + t.Run("primary laps failure is an error", func(t *testing.T) { + server := componentTestServer(t, nil, map[string]bool{ + "/v1/laps": true, + "/v1/stints": true, + "/v1/pit": true, + "/v1/race_control": true, + }) + recorder := httptest.NewRecorder() + server.handleLapsComparison(recorder, httptest.NewRequest(http.MethodGet, "/api/v1/laps/comparison?session_key=99", nil)) + if recorder.Code != http.StatusInternalServerError { + t.Fatalf("status = %d body=%s", recorder.Code, recorder.Body.String()) + } + }) +} diff --git a/internal/web/context.go b/internal/web/context.go index ee63deb..d13aa65 100644 --- a/internal/web/context.go +++ b/internal/web/context.go @@ -10,6 +10,7 @@ import ( func (s *Server) handleWeekendContext(w http.ResponseWriter, _ *http.Request) { if !s.hasLocalQuery() { + markDataResponse(w, "none", "limited") writeJSON(w, query.WeekendContext{TemporalState: query.TemporalNoSeason}) return } @@ -29,9 +30,33 @@ func (s *Server) handleWeekendContext(w http.ResponseWriter, _ *http.Request) { writeError(w, err, http.StatusInternalServerError, false) return } + if focus := focusedContextSession(context); focus != nil { + markDataResponse(w, focus.Availability.Source, focus.Availability.Freshness) + } else { + markDataResponse(w, "none", "limited") + } writeJSON(w, context) } +// focusedContextSession selects the session whose state the Weekend shell is +// presenting. An older terminal/default session must never override an +// upcoming focus session's metadata. +func focusedContextSession(context query.WeekendContext) *query.ContextSession { + if context.ActiveSession != nil { + return context.ActiveSession + } + if context.FocusMeeting == nil { + return nil + } + focusKey := context.FocusMeeting.MeetingKey + for _, ref := range []*query.ContextSession{context.NextSession, context.PreviousCompletedSession, context.DefaultAnalysisSession} { + if ref != nil && ref.Meeting != nil && ref.Meeting.MeetingKey == focusKey { + return ref + } + } + return nil +} + func liveEvidence(data *live.LiveStreamData, active, final bool) query.LiveEvidence { return query.LiveEvidence{Active: active, Final: final, MeetingName: data.Session.MeetingName, CircuitName: data.Session.CircuitName, SessionName: data.Session.SessionName, SessionType: data.Session.SessionType} } diff --git a/internal/web/context_test.go b/internal/web/context_test.go index f5afd0c..18dcd79 100644 --- a/internal/web/context_test.go +++ b/internal/web/context_test.go @@ -47,6 +47,19 @@ func TestWeekendContextHandlerWithoutStoreReturnsNoSeason(t *testing.T) { if got.TemporalState != query.TemporalNoSeason { t.Fatalf("state = %s", got.TemporalState) } + if rr.Header().Get(dataSourceHeader) != "none" || rr.Header().Get(dataFreshnessHeader) != "limited" { + t.Fatalf("missing context metadata = %q/%q", rr.Header().Get(dataSourceHeader), rr.Header().Get(dataFreshnessHeader)) + } +} + +func TestWeekendContextHandlerEmptyStoreReportsLimited(t *testing.T) { + st := openContextStore(t) + s := NewServer(nil, 0, st) + rr := httptest.NewRecorder() + s.handleWeekendContext(rr, httptest.NewRequest(http.MethodGet, "/api/v1/weekend-context", nil)) + if rr.Code != http.StatusOK || rr.Header().Get(dataSourceHeader) != "none" || rr.Header().Get(dataFreshnessHeader) != "limited" { + t.Fatalf("empty context = %d %q/%q body=%s", rr.Code, rr.Header().Get(dataSourceHeader), rr.Header().Get(dataFreshnessHeader), rr.Body.String()) + } } func TestWeekendContextHandlerUsesLiveHubIdentityWithoutOpenF1(t *testing.T) { @@ -76,6 +89,12 @@ func TestWeekendContextHandlerUsesLiveHubIdentityWithoutOpenF1(t *testing.T) { if got.ActiveSession.Availability.LiveTransport != "connected" || got.ActiveSession.Availability.LiveSession != "active" { t.Fatalf("availability = %+v", got.ActiveSession.Availability) } + if got.ActiveSession.Availability.Source != "mixed" || got.ActiveSession.Availability.Freshness != "live" { + t.Fatalf("live source/freshness = %+v", got.ActiveSession.Availability) + } + if rr.Header().Get(dataSourceHeader) != "mixed" || rr.Header().Get(dataFreshnessHeader) != "live" { + t.Fatalf("response source/freshness = %q/%q", rr.Header().Get(dataSourceHeader), rr.Header().Get(dataFreshnessHeader)) + } } func TestWeekendContextHandlerUsesTerminalArchiveAsCompletionEvidence(t *testing.T) { @@ -94,11 +113,48 @@ func TestWeekendContextHandlerUsesTerminalArchiveAsCompletionEvidence(t *testing if got.PreviousCompletedSession == nil || got.PreviousCompletedSession.Availability.Archive != "available" { t.Fatalf("archive context = %+v", got) } + if got.PreviousCompletedSession.Availability.Source != "mixed" || got.PreviousCompletedSession.Availability.Freshness != "archive" { + t.Fatalf("archive source/freshness = %+v", got.PreviousCompletedSession.Availability) + } + if rr.Header().Get(dataSourceHeader) != "mixed" || rr.Header().Get(dataFreshnessHeader) != "archive" { + t.Fatalf("archive response source/freshness = %q/%q", rr.Header().Get(dataSourceHeader), rr.Header().Get(dataFreshnessHeader)) + } if got.DefaultAnalysisSession != nil { t.Fatal("archive-only session must not become local default analysis") } } +func TestWeekendContextMetadataFollowsUpcomingFocusNotTerminalPrevious(t *testing.T) { + st := openContextStore(t) + seedContextHandler(t, st) + if err := st.UpsertMeeting(store.Meeting{MeetingKey: 2, MeetingName: "Belgian Grand Prix", CircuitShortName: "Spa", Year: 2026, DateStart: "2026-07-10T09:00:00Z", DateEnd: "2026-07-12T16:00:00Z"}); err != nil { + t.Fatal(err) + } + if err := st.UpsertSession(store.Session{SessionKey: 21, MeetingKey: 2, SessionName: "Practice 1", SessionType: "Practice", DateStart: "2026-07-10T09:00:00Z", DateEnd: "2026-07-10T10:00:00Z"}); err != nil { + t.Fatal(err) + } + now, _ := time.Parse(time.RFC3339, "2026-07-05T16:05:00Z") + s := NewServer(nil, 0, st) + s.query = query.NewServiceWithClock(st, func() time.Time { return now }) + s.hub.applySnapshot(live.LiveStreamData{SessionStatus: "Finished", Session: live.LiveSessionMeta{MeetingName: "British Grand Prix", CircuitName: "Silverstone", SessionName: "Race", SessionType: "Race"}}, now) + + rr := httptest.NewRecorder() + s.handleWeekendContext(rr, httptest.NewRequest(http.MethodGet, "/api/v1/weekend-context", nil)) + var got query.WeekendContext + if err := json.Unmarshal(rr.Body.Bytes(), &got); err != nil { + t.Fatal(err) + } + if got.FocusMeeting == nil || got.FocusMeeting.MeetingKey != 2 || got.NextSession == nil { + t.Fatalf("focus context = %+v", got) + } + if got.PreviousCompletedSession == nil || got.PreviousCompletedSession.Availability.Freshness != "archive" { + t.Fatalf("terminal previous missing = %+v", got.PreviousCompletedSession) + } + if rr.Header().Get(dataSourceHeader) != "local" || rr.Header().Get(dataFreshnessHeader) != "local" { + t.Fatalf("focus metadata was overridden by archive = %q/%q", rr.Header().Get(dataSourceHeader), rr.Header().Get(dataFreshnessHeader)) + } +} + func TestTerminalSessionStatus(t *testing.T) { for _, status := range []string{"Finished", "Finalised", "ENDED", "Aborted"} { if !terminalSessionStatus(status) { diff --git a/internal/web/driversummary.go b/internal/web/driversummary.go index 7c87e59..1b25c67 100644 --- a/internal/web/driversummary.go +++ b/internal/web/driversummary.go @@ -1,18 +1,25 @@ package web import ( + "context" "fmt" "net/http" "strconv" "time" + "github.com/AmanTahiliani/box-box/internal/api" "github.com/AmanTahiliani/box-box/internal/models" ) // --- /api/v1/driver/summary --- -// Per-driver season summary: championship standing plus per-round race results, -// aggregated server-side from the same sources as the championship hub. Caching -// relies on the OpenF1 client's HTTP cache TTLs — no extra layer here. +// Per-driver season summary: championship standing plus per-round race results. +// Current-season identity/results are local-first from the domain DB. Optional +// OpenF1 enrichment (headshot / polished identity) is bounded so it cannot hang +// the profile when remote data is slow or unavailable. + +// driverEnrichmentTimeout bounds optional remote enrichment so a hung OpenF1 +// call never blocks a local-first profile response. Overridable in tests. +var driverEnrichmentTimeout = 2 * time.Second type driverSummaryRound struct { MeetingKey int `json:"meeting_key"` @@ -49,6 +56,11 @@ type driverSummaryResponse struct { Cumulative []float64 `json:"cumulative"` RoundLabels []string `json:"round_labels"` Rounds []driverSummaryRound `json:"rounds"` + // Source is "local" when served from the domain DB, else "openf1". + Source string `json:"source,omitempty"` + // Enrichment is "full" when optional remote identity landed, "limited" + // when it timed out/failed, or "none" when no enrichment was attempted. + Enrichment string `json:"enrichment,omitempty"` } func (s *Server) handleDriverSummary(w http.ResponseWriter, r *http.Request) { @@ -61,12 +73,102 @@ func (s *Server) handleDriverSummary(w http.ResponseWriter, r *http.Request) { if year == 0 { year = time.Now().Year() } + mode := parseSourceMode(r) + // Driver summary is local-first for current-season identity/results. When the + // caller omits ?source=, prefer auto (local then OpenF1) rather than the + // package default of openf1-only. + if r.URL.Query().Get("source") == "" { + mode = sourceAuto + } + client := s.client.Scoped() - champ, err := s.client.GetDriverChampionshipForYear(year) + if mode == sourceLocal || mode == sourceAuto { + resp, sessionKey, ok, lerr := s.localDriverSummary(year, driverNumber) + if lerr != nil { + writeError(w, lerr, http.StatusInternalServerError, false) + return + } + if ok { + if mode != sourceLocal { + tryEnrichDriverSummary(r.Context(), client, &resp, sessionKey) + } + switch resp.Enrichment { + case "full": + markMixedResponse(w, client, false) + case "limited": + markDataResponse(w, "local", "limited") + default: + markLocalResponse(w, false) + } + writeJSON(w, resp) + return + } + if mode == sourceLocal { + http.Error(w, fmt.Sprintf("driver %d not found in %d championship", driverNumber, year), http.StatusNotFound) + return + } + } + + resp, incomplete, err := openF1DriverSummary(client, year, driverNumber) if err != nil { - writeError(w, err, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return } + if resp == nil { + http.Error(w, fmt.Sprintf("driver %d not found in %d championship", driverNumber, year), http.StatusNotFound) + return + } + markOpenF1AggregateResponse(w, client, incomplete) + writeJSON(w, resp) +} + +func (s *Server) localDriverSummary(year, driverNumber int) (driverSummaryResponse, int, bool, error) { + if !s.hasLocalQuery() { + return driverSummaryResponse{}, 0, false, nil + } + inputs, err := s.query.GetChampionshipInputs(year) + if err != nil { + return driverSummaryResponse{}, 0, false, err + } + if len(inputs.Champ) == 0 { + return driverSummaryResponse{}, 0, false, nil + } + + races := make([]meetingRace, 0, len(inputs.Races)) + for _, race := range inputs.Races { + races = append(races, meetingRace{ + Meeting: race.Meeting, + RaceSessionKey: race.RaceSessionKey, + Results: race.Results, + Grid: race.Grid, + }) + } + + resp, ok := aggregateDriverSummary(year, driverNumber, races, inputs.Champ, inputs.DriverMap) + if !ok { + return driverSummaryResponse{}, 0, false, nil + } + resp.Source = "local" + resp.Enrichment = "none" + + sessionKey := 0 + for _, c := range inputs.Champ { + if c.DriverNumber == driverNumber && c.SessionKey > 0 { + sessionKey = c.SessionKey + break + } + if sessionKey == 0 && c.SessionKey > 0 { + sessionKey = c.SessionKey + } + } + return resp, sessionKey, true, nil +} + +func openF1DriverSummary(client *api.OpenF1Client, year, driverNumber int) (*driverSummaryResponse, bool, error) { + champ, err := client.GetDriverChampionshipForYear(year) + if err != nil { + return nil, false, err + } var entry *models.ChampionshipDriver for i := range champ { if champ[i].DriverNumber == driverNumber { @@ -75,30 +177,83 @@ func (s *Server) handleDriverSummary(w http.ResponseWriter, r *http.Request) { } } if entry == nil { - http.Error(w, fmt.Sprintf("driver %d not found in %d championship", driverNumber, year), http.StatusNotFound) - return + return nil, false, nil } driverInfo := map[int]models.Driver{} - if ds, derr := s.client.GetDriversForSession(champ[0].SessionKey); derr == nil { + sessionKey := champ[0].SessionKey + if ds, derr := client.GetDriversForSession(sessionKey); derr == nil { driverInfo = buildDriverMapFirst(ds) } - if d, ok := s.championshipDriverInfo(entry.SessionKey, driverNumber, driverInfo); ok { - driverInfo[driverNumber] = d + d, directErr := client.GetDriver(entry.SessionKey, driverNumber) + if directErr == nil && d != nil { + driverInfo[driverNumber] = *d + } else if fallback, ok := driverInfo[driverNumber]; ok { + driverInfo[driverNumber] = fallback } - - races, _, err := s.fetchSeasonRaces(year) + identityIncomplete := !hasDriverPresentation(driverInfo[driverNumber]) + races, racesIncomplete, err := fetchSeasonRaces(client, year) if err != nil { - writeError(w, err, http.StatusInternalServerError, s.client.LastResponseWasStale()) - return + return nil, false, err } resp, ok := aggregateDriverSummary(year, driverNumber, races, champ, driverInfo) if !ok { - http.Error(w, fmt.Sprintf("driver %d not found in %d championship", driverNumber, year), http.StatusNotFound) + return nil, false, nil + } + incomplete := identityIncomplete || racesIncomplete + resp.Source = "openf1" + if identityIncomplete { + resp.Enrichment = "limited" + } else { + resp.Enrichment = "full" + } + return &resp, incomplete, nil +} + +func hasDriverPresentation(driver models.Driver) bool { + hasName := driver.FullName != "" || driver.NameAcronym != "" || driver.BroadcastName != "" + return hasName && driver.TeamName != "" && driver.TeamColour != "" +} + +// tryEnrichDriverSummary optionally fills headshot / polished identity from +// OpenF1. It never blocks longer than driverEnrichmentTimeout — on timeout or +// failure the local profile remains intact with enrichment=limited. +func tryEnrichDriverSummary(parent context.Context, client *api.OpenF1Client, resp *driverSummaryResponse, sessionKey int) { + if resp == nil || client == nil || sessionKey <= 0 { + if resp != nil && resp.Enrichment == "none" { + // No session to enrich from — leave as none (local identity only). + } return } - writeJSON(w, resp) + + ctx, cancel := context.WithTimeout(parent, driverEnrichmentTimeout) + defer cancel() + driver, err := client.GetDriverContext(ctx, sessionKey, resp.DriverNumber) + if err != nil || driver == nil { + resp.Enrichment = "limited" + return + } + applyDriverEnrichment(resp, *driver) + resp.Enrichment = "full" +} + +func applyDriverEnrichment(resp *driverSummaryResponse, d models.Driver) { + if d.HeadshotURL != "" { + resp.HeadshotURL = d.HeadshotURL + } + if d.FullName != "" { + resp.FullName = d.FullName + } + if d.NameAcronym != "" { + resp.NameAcronym = d.NameAcronym + } + if d.TeamName != "" { + resp.TeamName = d.TeamName + } + if d.TeamColour != "" { + resp.TeamColour = d.TeamColour + } } // aggregateDriverSummary is the pure aggregation core (no network) so it can be diff --git a/internal/web/driversummary_test.go b/internal/web/driversummary_test.go index 16e1dec..97a7037 100644 --- a/internal/web/driversummary_test.go +++ b/internal/web/driversummary_test.go @@ -1,11 +1,16 @@ package web import ( + "encoding/json" + "fmt" "net/http" "net/http/httptest" "testing" + "time" + "github.com/AmanTahiliani/box-box/internal/api" "github.com/AmanTahiliani/box-box/internal/models" + "github.com/AmanTahiliani/box-box/internal/store" ) func driverSummaryFixtures() ([]meetingRace, []models.ChampionshipDriver, map[int]models.Driver) { @@ -134,3 +139,234 @@ func TestHandleDriverSummaryBadRequest(t *testing.T) { } } } + +func seedDriverSummaryStore(t *testing.T, st *store.Store) { + t.Helper() + meetingKey := 1201 + sessionKey := 9901 + + if err := st.UpsertMeeting(store.Meeting{ + MeetingKey: meetingKey, + MeetingName: "Bahrain GP", + CountryCode: "BHR", + CountryName: "Bahrain", + Year: 2025, + DateStart: "2025-03-02", + }); err != nil { + t.Fatalf("UpsertMeeting: %v", err) + } + if err := st.UpsertSession(store.Session{ + SessionKey: sessionKey, + MeetingKey: meetingKey, + SessionName: "Race", + SessionType: "Race", + DateStart: "2025-03-02T15:00:00Z", + }); err != nil { + t.Fatalf("UpsertSession: %v", err) + } + if err := st.UpsertDriver(store.Driver{ + DriverNumber: 1, + FullName: "Max Verstappen", + NameAcronym: "VER", + TeamName: "Red Bull", + TeamColour: "3671c6", + }); err != nil { + t.Fatalf("UpsertDriver: %v", err) + } + if err := st.UpsertSessionDriver(store.SessionDriver{ + SessionKey: sessionKey, + DriverNumber: 1, + MeetingKey: meetingKey, + FullName: "Max Verstappen", + NameAcronym: "VER", + TeamName: "Red Bull", + TeamColour: "3671c6", + }); err != nil { + t.Fatalf("UpsertSessionDriver: %v", err) + } + if err := st.UpsertSessionResult(store.SessionResult{ + SessionKey: sessionKey, + DriverNumber: 1, + MeetingKey: meetingKey, + Position: 1, + Points: 25, + }); err != nil { + t.Fatalf("UpsertSessionResult: %v", err) + } + if err := st.UpsertStartingGridEntry(store.StartingGridEntry{ + SessionKey: sessionKey, + DriverNumber: 1, + MeetingKey: meetingKey, + Position: 1, + }); err != nil { + t.Fatalf("UpsertStartingGridEntry: %v", err) + } +} + +func TestHandleDriverSummaryLocalFirstIgnoresHangingEnrichment(t *testing.T) { + prev := driverEnrichmentTimeout + driverEnrichmentTimeout = 40 * time.Millisecond + t.Cleanup(func() { driverEnrichmentTimeout = prev }) + + st := openTestStore(t) + seedDriverSummaryStore(t, st) + + // Enrichment seam: OpenF1 hangs until released. Local summary must still return. + release := make(chan struct{}) + cancelObserved := make(chan struct{}) + hang := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + select { + case <-release: + case <-r.Context().Done(): + close(cancelObserved) + } + })) + t.Cleanup(func() { + close(release) + hang.Close() + }) + + client := api.NewOpenF1Client(hang.URL, 15*time.Second) + t.Cleanup(func() { _ = client.Close() }) + srv := NewServer(client, 8080, st) + + start := time.Now() + req := httptest.NewRequest(http.MethodGet, "/api/v1/driver/summary?year=2025&driver_number=1", nil) + rec := httptest.NewRecorder() + srv.handleDriverSummary(rec, req) + elapsed := time.Since(start) + + if rec.Code != http.StatusOK { + t.Fatalf("status = %d body=%s", rec.Code, rec.Body.String()) + } + if elapsed > 500*time.Millisecond { + t.Fatalf("handler blocked on enrichment for %v", elapsed) + } + select { + case <-cancelObserved: + case <-time.After(250 * time.Millisecond): + t.Fatal("timed-out enrichment did not cancel its upstream request") + } + + var resp driverSummaryResponse + if err := json.NewDecoder(rec.Body).Decode(&resp); err != nil { + t.Fatalf("decode: %v", err) + } + if resp.Source != "local" { + t.Errorf("source = %q, want local", resp.Source) + } + if resp.Enrichment != "limited" { + t.Errorf("enrichment = %q, want limited", resp.Enrichment) + } + if rec.Header().Get(dataSourceHeader) != "local" || rec.Header().Get(dataFreshnessHeader) != "limited" { + t.Errorf("limited metadata = %q/%q", rec.Header().Get(dataSourceHeader), rec.Header().Get(dataFreshnessHeader)) + } + if resp.DriverNumber != 1 || resp.NameAcronym != "VER" || resp.Points != 25 { + t.Errorf("local identity/results missing: %+v", resp) + } +} + +func TestHandleDriverSummaryLocalFirstWithFailingEnrichment(t *testing.T) { + prev := driverEnrichmentTimeout + driverEnrichmentTimeout = 200 * time.Millisecond + t.Cleanup(func() { driverEnrichmentTimeout = prev }) + + st := openTestStore(t) + seedDriverSummaryStore(t, st) + + fail := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + http.Error(w, "boom", http.StatusBadGateway) + })) + t.Cleanup(fail.Close) + + client := api.NewOpenF1Client(fail.URL, 2*time.Second) + t.Cleanup(func() { _ = client.Close() }) + srv := NewServer(client, 8080, st) + + req := httptest.NewRequest(http.MethodGet, "/api/v1/driver/summary?year=2025&driver_number=1&source=auto", nil) + rec := httptest.NewRecorder() + srv.handleDriverSummary(rec, req) + + if rec.Code != http.StatusOK { + t.Fatalf("status = %d body=%s", rec.Code, rec.Body.String()) + } + var resp driverSummaryResponse + if err := json.NewDecoder(rec.Body).Decode(&resp); err != nil { + t.Fatalf("decode: %v", err) + } + if resp.Source != "local" { + t.Errorf("source = %q, want local", resp.Source) + } + if resp.Enrichment != "limited" { + t.Errorf("enrichment = %q, want limited", resp.Enrichment) + } + if rec.Header().Get(dataSourceHeader) != "local" || rec.Header().Get(dataFreshnessHeader) != "limited" { + t.Errorf("limited metadata = %q/%q", rec.Header().Get(dataSourceHeader), rec.Header().Get(dataFreshnessHeader)) + } + if resp.FullName != "Max Verstappen" { + t.Errorf("full_name = %q, want local identity", resp.FullName) + } +} + +func TestHandleDriverSummarySourceLocalOnly(t *testing.T) { + st := openTestStore(t) + seedDriverSummaryStore(t, st) + + // Even with a broken OpenF1 client, source=local must succeed from the DB. + fail := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + http.Error(w, "nope", http.StatusInternalServerError) + })) + t.Cleanup(fail.Close) + client := api.NewOpenF1Client(fail.URL, time.Second) + t.Cleanup(func() { _ = client.Close() }) + srv := NewServer(client, 8080, st) + + req := httptest.NewRequest(http.MethodGet, "/api/v1/driver/summary?year=2025&driver_number=1&source=local", nil) + rec := httptest.NewRecorder() + srv.handleDriverSummary(rec, req) + if rec.Code != http.StatusOK { + t.Fatalf("status = %d body=%s", rec.Code, rec.Body.String()) + } + if rec.Header().Get(dataSourceHeader) != "local" || rec.Header().Get(dataFreshnessHeader) != "local" { + t.Fatalf("local metadata = %q/%q", rec.Header().Get(dataSourceHeader), rec.Header().Get(dataFreshnessHeader)) + } +} + +func TestHandleRemoteDriverSummaryReportsIdentityAndRoundLimitations(t *testing.T) { + tests := []struct { + name string + driversOK bool + meetingHasRace bool + wantEnrichment string + }{ + {name: "missing identity", driversOK: false, meetingHasRace: true, wantEnrichment: "limited"}, + {name: "missing race round", driversOK: true, meetingHasRace: false, wantEnrichment: "full"}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + upstream := championshipTestUpstream(t, tt.driversOK, tt.meetingHasRace) + defer upstream.Close() + client := api.NewOpenF1Client(upstream.URL, 2*time.Second) + defer client.Close() + server := NewServer(client, 0, nil) + year := time.Now().Year() + recorder := httptest.NewRecorder() + request := httptest.NewRequest(http.MethodGet, fmt.Sprintf("/api/v1/driver/summary?year=%d&driver_number=1&source=openf1", year), nil) + + server.handleDriverSummary(recorder, request) + if recorder.Code != http.StatusOK { + t.Fatalf("status = %d body=%s", recorder.Code, recorder.Body.String()) + } + var response driverSummaryResponse + if err := json.Unmarshal(recorder.Body.Bytes(), &response); err != nil { + t.Fatal(err) + } + if response.Enrichment != tt.wantEnrichment { + t.Fatalf("enrichment = %q, want %q", response.Enrichment, tt.wantEnrichment) + } + if recorder.Header().Get(dataSourceHeader) != "openf1" || recorder.Header().Get(dataFreshnessHeader) != "partial" { + t.Fatalf("remote limitation metadata = %q/%q", recorder.Header().Get(dataSourceHeader), recorder.Header().Get(dataFreshnessHeader)) + } + }) + } +} diff --git a/internal/web/freshness.go b/internal/web/freshness.go new file mode 100644 index 0000000..f6d9530 --- /dev/null +++ b/internal/web/freshness.go @@ -0,0 +1,65 @@ +package web + +import ( + "net/http" +) + +const ( + dataSourceHeader = "X-BoxBox-Data-Source" + dataFreshnessHeader = "X-BoxBox-Data-Freshness" +) + +type staleResponseReporter interface { + LastResponseWasStale() bool +} + +// markOpenF1Response publishes request-scoped success provenance. Callers must +// pass the scoped client used for this response, never Server.client. +func markOpenF1Response(w http.ResponseWriter, client staleResponseReporter) { + markOpenF1AggregateResponse(w, client, false) +} + +func markOpenF1AggregateResponse(w http.ResponseWriter, client staleResponseReporter, partial bool) { + freshness := "fresh" + if partial { + freshness = "partial" + } + markOpenF1Availability(w, client, freshness) +} + +func markOpenF1Availability(w http.ResponseWriter, client staleResponseReporter, freshness string) { + w.Header().Set(dataSourceHeader, "openf1") + if client != nil && client.LastResponseWasStale() { + w.Header().Set(dataFreshnessHeader, "stale") + return + } + if freshness == "" { + freshness = "fresh" + } + w.Header().Set(dataFreshnessHeader, freshness) +} + +func markDataResponse(w http.ResponseWriter, source, freshness string) { + w.Header().Set(dataSourceHeader, source) + w.Header().Set(dataFreshnessHeader, freshness) +} + +func markLocalResponse(w http.ResponseWriter, partial bool) { + w.Header().Set(dataSourceHeader, "local") + if partial { + w.Header().Set(dataFreshnessHeader, "partial") + return + } + w.Header().Set(dataFreshnessHeader, "local") +} + +func markMixedResponse(w http.ResponseWriter, client staleResponseReporter, partial bool) { + w.Header().Set(dataSourceHeader, "mixed") + if client != nil && client.LastResponseWasStale() { + w.Header().Set(dataFreshnessHeader, "stale") + } else if partial { + w.Header().Set(dataFreshnessHeader, "partial") + } else { + w.Header().Set(dataFreshnessHeader, "local") + } +} diff --git a/internal/web/freshness_test.go b/internal/web/freshness_test.go new file mode 100644 index 0000000..bb63870 --- /dev/null +++ b/internal/web/freshness_test.go @@ -0,0 +1,86 @@ +package web + +import ( + "database/sql" + "fmt" + "net/http" + "net/http/httptest" + "testing" + "time" + + "github.com/AmanTahiliani/box-box/internal/api" + _ "modernc.org/sqlite" +) + +type fakeStaleReporter bool + +func (f fakeStaleReporter) LastResponseWasStale() bool { return bool(f) } + +func TestStaleFreshnessTakesPrecedenceOverPartialAndLimited(t *testing.T) { + for _, fallback := range []string{"partial", "limited"} { + recorder := httptest.NewRecorder() + markOpenF1Availability(recorder, fakeStaleReporter(true), fallback) + if recorder.Header().Get(dataFreshnessHeader) != "stale" { + t.Fatalf("fallback %q overrode stale: %q", fallback, recorder.Header().Get(dataFreshnessHeader)) + } + } +} + +func TestOpenF1HandlerReportsFreshThenStaleSuccess(t *testing.T) { + year := time.Now().Year() + upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + _, _ = fmt.Fprintf(w, `[{"meeting_key":1,"meeting_name":"British Grand Prix","year":%d}]`, year) + })) + client := api.NewOpenF1Client(upstream.URL, time.Second) + t.Cleanup(func() { _ = client.Close() }) + server := NewServer(client, 0, nil) + cacheKey := fmt.Sprintf("%s/v1/meetings?year=%d", upstream.URL, year) + requestURL := fmt.Sprintf("/api/v1/meetings?year=%d&source=openf1", year) + defer func() { + // Keep the shared application cache clean even if this test fails. + db, err := sql.Open("sqlite", api.DefaultCacheDBPath()+"?_busy_timeout=5000") + if err == nil { + _, _ = db.Exec(`DELETE FROM cache WHERE key = ?`, cacheKey) + _ = db.Close() + } + }() + + fresh := httptest.NewRecorder() + server.handleMeetings(fresh, httptest.NewRequest(http.MethodGet, requestURL, nil)) + if fresh.Code != http.StatusOK || fresh.Header().Get(dataSourceHeader) != "openf1" || fresh.Header().Get(dataFreshnessHeader) != "fresh" { + t.Fatalf("fresh response status/metadata = %d %q/%q body=%s", fresh.Code, fresh.Header().Get(dataSourceHeader), fresh.Header().Get(dataFreshnessHeader), fresh.Body.String()) + } + + db, err := sql.Open("sqlite", api.DefaultCacheDBPath()+"?_busy_timeout=5000") + if err != nil { + t.Fatal(err) + } + if _, err := db.Exec(`UPDATE cache SET created_at = ? WHERE key = ?`, time.Now().Add(-48*time.Hour).Unix(), cacheKey); err != nil { + _ = db.Close() + t.Fatal(err) + } + _ = db.Close() + upstream.Close() + + stale := httptest.NewRecorder() + server.handleMeetings(stale, httptest.NewRequest(http.MethodGet, requestURL, nil)) + if stale.Code != http.StatusOK || stale.Header().Get(dataSourceHeader) != "openf1" || stale.Header().Get(dataFreshnessHeader) != "stale" { + t.Fatalf("stale response status/metadata = %d %q/%q body=%s", stale.Code, stale.Header().Get(dataSourceHeader), stale.Header().Get(dataFreshnessHeader), stale.Body.String()) + } +} + +func TestFreshnessHeadersAreExposedToBrowserClients(t *testing.T) { + handler := withCORS(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + markDataResponse(w, "local", "partial") + writeJSON(w, map[string]bool{"ok": true}) + })) + recorder := httptest.NewRecorder() + handler.ServeHTTP(recorder, httptest.NewRequest(http.MethodGet, "/api/v1/test", nil)) + + if got := recorder.Header().Get("Access-Control-Expose-Headers"); got != dataSourceHeader+", "+dataFreshnessHeader { + t.Fatalf("exposed headers = %q", got) + } + if recorder.Header().Get(dataSourceHeader) != "local" || recorder.Header().Get(dataFreshnessHeader) != "partial" { + t.Fatalf("data metadata = %q/%q", recorder.Header().Get(dataSourceHeader), recorder.Header().Get(dataFreshnessHeader)) + } +} diff --git a/internal/web/live.go b/internal/web/live.go index a390010..7f56f42 100644 --- a/internal/web/live.go +++ b/internal/web/live.go @@ -364,7 +364,15 @@ func cloneLivePositions(in map[string]live.LivePositionData) map[string]live.Liv // handleLiveState returns the current live data snapshot as JSON. func (s *Server) handleLiveState(w http.ResponseWriter, r *http.Request) { - writeJSON(w, s.hub.State()) + state := s.hub.State() + freshness := "limited" + if state.IsLive { + freshness = "live" + } else if state.LastSnapshot != nil { + freshness = "archive" + } + markDataResponse(w, "fia", freshness) + writeJSON(w, state) } // handleSSEStream is the persistent SSE endpoint for live data. diff --git a/internal/web/navigation.go b/internal/web/navigation.go index d6b9830..d9ffd8c 100644 --- a/internal/web/navigation.go +++ b/internal/web/navigation.go @@ -9,6 +9,7 @@ import ( ) func (s *Server) handleSeasons(w http.ResponseWriter, r *http.Request) { + markLocalResponse(w, false) if !s.hasLocalQuery() { writeJSON(w, []int{}) return @@ -26,6 +27,7 @@ func (s *Server) handleSeasons(w http.ResponseWriter, r *http.Request) { } func (s *Server) handleWeekend(w http.ResponseWriter, r *http.Request) { + markLocalResponse(w, false) meetingKey, err := strconv.Atoi(r.URL.Query().Get("meeting_key")) if err != nil || meetingKey == 0 { http.Error(w, "meeting_key required", http.StatusBadRequest) @@ -46,5 +48,6 @@ func (s *Server) handleWeekend(w http.ResponseWriter, r *http.Request) { writeError(w, err, http.StatusInternalServerError, false) return } + markLocalResponse(w, weekend.Source == query.ResponseSourcePartial) writeJSON(w, weekend) } diff --git a/internal/web/racehub.go b/internal/web/racehub.go index db8725e..b00647a 100644 --- a/internal/web/racehub.go +++ b/internal/web/racehub.go @@ -17,6 +17,7 @@ func (s *Server) handleRaceHub(w http.ResponseWriter, r *http.Request) { } if !s.hasLocalQuery() { + markDataResponse(w, "none", "limited") writeJSON(w, emptyRaceHub(sessionKey)) return } @@ -26,6 +27,14 @@ func (s *Server) handleRaceHub(w http.ResponseWriter, r *http.Request) { writeError(w, err, http.StatusInternalServerError, false) return } + switch hub.Source { + case query.ResponseSourceNone: + markDataResponse(w, "none", "limited") + case query.ResponseSourcePartial: + markLocalResponse(w, true) + default: + markLocalResponse(w, false) + } writeJSON(w, hub) } diff --git a/internal/web/racehub_test.go b/internal/web/racehub_test.go index 50755d5..58babde 100644 --- a/internal/web/racehub_test.go +++ b/internal/web/racehub_test.go @@ -86,11 +86,29 @@ func TestHandleRaceHubWithoutStore(t *testing.T) { if hub.Source != query.ResponseSourceNone { t.Fatalf("source = %q, want %q", hub.Source, query.ResponseSourceNone) } + if rec.Header().Get(dataSourceHeader) != "none" || rec.Header().Get(dataFreshnessHeader) != "limited" { + t.Fatalf("missing hub metadata = %q/%q", rec.Header().Get(dataSourceHeader), rec.Header().Get(dataFreshnessHeader)) + } if hub.Datasets["session"].Status != query.DatasetStatusMissing { t.Fatalf("session dataset = %+v, want missing", hub.Datasets["session"]) } } +func TestHandleRaceHubUnknownSessionWithStoreReportsLimited(t *testing.T) { + st := openTestStore(t) + srv := testServer(t, st) + rec := httptest.NewRecorder() + srv.handleRaceHub(rec, httptest.NewRequest(http.MethodGet, "/api/v1/race-hub?session_key=999999", nil)) + + var hub query.RaceHub + if err := json.Unmarshal(rec.Body.Bytes(), &hub); err != nil { + t.Fatal(err) + } + if hub.Source != query.ResponseSourceNone || rec.Header().Get(dataSourceHeader) != "none" || rec.Header().Get(dataFreshnessHeader) != "limited" { + t.Fatalf("unknown session = source %q, metadata %q/%q", hub.Source, rec.Header().Get(dataSourceHeader), rec.Header().Get(dataFreshnessHeader)) + } +} + func TestHandleRaceHubWithLocalData(t *testing.T) { st := openTestStore(t) seedRaceHubStore(t, st) @@ -117,6 +135,12 @@ func TestHandleRaceHubWithLocalData(t *testing.T) { if hub.Datasets["results"].Status != query.DatasetStatusMissing { t.Fatalf("results dataset = %+v, want missing", hub.Datasets["results"]) } + if hub.Source != query.ResponseSourcePartial { + t.Fatalf("source = %q, want partial", hub.Source) + } + if rec.Header().Get(dataSourceHeader) != "local" || rec.Header().Get(dataFreshnessHeader) != "partial" { + t.Fatalf("partial hub metadata = %q/%q", rec.Header().Get(dataSourceHeader), rec.Header().Get(dataFreshnessHeader)) + } } func TestHandleRaceHubIncludesChapters(t *testing.T) { diff --git a/internal/web/replay.go b/internal/web/replay.go index 9caa7bc..7277013 100644 --- a/internal/web/replay.go +++ b/internal/web/replay.go @@ -61,15 +61,27 @@ func (s *Server) handleReplayFrames(w http.ResponseWriter, r *http.Request) { } } - resp, err := assembleReplayFrames(r.Context(), s.client, sessionKey, intervalMS) + client := s.client.Scoped() + resp, incomplete, err := assembleReplayFrames(r.Context(), client, sessionKey, intervalMS) if err != nil { - writeError(w, err, http.StatusInternalServerError, s.client.LastResponseWasStale()) + writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return } + markOpenF1Availability(w, client, replayResponseFreshness(resp, incomplete)) writeJSON(w, resp) } -func assembleReplayFrames(ctx context.Context, client replayDataClient, sessionKey, intervalMS int) (replayFramesResponse, error) { +func replayResponseFreshness(resp replayFramesResponse, incomplete bool) string { + if !incomplete { + return "fresh" + } + if len(resp.Frames) == 0 { + return "limited" + } + return "partial" +} + +func assembleReplayFrames(ctx context.Context, client replayDataClient, sessionKey, intervalMS int) (replayFramesResponse, bool, error) { if intervalMS < defaultReplayIntervalMS { intervalMS = defaultReplayIntervalMS } @@ -82,26 +94,26 @@ func assembleReplayFrames(ctx context.Context, client replayDataClient, sessionK drivers, err := client.GetDriversForSession(sessionKey) if err != nil { - return resp, err + return resp, false, err } driverNumbers := uniqueDriverNumbers(drivers) if len(driverNumbers) == 0 { - return resp, nil + return resp, true, nil } series, err := fetchReplayLocationSeries(ctx, client, sessionKey, driverNumbers) if err != nil && len(series) == 0 { - return resp, err + return resp, false, err } start, ok := earliestReplayLocationTime(series) if !ok { - return resp, nil + return resp, true, nil } resp.StartTime = start.Format(time.RFC3339Nano) resp.Frames = snapReplayFrames(series, start, intervalMS) - return resp, nil + return resp, err != nil || len(series) < len(driverNumbers) || len(resp.Frames) == 0, nil } func uniqueDriverNumbers(drivers []models.Driver) []int { diff --git a/internal/web/replay_test.go b/internal/web/replay_test.go index 25e7437..db24b86 100644 --- a/internal/web/replay_test.go +++ b/internal/web/replay_test.go @@ -3,6 +3,7 @@ package web import ( "context" "encoding/json" + "errors" "net/http" "net/http/httptest" "sync" @@ -16,6 +17,7 @@ type fakeReplayClient struct { drivers []models.Driver locs map[int][]models.Location err error + locErrs map[int]error mu sync.Mutex inFlight int @@ -46,7 +48,38 @@ func (f *fakeReplayClient) GetLocation(sessionKey, driverNumber int) ([]models.L f.inFlight-- f.mu.Unlock() - return f.locs[driverNumber], nil + return f.locs[driverNumber], f.locErrs[driverNumber] +} + +func TestAssembleReplayFramesReportsPartialDriverSeries(t *testing.T) { + start := time.Date(2025, 5, 25, 13, 0, 0, 0, time.UTC) + client := &fakeReplayClient{ + drivers: []models.Driver{{DriverNumber: 1}, {DriverNumber: 4}}, + locs: map[int][]models.Location{ + 1: {{Date: start.Format(time.RFC3339Nano), X: 1, Y: 2}}, + }, + locErrs: map[int]error{4: errors.New("location unavailable")}, + } + resp, incomplete, err := assembleReplayFrames(context.Background(), client, 99, defaultReplayIntervalMS) + if err != nil { + t.Fatalf("partial replay should remain usable: %v", err) + } + if !incomplete || len(resp.Frames) != 1 { + t.Fatalf("partial replay = incomplete %v, frames %+v", incomplete, resp.Frames) + } + if got := replayResponseFreshness(resp, incomplete); got != "partial" { + t.Fatalf("partial replay freshness = %q", got) + } +} + +func TestAssembleReplayFramesEmptyDriverSetIsLimited(t *testing.T) { + resp, incomplete, err := assembleReplayFrames(context.Background(), &fakeReplayClient{}, 99, defaultReplayIntervalMS) + if err != nil { + t.Fatal(err) + } + if !incomplete || replayResponseFreshness(resp, incomplete) != "limited" { + t.Fatalf("empty replay = incomplete %v, freshness %q", incomplete, replayResponseFreshness(resp, incomplete)) + } } func TestAssembleReplayFramesSnapsNearestSamplesAndOmitsEmptyDrivers(t *testing.T) { @@ -70,10 +103,16 @@ func TestAssembleReplayFramesSnapsNearestSamplesAndOmitsEmptyDrivers(t *testing. }, } - resp, err := assembleReplayFrames(context.Background(), client, 99, 5000) + resp, incomplete, err := assembleReplayFrames(context.Background(), client, 99, 5000) if err != nil { t.Fatalf("assembleReplayFrames() error = %v", err) } + if !incomplete { + t.Fatal("empty entrant location series was labelled complete") + } + if got := replayResponseFreshness(resp, incomplete); got != "partial" { + t.Fatalf("empty entrant freshness = %q", got) + } if resp.SessionKey != 99 || resp.Interval != 5000 { t.Fatalf("response metadata = %+v", resp) } @@ -112,7 +151,7 @@ func TestAssembleReplayFramesCapsFrameCount(t *testing.T) { locs: map[int][]models.Location{1: locs}, } - resp, err := assembleReplayFrames(context.Background(), client, 99, defaultReplayIntervalMS) + resp, _, err := assembleReplayFrames(context.Background(), client, 99, defaultReplayIntervalMS) if err != nil { t.Fatalf("assembleReplayFrames() error = %v", err) } @@ -139,7 +178,7 @@ func TestAssembleReplayFramesBoundsLocationFanOut(t *testing.T) { delay: 5 * time.Millisecond, } - if _, err := assembleReplayFrames(context.Background(), client, 99, defaultReplayIntervalMS); err != nil { + if _, _, err := assembleReplayFrames(context.Background(), client, 99, defaultReplayIntervalMS); err != nil { t.Fatalf("assembleReplayFrames() error = %v", err) } if client.maxInFlight > replayFetchConcurrency { @@ -167,7 +206,7 @@ func TestHandleReplayFramesValidatesParamsAndFloorsInterval(t *testing.T) { client := &fakeReplayClient{drivers: []models.Driver{{DriverNumber: 1}}, locs: map[int][]models.Location{ 1: {{Date: time.Date(2025, 5, 25, 13, 0, 0, 0, time.UTC).Format(time.RFC3339Nano), X: 1, Y: 2}}, }} - resp, err := assembleReplayFrames(context.Background(), client, 99, 1000) + resp, _, err := assembleReplayFrames(context.Background(), client, 99, 1000) if err != nil { t.Fatalf("assembleReplayFrames() error = %v", err) } diff --git a/internal/web/server.go b/internal/web/server.go index 1612966..8aae981 100644 --- a/internal/web/server.go +++ b/internal/web/server.go @@ -159,6 +159,7 @@ func withCORS(next http.Handler) http.Handler { w.Header().Set("Access-Control-Allow-Origin", "*") w.Header().Set("Access-Control-Allow-Methods", "GET, OPTIONS") w.Header().Set("Access-Control-Allow-Headers", "Content-Type") + w.Header().Set("Access-Control-Expose-Headers", dataSourceHeader+", "+dataFreshnessHeader) if r.Method == http.MethodOptions { w.WriteHeader(http.StatusNoContent) return