From 0bba21365290c61bf372d855a9a5089dba9215b6 Mon Sep 17 00:00:00 2001 From: AmanTahiliani Date: Sun, 12 Jul 2026 20:50:18 -0400 Subject: [PATCH] fix(#76): prevent false-fresh aggregate responses --- internal/api/client.go | 21 +++- internal/api/openf1.go | 46 ++++++- internal/api/pacing_test.go | 23 ++++ internal/web/api.go | 148 ++++++++++++++++------- internal/web/championship_hub_test.go | 78 ++++++++++++ internal/web/component_freshness_test.go | 107 ++++++++++++++++ internal/web/context.go | 30 ++++- internal/web/context_test.go | 44 +++++++ internal/web/driversummary.go | 71 +++++------ internal/web/driversummary_test.go | 47 +++++++ internal/web/freshness.go | 15 ++- internal/web/racehub.go | 11 +- internal/web/racehub_test.go | 17 ++- internal/web/replay.go | 26 ++-- internal/web/replay_test.go | 46 ++++++- 15 files changed, 620 insertions(+), 110 deletions(-) create mode 100644 internal/web/component_freshness_test.go diff --git a/internal/api/client.go b/internal/api/client.go index 5bcbfe9..2951519 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,20 @@ 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(): + // Return the unused reservation so repeated bounded enrichment + // cancellations do not leave pacing debt for later real requests. + p.mu.Lock() + p.next = p.next.Add(-p.interval) + p.mu.Unlock() + return ctx.Err() + } } + return nil } type OpenF1Client struct { 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..20d3572 100644 --- a/internal/api/pacing_test.go +++ b/internal/api/pacing_test.go @@ -1,6 +1,7 @@ package api import ( + "context" "encoding/json" "io" "net/http" @@ -49,6 +50,28 @@ func TestRequestPacerNilSafe(t *testing.T) { p.wait() // must not panic } +func TestRequestPacerCancellationReturnsUnusedReservation(t *testing.T) { + p := &requestPacer{interval: 100 * time.Millisecond} + if err := p.waitContext(context.Background()); err != nil { + t.Fatal(err) + } + p.mu.Lock() + wantNext := p.next + p.mu.Unlock() + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Millisecond) + defer cancel() + if err := p.waitContext(ctx); err == nil { + t.Fatal("expected paced wait cancellation") + } + p.mu.Lock() + gotNext := p.next + p.mu.Unlock() + if !gotNext.Equal(wantNext) { + t.Fatalf("cancelled reservation left pacing debt: next %v, want %v", gotNext, wantNext) + } +} + 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/web/api.go b/internal/web/api.go index 07d6490..b2fbc93 100644 --- a/internal/web/api.go +++ b/internal/web/api.go @@ -324,12 +324,13 @@ 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 = client.GetSessionResult(sessionKey) }() - go func() { defer wg.Done(); drivers, _ = client.GetDriversForSession(sessionKey) }() + go func() { defer wg.Done(); drivers, driversErr = client.GetDriversForSession(sessionKey) }() wg.Wait() if resultsErr != nil { @@ -338,18 +339,27 @@ func (s *Server) handleResults(w http.ResponseWriter, r *http.Request) { } 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) } - markOpenF1Response(w, client) + resultsFreshness := "fresh" + if len(results) == 0 { + resultsFreshness = "limited" + } else if incomplete { + resultsFreshness = "partial" + } + markOpenF1Availability(w, client, resultsFreshness) writeJSON(w, enriched) } @@ -400,15 +410,16 @@ 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 = client.GetStartingGrid(sessionKey) }() - go func() { defer wg.Done(); drivers, _ = client.GetDriversForSession(sessionKey) }() + go func() { defer wg.Done(); drivers, driversErr = client.GetDriversForSession(sessionKey) }() wg.Wait() if gridErr != nil { @@ -417,18 +428,27 @@ func (s *Server) handleGrid(w http.ResponseWriter, r *http.Request) { } 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) } - markOpenF1Response(w, client) + gridFreshness := "fresh" + if len(grid) == 0 { + gridFreshness = "limited" + } else if incomplete { + gridFreshness = "partial" + } + markOpenF1Availability(w, client, gridFreshness) writeJSON(w, enriched) } @@ -589,27 +609,30 @@ func (s *Server) handleChampionshipDrivers(w http.ResponseWriter, r *http.Reques return } if len(champ) == 0 { - markOpenF1Response(w, client) + markOpenF1Availability(w, client, "limited") writeJSON(w, []any{}) return } - drivers, _ := 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 := championshipDriverInfo(client, c.SessionKey, c.DriverNumber, driverMap) - if ok { + 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) } - markOpenF1Response(w, client) + markOpenF1AggregateResponse(w, client, incomplete) writeJSON(w, enriched) } @@ -634,7 +657,11 @@ func (s *Server) handleChampionshipTeams(w http.ResponseWriter, r *http.Request) writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return } - markOpenF1Response(w, client) + if len(teams) == 0 { + markOpenF1Availability(w, client, "limited") + } else { + markOpenF1Response(w, client) + } writeJSON(w, teams) } @@ -864,20 +891,29 @@ func (s *Server) openF1ChampionshipHub(client *api.OpenF1Client, year int) (cham return champHubResponse{}, false, err } if len(champ) == 0 { - return champHubResponse{Season: year, RoundLabels: []string{}, Drivers: []champHubDriver{}, Teams: []champHubTeam{}}, false, nil + return champHubResponse{Season: year, RoundLabels: []string{}, Drivers: []champHubDriver{}, Teams: []champHubTeam{}}, true, nil } teams, teamsErr := client.GetTeamChampionshipForYear(year) driverInfo := map[int]models.Driver{} + 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 := fetchSeasonRaces(client, year) if err != nil { return champHubResponse{}, false, err } - incomplete = incomplete || teamsErr != nil + incomplete = incomplete || teamsErr != nil || driversIncomplete resp := aggregateChampionshipHub(year, races, champ, teams, driverInfo) ttl := champHubTTL(year, time.Now()) @@ -923,7 +959,11 @@ func fetchSeasonRaces(client *api.OpenF1Client, year int) (races []meetingRace, } } if raceKey == 0 { - return meetingRace{}, false // not a GP meeting (e.g. pre-season testing) + 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 := client.GetSessionResult(raceKey) grid, gerr := client.GetStartingGrid(raceKey) @@ -1342,23 +1382,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 = 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, _ = client.GetDriversForSession(sessionKey) }() - go func() { defer wg.Done(); rc, _ = client.GetRaceControl(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 { @@ -1375,12 +1417,13 @@ func (s *Server) handleStrategy(w http.ResponseWriter, r *http.Request) { // Non-race sessions have no stints. if len(stints) == 0 { - markOpenF1Response(w, client) + markOpenF1AggregateResponse(w, client, driversErr != nil || rcErr != nil) 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 @@ -1413,6 +1456,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{ @@ -1467,7 +1513,7 @@ func (s *Server) handleStrategy(w http.ResponseWriter, r *http.Request) { return pi < pj }) - markOpenF1Response(w, client) + markOpenF1AggregateResponse(w, client, incomplete) writeJSON(w, strategyResponse{ SessionKey: sessionKey, TotalLaps: totalLaps, @@ -1563,22 +1609,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, _ = client.GetLapsForSession(sessionKey) }() - go func() { defer wg.Done(); stints, _ = client.GetStintsForSession(sessionKey) }() - go func() { defer wg.Done(); pits, _ = client.GetPitStopsForSession(sessionKey) }() - go func() { defer wg.Done(); rc, _ = 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, _ := 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 { @@ -1616,6 +1671,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, @@ -1631,7 +1689,13 @@ func (s *Server) handleLapsComparison(w http.ResponseWriter, r *http.Request) { compDrivers = append(compDrivers, cd) } - markOpenF1Response(w, client) + 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..cb4a23e 100644 --- a/internal/web/championship_hub_test.go +++ b/internal/web/championship_hub_test.go @@ -1,13 +1,91 @@ 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":"Test GP"}]`)) + 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 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..6104e99 --- /dev/null +++ b/internal/web/component_freshness_test.go @@ -0,0 +1,107 @@ +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 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 996489d..d13aa65 100644 --- a/internal/web/context.go +++ b/internal/web/context.go @@ -9,8 +9,8 @@ import ( ) func (s *Server) handleWeekendContext(w http.ResponseWriter, _ *http.Request) { - markLocalResponse(w, false) if !s.hasLocalQuery() { + markDataResponse(w, "none", "limited") writeJSON(w, query.WeekendContext{TemporalState: query.TemporalNoSeason}) return } @@ -30,15 +30,33 @@ func (s *Server) handleWeekendContext(w http.ResponseWriter, _ *http.Request) { writeError(w, err, http.StatusInternalServerError, false) return } - for _, ref := range []*query.ContextSession{context.ActiveSession, context.PreviousCompletedSession, context.DefaultAnalysisSession, context.NextSession} { - if ref != nil && ref.Availability.Freshness != "local" { - markDataResponse(w, ref.Availability.Source, ref.Availability.Freshness) - break - } + 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 61e8679..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) { @@ -111,6 +124,37 @@ func TestWeekendContextHandlerUsesTerminalArchiveAsCompletionEvidence(t *testing } } +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 ebb874c..1b25c67 100644 --- a/internal/web/driversummary.go +++ b/internal/web/driversummary.go @@ -1,6 +1,7 @@ package web import ( + "context" "fmt" "net/http" "strconv" @@ -89,7 +90,7 @@ func (s *Server) handleDriverSummary(w http.ResponseWriter, r *http.Request) { } if ok { if mode != sourceLocal { - tryEnrichDriverSummary(client, &resp, sessionKey) + tryEnrichDriverSummary(r.Context(), client, &resp, sessionKey) } switch resp.Enrichment { case "full": @@ -108,7 +109,7 @@ func (s *Server) handleDriverSummary(w http.ResponseWriter, r *http.Request) { } } - resp, err := openF1DriverSummary(client, year, driverNumber) + resp, incomplete, err := openF1DriverSummary(client, year, driverNumber) if err != nil { writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return @@ -117,7 +118,7 @@ func (s *Server) handleDriverSummary(w http.ResponseWriter, r *http.Request) { http.Error(w, fmt.Sprintf("driver %d not found in %d championship", driverNumber, year), http.StatusNotFound) return } - markOpenF1Response(w, client) + markOpenF1AggregateResponse(w, client, incomplete) writeJSON(w, resp) } @@ -163,10 +164,10 @@ func (s *Server) localDriverSummary(year, driverNumber int) (driverSummaryRespon return resp, sessionKey, true, nil } -func openF1DriverSummary(client *api.OpenF1Client, year, driverNumber int) (*driverSummaryResponse, error) { +func openF1DriverSummary(client *api.OpenF1Client, year, driverNumber int) (*driverSummaryResponse, bool, error) { champ, err := client.GetDriverChampionshipForYear(year) if err != nil { - return nil, err + return nil, false, err } var entry *models.ChampionshipDriver for i := range champ { @@ -176,7 +177,7 @@ func openF1DriverSummary(client *api.OpenF1Client, year, driverNumber int) (*dri } } if entry == nil { - return nil, nil + return nil, false, nil } driverInfo := map[int]models.Driver{} @@ -184,28 +185,41 @@ func openF1DriverSummary(client *api.OpenF1Client, year, driverNumber int) (*dri if ds, derr := client.GetDriversForSession(sessionKey); derr == nil { driverInfo = buildDriverMapFirst(ds) } - if d, ok := championshipDriverInfo(client, 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 := fetchSeasonRaces(client, year) + identityIncomplete := !hasDriverPresentation(driverInfo[driverNumber]) + races, racesIncomplete, err := fetchSeasonRaces(client, year) if err != nil { - return nil, err + return nil, false, err } resp, ok := aggregateDriverSummary(year, driverNumber, races, champ, driverInfo) if !ok { - return nil, nil + return nil, false, nil } + incomplete := identityIncomplete || racesIncomplete resp.Source = "openf1" - resp.Enrichment = "full" - return &resp, nil + 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(client *api.OpenF1Client, resp *driverSummaryResponse, sessionKey int) { +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). @@ -213,28 +227,15 @@ func tryEnrichDriverSummary(client *api.OpenF1Client, resp *driverSummaryRespons return } - type enrichResult struct { - driver models.Driver - ok bool - } - - done := make(chan enrichResult, 1) - go func() { - d, ok := championshipDriverInfo(client, sessionKey, resp.DriverNumber, nil) - done <- enrichResult{driver: d, ok: ok} - }() - - select { - case result := <-done: - if !result.ok { - resp.Enrichment = "limited" - return - } - applyDriverEnrichment(resp, result.driver) - resp.Enrichment = "full" - case <-time.After(driverEnrichmentTimeout): + 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) { diff --git a/internal/web/driversummary_test.go b/internal/web/driversummary_test.go index afb9ed3..97a7037 100644 --- a/internal/web/driversummary_test.go +++ b/internal/web/driversummary_test.go @@ -2,6 +2,7 @@ package web import ( "encoding/json" + "fmt" "net/http" "net/http/httptest" "testing" @@ -212,10 +213,12 @@ func TestHandleDriverSummaryLocalFirstIgnoresHangingEnrichment(t *testing.T) { // 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() { @@ -239,6 +242,11 @@ func TestHandleDriverSummaryLocalFirstIgnoresHangingEnrichment(t *testing.T) { 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 { @@ -323,3 +331,42 @@ func TestHandleDriverSummarySourceLocalOnly(t *testing.T) { 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 index 969bf96..523c881 100644 --- a/internal/web/freshness.go +++ b/internal/web/freshness.go @@ -18,16 +18,23 @@ func markOpenF1Response(w http.ResponseWriter, client *api.OpenF1Client) { } func markOpenF1AggregateResponse(w http.ResponseWriter, client *api.OpenF1Client, partial bool) { + freshness := "fresh" + if partial { + freshness = "partial" + } + markOpenF1Availability(w, client, freshness) +} + +func markOpenF1Availability(w http.ResponseWriter, client *api.OpenF1Client, freshness string) { w.Header().Set(dataSourceHeader, "openf1") if client != nil && client.LastResponseWasStale() { w.Header().Set(dataFreshnessHeader, "stale") return } - if partial { - w.Header().Set(dataFreshnessHeader, "partial") - return + if freshness == "" { + freshness = "fresh" } - w.Header().Set(dataFreshnessHeader, "fresh") + w.Header().Set(dataFreshnessHeader, freshness) } func markDataResponse(w http.ResponseWriter, source, freshness string) { diff --git a/internal/web/racehub.go b/internal/web/racehub.go index 148681f..b00647a 100644 --- a/internal/web/racehub.go +++ b/internal/web/racehub.go @@ -17,7 +17,7 @@ func (s *Server) handleRaceHub(w http.ResponseWriter, r *http.Request) { } if !s.hasLocalQuery() { - markDataResponse(w, "local", "limited") + markDataResponse(w, "none", "limited") writeJSON(w, emptyRaceHub(sessionKey)) return } @@ -27,7 +27,14 @@ func (s *Server) handleRaceHub(w http.ResponseWriter, r *http.Request) { writeError(w, err, http.StatusInternalServerError, false) return } - markLocalResponse(w, hub.Source == query.ResponseSourcePartial) + 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 14b82ce..58babde 100644 --- a/internal/web/racehub_test.go +++ b/internal/web/racehub_test.go @@ -86,7 +86,7 @@ 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) != "local" || rec.Header().Get(dataFreshnessHeader) != "limited" { + 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 { @@ -94,6 +94,21 @@ func TestHandleRaceHubWithoutStore(t *testing.T) { } } +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) diff --git a/internal/web/replay.go b/internal/web/replay.go index f98a2ec..c4d681f 100644 --- a/internal/web/replay.go +++ b/internal/web/replay.go @@ -62,16 +62,26 @@ func (s *Server) handleReplayFrames(w http.ResponseWriter, r *http.Request) { } client := s.client.Scoped() - resp, err := assembleReplayFrames(r.Context(), client, sessionKey, intervalMS) + resp, incomplete, err := assembleReplayFrames(r.Context(), client, sessionKey, intervalMS) if err != nil { writeError(w, err, http.StatusInternalServerError, client.LastResponseWasStale()) return } - markOpenF1Response(w, client) + 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 } @@ -84,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(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..360b6ce 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,13 @@ 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("complete driver series reported incomplete") + } if resp.SessionKey != 99 || resp.Interval != 5000 { t.Fatalf("response metadata = %+v", resp) } @@ -112,7 +148,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 +175,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 +203,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) }