mirror of
https://github.com/AmanTahiliani/box-box.git
synced 2026-08-07 19:56:18 -04:00
feat(#76): harvest request-scoped availability and freshness truth
Backend-only re-cut of the #76 availability work onto main, stacked on the canonical Weekend Context API. Adds request-scoped freshness reporting so aggregate responses cannot report fresh when a component is stale, plus local-first driver summary resolution and cache/pacing truth. The frontend half of #76 is deliberately excluded: it is built on the Weekend shell that failed owner review, including the full-width Partial banner treatment. Availability presentation is re-cut with the shell in #89. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
@@ -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
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
123
internal/api/freshness_test.go
Normal file
123
internal/api/freshness_test.go
Normal file
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user