mirror of
https://github.com/AmanTahiliani/box-box.git
synced 2026-08-07 19:56:18 -04:00
fix: separate live archive snapshots
Parse and expose SessionStatus from the official live feed, treating Started/Resumed as active and terminal or missing statuses as inactive archive candidates. Keep /api/v1/live/state and SSE snapshot data reserved for active sessions while exposing memory-only last_snapshot, last_positions, and last_snapshot_at for the explicit Live tab archive view.
This commit is contained in:
@@ -29,10 +29,21 @@ type SSEHub struct {
|
||||
deregister chan *sseClient
|
||||
broadcast chan sseEvent
|
||||
|
||||
mu sync.RWMutex
|
||||
lastSnapshot *live.LiveStreamData
|
||||
lastPositions map[string]live.LivePositionData
|
||||
isLive bool
|
||||
mu sync.RWMutex
|
||||
activeSnapshot *live.LiveStreamData
|
||||
activePositions map[string]live.LivePositionData
|
||||
lastSnapshot *live.LiveStreamData
|
||||
lastPositions map[string]live.LivePositionData
|
||||
lastSnapshotAt time.Time
|
||||
isLive bool
|
||||
}
|
||||
|
||||
type liveStatePayload struct {
|
||||
IsLive bool `json:"is_live"`
|
||||
Data *live.LiveStreamData `json:"data"`
|
||||
LastSnapshot *live.LiveStreamData `json:"last_snapshot,omitempty"`
|
||||
LastPositions map[string]live.LivePositionData `json:"last_positions,omitempty"`
|
||||
LastSnapshotAt *time.Time `json:"last_snapshot_at,omitempty"`
|
||||
}
|
||||
|
||||
func newSSEHub() *SSEHub {
|
||||
@@ -51,19 +62,16 @@ func (h *SSEHub) run() {
|
||||
case c := <-h.register:
|
||||
clients[c] = true
|
||||
// Send catch-up snapshot so new clients see current state immediately.
|
||||
h.mu.RLock()
|
||||
snap := h.lastSnapshot
|
||||
positions := cloneLivePositions(h.lastPositions)
|
||||
live := h.isLive
|
||||
h.mu.RUnlock()
|
||||
if snap != nil {
|
||||
if data, err := json.Marshal(map[string]any{"data": snap, "is_live": live}); err == nil {
|
||||
state := h.State()
|
||||
if state.Data != nil || state.LastSnapshot != nil {
|
||||
if data, err := json.Marshal(state); err == nil {
|
||||
select {
|
||||
case c.ch <- formatSSEFrame("snapshot", data):
|
||||
default:
|
||||
}
|
||||
}
|
||||
}
|
||||
positions := h.ActivePositions()
|
||||
if len(positions) > 0 {
|
||||
if data, err := json.Marshal(positions); err == nil {
|
||||
select {
|
||||
@@ -96,11 +104,125 @@ func formatSSEFrame(event string, data []byte) []byte {
|
||||
return []byte(fmt.Sprintf("event: %s\ndata: %s\n\n", event, data))
|
||||
}
|
||||
|
||||
// Snapshot returns the latest live data snapshot and whether a session is active.
|
||||
func (h *SSEHub) Snapshot() (*live.LiveStreamData, bool) {
|
||||
// State returns the active live snapshot and the retained in-memory archive.
|
||||
func (h *SSEHub) State() liveStatePayload {
|
||||
h.mu.RLock()
|
||||
defer h.mu.RUnlock()
|
||||
return h.lastSnapshot, h.isLive
|
||||
payload := liveStatePayload{
|
||||
IsLive: h.isLive,
|
||||
}
|
||||
if h.isLive {
|
||||
payload.Data = h.activeSnapshot
|
||||
} else if h.lastSnapshot != nil {
|
||||
payload.LastSnapshot = h.lastSnapshot
|
||||
payload.LastPositions = cloneLivePositions(h.lastPositions)
|
||||
if !h.lastSnapshotAt.IsZero() {
|
||||
at := h.lastSnapshotAt
|
||||
payload.LastSnapshotAt = &at
|
||||
}
|
||||
}
|
||||
return payload
|
||||
}
|
||||
|
||||
func (h *SSEHub) ActivePositions() map[string]live.LivePositionData {
|
||||
h.mu.RLock()
|
||||
defer h.mu.RUnlock()
|
||||
return cloneLivePositions(h.activePositions)
|
||||
}
|
||||
|
||||
func (h *SSEHub) applySnapshot(data live.LiveStreamData, now time.Time) liveStatePayload {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
|
||||
if live.SessionStatusIsActive(data.SessionStatus) {
|
||||
h.isLive = true
|
||||
h.activeSnapshot = &data
|
||||
if data.PositionUpdated && len(data.Positions) > 0 {
|
||||
h.activePositions = cloneLivePositions(data.Positions)
|
||||
}
|
||||
return liveStatePayload{IsLive: true, Data: h.activeSnapshot}
|
||||
}
|
||||
|
||||
archivePositions := cloneLivePositions(h.activePositions)
|
||||
if data.PositionUpdated && len(data.Positions) > 0 {
|
||||
archivePositions = cloneLivePositions(data.Positions)
|
||||
}
|
||||
h.isLive = false
|
||||
h.activeSnapshot = nil
|
||||
h.activePositions = nil
|
||||
if hasLiveSnapshotData(data) {
|
||||
h.lastSnapshot = &data
|
||||
h.lastSnapshotAt = now
|
||||
h.lastPositions = archivePositions
|
||||
}
|
||||
return h.stateLocked()
|
||||
}
|
||||
|
||||
func (h *SSEHub) applyPositions(data live.LiveStreamData) map[string]live.LivePositionData {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
|
||||
if len(data.Positions) == 0 {
|
||||
return nil
|
||||
}
|
||||
positions := cloneLivePositions(data.Positions)
|
||||
if h.isLive {
|
||||
h.activePositions = positions
|
||||
return positions
|
||||
}
|
||||
if h.lastSnapshot != nil {
|
||||
h.lastPositions = positions
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (h *SSEHub) deactivate(now time.Time) liveStatePayload {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
|
||||
if h.activeSnapshot != nil {
|
||||
h.lastSnapshot = h.activeSnapshot
|
||||
h.lastPositions = cloneLivePositions(h.activePositions)
|
||||
h.lastSnapshotAt = now
|
||||
}
|
||||
h.isLive = false
|
||||
h.activeSnapshot = nil
|
||||
h.activePositions = nil
|
||||
return h.stateLocked()
|
||||
}
|
||||
|
||||
func (h *SSEHub) stateLocked() liveStatePayload {
|
||||
payload := liveStatePayload{IsLive: h.isLive}
|
||||
if h.isLive {
|
||||
payload.Data = h.activeSnapshot
|
||||
return payload
|
||||
}
|
||||
if h.lastSnapshot != nil {
|
||||
payload.LastSnapshot = h.lastSnapshot
|
||||
payload.LastPositions = cloneLivePositions(h.lastPositions)
|
||||
if !h.lastSnapshotAt.IsZero() {
|
||||
at := h.lastSnapshotAt
|
||||
payload.LastSnapshotAt = &at
|
||||
}
|
||||
}
|
||||
return payload
|
||||
}
|
||||
|
||||
func hasLiveSnapshotData(data live.LiveStreamData) bool {
|
||||
return len(data.Drivers) > 0 ||
|
||||
len(data.DriverInfo) > 0 ||
|
||||
len(data.Tyres) > 0 ||
|
||||
len(data.Telemetry) > 0 ||
|
||||
len(data.RCMessages) > 0 ||
|
||||
len(data.TeamRadio) > 0 ||
|
||||
len(data.Stints) > 0 ||
|
||||
data.Session.MeetingName != "" ||
|
||||
data.Session.SessionName != "" ||
|
||||
data.Session.SessionType != "" ||
|
||||
data.TrackStatus != "" ||
|
||||
data.CurrentLap != 0 ||
|
||||
data.TotalLaps != 0 ||
|
||||
data.Clock != ""
|
||||
}
|
||||
|
||||
// runLiveFeeds launches background goroutines for the F1 SignalR feed and keepalive.
|
||||
@@ -129,13 +251,8 @@ func (s *Server) signalRLoop() {
|
||||
log.Printf("web: live feed ended: %v", err)
|
||||
}
|
||||
|
||||
s.hub.mu.Lock()
|
||||
s.hub.isLive = false
|
||||
s.hub.lastSnapshot = nil
|
||||
s.hub.lastPositions = nil
|
||||
s.hub.mu.Unlock()
|
||||
|
||||
if payload, err := json.Marshal(map[string]any{"data": nil, "is_live": false}); err == nil {
|
||||
state := s.hub.deactivate(time.Now())
|
||||
if payload, err := json.Marshal(state); err == nil {
|
||||
s.hub.broadcast <- sseEvent{name: "snapshot", data: payload}
|
||||
}
|
||||
|
||||
@@ -167,24 +284,20 @@ func (s *Server) connectAndDrain() error {
|
||||
select {
|
||||
case data := <-dataChan:
|
||||
now := time.Now()
|
||||
s.hub.mu.Lock()
|
||||
if data.SnapshotUpdated {
|
||||
s.hub.lastSnapshot = &data
|
||||
}
|
||||
if data.PositionUpdated && len(data.Positions) > 0 {
|
||||
s.hub.lastPositions = cloneLivePositions(data.Positions)
|
||||
}
|
||||
s.hub.isLive = true
|
||||
s.hub.mu.Unlock()
|
||||
|
||||
if data.SnapshotUpdated {
|
||||
if payload, err := json.Marshal(map[string]any{"data": data, "is_live": true}); err == nil {
|
||||
state := s.hub.applySnapshot(data, now)
|
||||
if payload, err := json.Marshal(state); err == nil {
|
||||
s.hub.broadcast <- sseEvent{name: "snapshot", data: payload}
|
||||
}
|
||||
}
|
||||
|
||||
if data.PositionUpdated && len(data.Positions) > 0 && now.Sub(lastPositionBroadcast) >= 250*time.Millisecond {
|
||||
if payload, err := json.Marshal(data.Positions); err == nil {
|
||||
positions := s.hub.applyPositions(data)
|
||||
if len(positions) > 0 {
|
||||
payload, err := json.Marshal(positions)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
s.hub.broadcast <- sseEvent{name: "positions", data: payload}
|
||||
lastPositionBroadcast = now
|
||||
}
|
||||
@@ -217,11 +330,7 @@ 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) {
|
||||
snap, isLive := s.hub.Snapshot()
|
||||
writeJSON(w, map[string]any{
|
||||
"is_live": isLive,
|
||||
"data": snap,
|
||||
})
|
||||
writeJSON(w, s.hub.State())
|
||||
}
|
||||
|
||||
// handleSSEStream is the persistent SSE endpoint for live data.
|
||||
|
||||
90
internal/web/live_archive_test.go
Normal file
90
internal/web/live_archive_test.go
Normal file
@@ -0,0 +1,90 @@
|
||||
package web
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/AmanTahiliani/box-box/internal/live"
|
||||
)
|
||||
|
||||
func TestSSEHubArchivesTerminalSessionSnapshot(t *testing.T) {
|
||||
hub := newSSEHub()
|
||||
now := time.Date(2026, 7, 4, 14, 0, 0, 0, time.UTC)
|
||||
active := live.LiveStreamData{
|
||||
SessionStatus: "Started",
|
||||
Drivers: map[string]live.LiveDriverData{
|
||||
"1": {RacingNumber: "1", Position: 1},
|
||||
},
|
||||
Positions: map[string]live.LivePositionData{
|
||||
"1": {X: 100, Y: -50, Z: 2, Status: "OnTrack"},
|
||||
},
|
||||
PositionUpdated: true,
|
||||
SnapshotUpdated: true,
|
||||
}
|
||||
if state := hub.applySnapshot(active, now); !state.IsLive || state.Data == nil {
|
||||
t.Fatalf("active state = %+v, want live data", state)
|
||||
}
|
||||
|
||||
terminal := active
|
||||
terminal.SessionStatus = "Finished"
|
||||
terminal.Positions = nil
|
||||
terminal.PositionUpdated = false
|
||||
state := hub.applySnapshot(terminal, now.Add(time.Minute))
|
||||
|
||||
if state.IsLive {
|
||||
t.Fatal("terminal SessionStatus should not be live")
|
||||
}
|
||||
if state.Data != nil {
|
||||
t.Fatalf("inactive state data = %+v, want nil", state.Data)
|
||||
}
|
||||
if state.LastSnapshot == nil || state.LastSnapshot.SessionStatus != "Finished" {
|
||||
t.Fatalf("last snapshot = %+v, want terminal snapshot", state.LastSnapshot)
|
||||
}
|
||||
if got := state.LastPositions["1"]; got.X != 100 || got.Status != "OnTrack" {
|
||||
t.Fatalf("last positions = %+v, want carried active positions", state.LastPositions)
|
||||
}
|
||||
if state.LastSnapshotAt == nil || !state.LastSnapshotAt.Equal(now.Add(time.Minute)) {
|
||||
t.Fatalf("last snapshot time = %v, want %v", state.LastSnapshotAt, now.Add(time.Minute))
|
||||
}
|
||||
}
|
||||
|
||||
func TestHandleLiveStateKeepsArchiveOutOfActiveData(t *testing.T) {
|
||||
hub := newSSEHub()
|
||||
now := time.Date(2026, 7, 4, 14, 0, 0, 0, time.UTC)
|
||||
hub.applySnapshot(live.LiveStreamData{
|
||||
SessionStatus: "Finished",
|
||||
Session: live.LiveSessionMeta{MeetingName: "British Grand Prix", SessionName: "Race"},
|
||||
Drivers: map[string]live.LiveDriverData{
|
||||
"44": {RacingNumber: "44", Position: 1},
|
||||
},
|
||||
SnapshotUpdated: true,
|
||||
}, now)
|
||||
|
||||
srv := &Server{hub: hub}
|
||||
req := httptest.NewRequest(http.MethodGet, "/api/v1/live/state", nil)
|
||||
rec := httptest.NewRecorder()
|
||||
srv.handleLiveState(rec, req)
|
||||
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("status = %d, want 200", rec.Code)
|
||||
}
|
||||
var resp liveStatePayload
|
||||
if err := json.Unmarshal(rec.Body.Bytes(), &resp); err != nil {
|
||||
t.Fatalf("decode response: %v", err)
|
||||
}
|
||||
if resp.IsLive {
|
||||
t.Fatal("archived snapshot should report is_live=false")
|
||||
}
|
||||
if resp.Data != nil {
|
||||
t.Fatalf("archived snapshot leaked into data: %+v", resp.Data)
|
||||
}
|
||||
if resp.LastSnapshot == nil || resp.LastSnapshot.Session.MeetingName != "British Grand Prix" {
|
||||
t.Fatalf("last snapshot = %+v, want archived race", resp.LastSnapshot)
|
||||
}
|
||||
if resp.LastSnapshotAt == nil {
|
||||
t.Fatal("last_snapshot_at should be present for archived snapshots")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user