From fc65dcaf161fb15a62bd19ebffe2442e2d84bece Mon Sep 17 00:00:00 2001 From: CoffeeKnyte <67730400+CoffeeKnyte@users.noreply.github.com> Date: Fri, 31 Jul 2026 00:03:44 +0000 Subject: [PATCH] fix(playback): make session liveness server-observed Client progress reports could keep a zero-byte phantom session alive and, worse, make it outrank a genuinely-serving stream when the enforcer picked over-cap victims. `streammonitor.LiveLocalSessions` substituted `Session.LastActivityAt` for a zero `LastServedAt`, and `LastActivityAt` is advanced by `UpdateProgress` and by the realtime WebSocket hello/ack/result handlers. Since `streamenforcer.selectVictims` keeps the `limit` most-recently-served streams, a progress-only phantom sorted ahead of a real stream and the real one was trimmed instead. Reaping had the same root cause: `sessionIsInactiveLocked` keyed idleness on `LastActivityAt`, so a client that kept pinging held a session open forever. Per decision A5 (Option C), client progress is now UI metadata only and never feeds enforcement or reaping: - `LiveLocalSessions` projects `LastServedAt` verbatim, emitting an empty timestamp when the session has never served, so it sorts as the stalest over-cap victim. - `sessionIsInactiveLocked` measures idleness from `LastServedAt`, falling back only to `StartedAt`. - A configurable never-served window (`DefaultUnservedSessionGrace`, 2m, via `SetUnservedSessionGrace`) keeps a legitimately slow start from being reaped before its first byte, without granting a phantom unbounded life. It is a separate knob rather than a hardcoded floor so it cannot silently override `SetLivenessGracePeriods`. In-flight transports remain exempt, so direct-play and remux long pours and per-segment HLS serves are unaffected. Paused sessions with an open realtime/WebSocket connection are exempt from reaping. That preserves the issue #243 fix (reaping a paused transcode froze clients) while staying within Option C: an open, ping-checked connection is server-observed, unlike a client's reported progress, and the session still consumes one of the user's cap slots. Part 1 of 3 for the Batch 4 liveness/replica work. Part of #305 --- internal/playback/session.go | 56 ++++++--- .../playback/session_paused_grace_test.go | 34 +++-- internal/playback/session_served_test.go | 2 +- .../playback/session_server_liveness_test.go | 118 ++++++++++++++++++ internal/playback/session_test.go | 18 +-- internal/streamenforcer/enforcer_test.go | 41 ++++++ internal/streammonitor/monitor.go | 17 ++- internal/streammonitor/monitor_test.go | 33 ++++- 8 files changed, 272 insertions(+), 47 deletions(-) create mode 100644 internal/playback/session_server_liveness_test.go diff --git a/internal/playback/session.go b/internal/playback/session.go index 56002ca1..b47638d1 100644 --- a/internal/playback/session.go +++ b/internal/playback/session.go @@ -182,6 +182,7 @@ type SessionManager struct { activeGrace time.Duration pausedGrace time.Duration transcodeGrace time.Duration + unservedGrace time.Duration expireHook func(*Session) } @@ -248,6 +249,10 @@ const ( // DefaultTranscodeSessionGrace covers buffer-ahead gaps between segment // requests without depending on local ffmpeg process state. DefaultTranscodeSessionGrace = 10 * time.Minute + + // DefaultUnservedSessionGrace gives a newly-created session time to begin + // serving media before admission and cleanup treat it as inactive. + DefaultUnservedSessionGrace = 2 * time.Minute ) // NewSessionManager creates a SessionManager with the given concurrency limits. @@ -261,6 +266,7 @@ func NewSessionManager(maxStreams, maxTranscodes int) *SessionManager { activeGrace: DefaultActiveSessionGrace, pausedGrace: DefaultPausedSessionGrace, transcodeGrace: DefaultTranscodeSessionGrace, + unservedGrace: DefaultUnservedSessionGrace, } } @@ -280,8 +286,9 @@ func (m *SessionManager) SetAdmissionDecider(decider AdmissionDecider) { m.admissionDecider = decider } -// SetLivenessGracePeriods overrides the grace periods used by admission -// control and stale-session cleanup. +// SetLivenessGracePeriods overrides the active and paused grace periods used by +// admission control and stale-session cleanup. The never-served grace is +// configured independently with SetUnservedSessionGrace. func (m *SessionManager) SetLivenessGracePeriods(active, paused time.Duration) { m.mu.Lock() defer m.mu.Unlock() @@ -303,6 +310,14 @@ func (m *SessionManager) SetTranscodeLivenessGrace(grace time.Duration) { } } +// SetUnservedSessionGrace overrides the idle window for sessions that have not +// served media. A non-positive duration disables this additional grace. +func (m *SessionManager) SetUnservedSessionGrace(grace time.Duration) { + m.mu.Lock() + defer m.mu.Unlock() + m.unservedGrace = grace +} + // SetExpirationHook registers a callback that runs after a session is removed // by stale cleanup. The hook executes outside the manager lock. func (m *SessionManager) SetExpirationHook(fn func(*Session)) { @@ -1305,8 +1320,9 @@ func (m *SessionManager) AllSessions() []*Session { return result } -// CleanExpired removes sessions whose last playback activity exceeds maxIdle. -// Paused sessions receive a 3x grace period for backwards compatibility. +// CleanExpired removes sessions whose last server-observed liveness exceeds +// maxIdle. Paused sessions receive a 3x grace period for backwards +// compatibility. func (m *SessionManager) CleanExpired(maxIdle time.Duration) []*Session { return m.CleanInactive(maxIdle, maxIdle*3) } @@ -1321,9 +1337,9 @@ func (m *SessionManager) CleanStale() []*Session { return m.CleanInactive(active, paused) } -// CleanInactive removes sessions whose last playback activity exceeds the -// provided grace period. Sessions with an active media transport request are -// preserved even if they have not emitted a recent heartbeat yet. +// CleanInactive removes sessions whose last server-observed liveness exceeds +// the provided grace period. Sessions with an active media transport request +// are preserved even if they have not emitted a recent heartbeat yet. func (m *SessionManager) CleanInactive(activeIdle, pausedIdle time.Duration) []*Session { m.mu.Lock() @@ -1371,14 +1387,12 @@ func (m *SessionManager) sessionIsInactiveLocked(s *Session, now time.Time, acti return true } - lastActivity := s.LastActivityAt - if lastActivity.IsZero() { - lastActivity = s.UpdatedAt + if s.IsPaused && (s.HasRealtimeConnection || s.HasWebSocket) { + return false } - if lastActivity.IsZero() { - lastActivity = s.StartedAt - } - if lastActivity.IsZero() { + + livenessAt := s.serverObservedLivenessAt() + if livenessAt.IsZero() { return false } @@ -1388,11 +1402,21 @@ func (m *SessionManager) sessionIsInactiveLocked(s *Session, now time.Time, acti } else if s.PlayMethod == PlayTranscode && m.transcodeGrace > grace { grace = m.transcodeGrace } + if s.LastServedAt.IsZero() && s.activeTransportCount == 0 && m.unservedGrace > grace { + grace = m.unservedGrace + } if grace <= 0 { - return !lastActivity.After(now) + return !livenessAt.After(now) } - return !lastActivity.Add(grace).After(now) + return !livenessAt.Add(grace).After(now) +} + +func (s *Session) serverObservedLivenessAt() time.Time { + if !s.LastServedAt.IsZero() { + return s.LastServedAt + } + return s.StartedAt } // String returns a human-readable summary of a session. diff --git a/internal/playback/session_paused_grace_test.go b/internal/playback/session_paused_grace_test.go index b8f5231e..e794ebcc 100644 --- a/internal/playback/session_paused_grace_test.go +++ b/internal/playback/session_paused_grace_test.go @@ -24,23 +24,34 @@ func TestPausedSessionSurvivesIntentionalPause(t *testing.T) { t.Fatalf("UpdateProgress(paused): %v", err) } - setLastActivity := func(age time.Duration) { + setLastServed := func(age time.Duration) { m.mu.Lock() s := m.sessions[session.ID] - s.LastActivityAt = time.Now().Add(-age) - s.UpdatedAt = s.LastActivityAt + s.LastServedAt = time.Now().Add(-age) m.mu.Unlock() } + if err := m.SetRealtimeConnection(session.ID, true); err != nil { + t.Fatalf("SetRealtimeConnection(true): %v", err) + } + setLastServed(DefaultPausedSessionGrace + time.Minute) + m.CleanStale() + if _, err := m.GetSession(session.ID); err != nil { + t.Fatalf("realtime-connected paused session was reaped: %v", err) + } + if err := m.SetRealtimeConnection(session.ID, false); err != nil { + t.Fatalf("SetRealtimeConnection(false): %v", err) + } + // Paused for 10 minutes: must survive. - setLastActivity(10 * time.Minute) + setLastServed(10 * time.Minute) m.CleanStale() if _, err := m.GetSession(session.ID); err != nil { t.Fatalf("session reaped after 10-minute pause; paused grace must allow intentional pauses (err: %v)", err) } // Abandoned well past the paused grace: must still be reaped. - setLastActivity(DefaultPausedSessionGrace + time.Minute) + setLastServed(DefaultPausedSessionGrace + time.Minute) m.CleanStale() if _, err := m.GetSession(session.ID); err == nil { t.Fatal("session survived past the paused grace; abandoned sessions must still be reaped") @@ -54,12 +65,11 @@ func TestSequentialRangedTransportsSurviveIdleAndPausedGrace(t *testing.T) { t.Fatalf("StartSession: %v", err) } - setLastActivity := func(age time.Duration) { + setLastServed := func(age time.Duration) { t.Helper() m.mu.Lock() s := m.sessions[session.ID] - s.LastActivityAt = time.Now().Add(-age) - s.UpdatedAt = s.LastActivityAt + s.LastServedAt = time.Now().Add(-age) m.mu.Unlock() } assertPresent := func(stage string) { @@ -71,14 +81,14 @@ func TestSequentialRangedTransportsSurviveIdleAndPausedGrace(t *testing.T) { const activeGrace = 2 * time.Minute for cycle := 1; cycle <= 3; cycle++ { - setLastActivity(activeGrace / 2) + setLastServed(activeGrace / 2) m.CleanInactive(activeGrace, DefaultPausedSessionGrace) assertPresent(fmt.Sprintf("idle gap before transport %d", cycle)) if err := m.BeginTransport(session.ID); err != nil { t.Fatalf("BeginTransport(%d): %v", cycle, err) } - setLastActivity(activeGrace + time.Minute) + setLastServed(activeGrace + time.Minute) m.CleanInactive(activeGrace, DefaultPausedSessionGrace) assertPresent(fmt.Sprintf("active transport %d", cycle)) if err := m.EndTransport(session.ID); err != nil { @@ -89,14 +99,14 @@ func TestSequentialRangedTransportsSurviveIdleAndPausedGrace(t *testing.T) { if err := m.UpdateProgress(session.ID, 42, true); err != nil { t.Fatalf("UpdateProgress(paused): %v", err) } - setLastActivity(DefaultPausedSessionGrace - time.Minute) + setLastServed(DefaultPausedSessionGrace - time.Minute) m.CleanInactive(activeGrace, DefaultPausedSessionGrace) assertPresent("late paused idle gap") if err := m.BeginTransport(session.ID); err != nil { t.Fatalf("BeginTransport(late range): %v", err) } - setLastActivity(DefaultPausedSessionGrace + time.Minute) + setLastServed(DefaultPausedSessionGrace + time.Minute) m.CleanInactive(activeGrace, DefaultPausedSessionGrace) assertPresent("late ranged transport active") if err := m.EndTransport(session.ID); err != nil { diff --git a/internal/playback/session_served_test.go b/internal/playback/session_served_test.go index c2606d13..6fe92be4 100644 --- a/internal/playback/session_served_test.go +++ b/internal/playback/session_served_test.go @@ -69,7 +69,7 @@ func TestTranscodeLivenessGraceAndPausedPrecedence(t *testing.T) { } m.mu.Lock() m.sessions[s.ID].IsPaused = tt.paused - m.sessions[s.ID].LastActivityAt = time.Now().Add(-tt.idle) + m.sessions[s.ID].LastServedAt = time.Now().Add(-tt.idle) m.mu.Unlock() reaped := m.CleanStale() if (len(reaped) > 0) != tt.wantReaped { diff --git a/internal/playback/session_server_liveness_test.go b/internal/playback/session_server_liveness_test.go new file mode 100644 index 00000000..635387c0 --- /dev/null +++ b/internal/playback/session_server_liveness_test.go @@ -0,0 +1,118 @@ +package playback + +import ( + "testing" + "time" +) + +func TestNeverServedProgressPingDoesNotPreventReaping(t *testing.T) { + m := NewSessionManager(0, 0) + m.SetLivenessGracePeriods(time.Minute, time.Hour) + m.SetUnservedSessionGrace(2 * time.Minute) + + session, err := m.StartSession(1, "profile", 100, PlayDirect, false) + if err != nil { + t.Fatalf("StartSession: %v", err) + } + m.mu.Lock() + m.sessions[session.ID].StartedAt = time.Now().Add(-90 * time.Second) + m.mu.Unlock() + if err := m.UpdateProgress(session.ID, 42, false); err != nil { + t.Fatalf("UpdateProgress: %v", err) + } + if expired := m.CleanStale(); len(expired) != 0 { + t.Fatalf("expired within configured unserved grace = %+v, want none", expired) + } + + m.mu.Lock() + m.sessions[session.ID].StartedAt = time.Now().Add(-3 * time.Minute) + m.mu.Unlock() + if err := m.UpdateProgress(session.ID, 43, false); err != nil { + t.Fatalf("UpdateProgress after grace: %v", err) + } + expired := m.CleanStale() + if len(expired) != 1 || expired[0].ID != session.ID { + t.Fatalf("expired = %+v, want progress-pinged never-served session %q", expired, session.ID) + } +} + +func TestActiveTransportPreventsServerObservedReaping(t *testing.T) { + m := NewSessionManager(0, 0) + m.SetUnservedSessionGrace(0) + + session, err := m.StartSession(1, "profile", 100, PlayDirect, false) + if err != nil { + t.Fatalf("StartSession: %v", err) + } + if err := m.BeginTransport(session.ID); err != nil { + t.Fatalf("BeginTransport: %v", err) + } + + if expired := m.CleanInactive(0, 0); len(expired) != 0 { + t.Fatalf("expired active transport = %+v, want none", expired) + } +} + +func TestFreshLastServedSurvivesStaleClientActivity(t *testing.T) { + m := NewSessionManager(0, 0) + m.SetUnservedSessionGrace(0) + + session, err := m.StartSession(1, "profile", 100, PlayDirect, false) + if err != nil { + t.Fatalf("StartSession: %v", err) + } + m.mu.Lock() + s := m.sessions[session.ID] + s.LastServedAt = time.Now() + s.LastActivityAt = time.Now().Add(-time.Hour) + s.UpdatedAt = s.LastActivityAt + m.mu.Unlock() + + if expired := m.CleanInactive(time.Minute, time.Hour); len(expired) != 0 { + t.Fatalf("expired freshly-served session = %+v, want none", expired) + } +} + +func TestPausedWebSocketConnectionPreventsServerObservedReaping(t *testing.T) { + m := NewSessionManager(0, 0) + m.SetUnservedSessionGrace(0) + + session, err := m.StartSession(1, "profile", 100, PlayTranscode, false) + if err != nil { + t.Fatalf("StartSession: %v", err) + } + if err := m.UpdateProgress(session.ID, 42, true); err != nil { + t.Fatalf("UpdateProgress(paused): %v", err) + } + if err := m.SetWebSocket(session.ID, true); err != nil { + t.Fatalf("SetWebSocket: %v", err) + } + m.mu.Lock() + m.sessions[session.ID].LastServedAt = time.Now().Add(-time.Hour) + m.mu.Unlock() + + if expired := m.CleanInactive(time.Minute, 30*time.Minute); len(expired) != 0 { + t.Fatalf("expired WebSocket-connected paused session = %+v, want none", expired) + } +} + +func TestPausedWithoutRealtimeConnectionReapsFromLastServed(t *testing.T) { + m := NewSessionManager(0, 0) + m.SetUnservedSessionGrace(0) + + session, err := m.StartSession(1, "profile", 100, PlayTranscode, false) + if err != nil { + t.Fatalf("StartSession: %v", err) + } + m.mu.Lock() + m.sessions[session.ID].LastServedAt = time.Now().Add(-time.Hour) + m.mu.Unlock() + if err := m.UpdateProgress(session.ID, 42, true); err != nil { + t.Fatalf("UpdateProgress(paused): %v", err) + } + + expired := m.CleanInactive(time.Minute, 30*time.Minute) + if len(expired) != 1 || expired[0].ID != session.ID { + t.Fatalf("expired = %+v, want paused disconnected session %q", expired, session.ID) + } +} diff --git a/internal/playback/session_test.go b/internal/playback/session_test.go index 5f6e880e..d20a733d 100644 --- a/internal/playback/session_test.go +++ b/internal/playback/session_test.go @@ -1093,6 +1093,7 @@ func TestSessionManager_LimitCountsIgnoreStaleSessions(t *testing.T) { sm := playback.NewSessionManager(5, 2) sm.SetLivenessGracePeriods(20*time.Millisecond, 40*time.Millisecond) sm.SetTranscodeLivenessGrace(20 * time.Millisecond) + sm.SetUnservedSessionGrace(20 * time.Millisecond) if _, err := sm.StartSession(1, "profile-1", 100, playback.PlayDirect, false); err != nil { t.Fatalf("StartSession direct: %v", err) @@ -1248,6 +1249,7 @@ func TestSessionManager_ZeroLimits(t *testing.T) { func TestSessionManager_CleanExpired(t *testing.T) { sm := playback.NewSessionManager(0, 0) + sm.SetUnservedSessionGrace(0) // Create three sessions. idle, _ := sm.StartSession(1, "prof", 100, playback.PlayDirect, false) @@ -1257,21 +1259,20 @@ func TestSessionManager_CleanExpired(t *testing.T) { // Mark the third session as actively transporting media. _ = sm.BeginTransport(transportActive.ID) - // Send a recent progress update on the "active" session so it stays fresh. + // Send a recent progress update on the "active" session. Client-reported + // progress does not confer server-observed liveness. _ = sm.UpdateProgress(active.ID, 42.0, false) // The "idle" session has no update since StartSession. // The "transportActive" session has an open media transport, so it should // be exempt. - // CleanExpired with 0 duration expires everything without a WebSocket - // that hasn't been updated "recently". Since all sessions were just - // created, use a very short maxIdle to ensure only truly idle ones are - // caught. We simulate staleness by using a large duration that none can - // exceed. + // CleanExpired with 0 duration expires everything without an active + // transport. Disabling the never-served grace above preserves exact control + // for this explicit zero-window cleanup. expired := sm.CleanExpired(0) // 0 means "expire anything older than now" - // Both inactive sessions should be expired (UpdatedAt <= time.Now()). + // Both sessions without a server-observed transport should be expired. // The transport-active session should survive regardless of staleness. if len(expired) != 2 { t.Fatalf("CleanExpired(0) removed %d sessions, want 2", len(expired)) @@ -1319,6 +1320,7 @@ func TestSessionManager_CleanExpired_RespectsMaxIdle(t *testing.T) { func TestSessionManager_CleanInactive_TriggersExpirationHook(t *testing.T) { sm := playback.NewSessionManager(0, 0) + sm.SetUnservedSessionGrace(0) session, err := sm.StartSession(1, "prof", 100, playback.PlayDirect, false) if err != nil { @@ -1347,6 +1349,7 @@ func TestSessionManager_CleanInactive_TriggersExpirationHook(t *testing.T) { func TestSessionManager_CleanExpired_PausedGracePeriod(t *testing.T) { sm := playback.NewSessionManager(0, 0) + sm.SetUnservedSessionGrace(0) // Create two sessions and mark one as paused. playing, _ := sm.StartSession(1, "prof", 100, playback.PlayDirect, false) @@ -1384,6 +1387,7 @@ func TestCheckReplacementAllowedExcludesOnlyTheReplacedSession(t *testing.T) { sm := playback.NewSessionManager(10, 2) sm.SetLivenessGracePeriods(25*time.Millisecond, time.Hour) sm.SetTranscodeLivenessGrace(25 * time.Millisecond) + sm.SetUnservedSessionGrace(25 * time.Millisecond) failed, err := sm.StartSession(1, "profile-1", 100, playback.PlayTranscode, false) if err != nil { diff --git a/internal/streamenforcer/enforcer_test.go b/internal/streamenforcer/enforcer_test.go index d13f08ba..e8d590a8 100644 --- a/internal/streamenforcer/enforcer_test.go +++ b/internal/streamenforcer/enforcer_test.go @@ -142,6 +142,47 @@ func TestEvaluateOnce(t *testing.T) { } } +func TestEvaluateOnceProgressOnlySessionIsVictimBeforeServingStreams(t *testing.T) { + sm := playback.NewSessionManager(0, 0) + phantom, err := sm.StartSession(7, "profile", 1, playback.PlayDirect, false) + if err != nil { + t.Fatalf("StartSession(phantom): %v", err) + } + + for i := 0; i < 2; i++ { + session, startErr := sm.StartSession(7, "profile", 10+i, playback.PlayDirect, false) + if startErr != nil { + t.Fatalf("StartSession(serving %d): %v", i, startErr) + } + if beginErr := sm.BeginTransport(session.ID); beginErr != nil { + t.Fatalf("BeginTransport(serving %d): %v", i, beginErr) + } + if bytesErr := sm.AddServedBytes(session.ID, 1024); bytesErr != nil { + t.Fatalf("AddServedBytes(serving %d): %v", i, bytesErr) + } + if endErr := sm.EndTransport(session.ID); endErr != nil { + t.Fatalf("EndTransport(serving %d): %v", i, endErr) + } + } + + time.Sleep(time.Until(time.Now().Truncate(time.Second).Add(time.Second)) + 10*time.Millisecond) + if err := sm.UpdateProgress(phantom.ID, 30, false); err != nil { + t.Fatalf("UpdateProgress(phantom): %v", err) + } + + source := streammonitor.NewFuncSource(func(context.Context) ([]nodesessions.SessionInfo, error) { + return streammonitor.LiveLocalSessions(sm, "local"), nil + }) + rev := &fakeRevoker{} + e := New(source, func(context.Context, int) (int, error) { return 2, nil }, rev, 0) + if err := e.EvaluateOnce(context.Background()); err != nil { + t.Fatalf("EvaluateOnce: %v", err) + } + if len(rev.revoked) != 1 || rev.revoked[0] != phantom.ID { + t.Fatalf("revoked = %v, want progress-only phantom %q", rev.revoked, phantom.ID) + } +} + func TestEvaluateOnceSourceError(t *testing.T) { rev := &fakeRevoker{} e := New(fakeSource{err: errors.New("scan failed")}, func(context.Context, int) (int, error) { return 1, nil }, rev, 0) diff --git a/internal/streammonitor/monitor.go b/internal/streammonitor/monitor.go index 0c633ef4..ca5ef750 100644 --- a/internal/streammonitor/monitor.go +++ b/internal/streammonitor/monitor.go @@ -13,11 +13,10 @@ // (BeginTransport/EndTransport around every byte pour) — no client progress // required to stay visible. // -// TIMING is a secondary signal. On the edge, LastServedAt is purely byte-observed. -// In integrated mode LastServedAt is mapped from SessionManager.LastActivityAt, -// which client progress reports also advance; that is acceptable because it is -// used only to order over-cap victims (selectVictims), never to decide whether a -// stream exists or is counted. See internal/nodesessions for record production. +// TIMING is server-observed on every path. LastServedAt advances only when the +// server begins, ends, or writes a media transport. A session that has not +// served media projects an empty LastServedAt and sorts as the stalest over-cap +// victim. See internal/nodesessions for record production. package streammonitor import ( @@ -315,9 +314,9 @@ func LiveLocalSessions(sm *playback.SessionManager, nodeName string) []nodesessi live := sm.AllSessions() out := make([]nodesessions.SessionInfo, 0, len(live)) for _, s := range live { - lastServedAt := s.LastServedAt - if lastServedAt.IsZero() { - lastServedAt = s.LastActivityAt + lastServedAt := "" + if !s.LastServedAt.IsZero() { + lastServedAt = s.LastServedAt.UTC().Format(time.RFC3339) } out = append(out, nodesessions.SessionInfo{ SessionID: s.ID, @@ -333,7 +332,7 @@ func LiveLocalSessions(sm *playback.SessionManager, nodeName string) []nodesessi Resolution: s.TargetResolution, HWAccel: s.TranscodeHWAccel, StartedAt: s.StartedAt.UTC().Format(time.RFC3339), - LastServedAt: lastServedAt.UTC().Format(time.RFC3339), + LastServedAt: lastServedAt, BytesServed: s.BytesServed, }) } diff --git a/internal/streammonitor/monitor_test.go b/internal/streammonitor/monitor_test.go index 31dd9a77..5be39585 100644 --- a/internal/streammonitor/monitor_test.go +++ b/internal/streammonitor/monitor_test.go @@ -363,6 +363,16 @@ func TestLiveLocalSessionsMapping(t *testing.T) { t.Fatalf("StartSessionWithContext: %v", err) } session.ClientIP = "192.0.2.10" + if err := sm.BeginTransport(session.ID); err != nil { + t.Fatalf("BeginTransport: %v", err) + } + if err := sm.EndTransport(session.ID); err != nil { + t.Fatalf("EndTransport: %v", err) + } + session, err = sm.GetSession(session.ID) + if err != nil { + t.Fatalf("GetSession: %v", err) + } got := LiveLocalSessions(sm, "local") if len(got) != 1 { t.Fatalf("sessions = %+v", got) @@ -375,7 +385,26 @@ func TestLiveLocalSessionsMapping(t *testing.T) { info.ClientIP != "192.0.2.10" { t.Fatalf("mapped session = %+v", info) } - if info.LastServedAt != session.LastActivityAt.UTC().Format(time.RFC3339) { - t.Fatalf("LastServedAt = %q, want LastActivityAt fallback", info.LastServedAt) + if info.LastServedAt != session.LastServedAt.UTC().Format(time.RFC3339) { + t.Fatalf("LastServedAt = %q, want server-observed timestamp", info.LastServedAt) + } +} + +func TestLiveLocalSessionsDoesNotProjectClientActivityAsLastServed(t *testing.T) { + sm := playback.NewSessionManager(0, 0) + session, err := sm.StartSession(42, "profile-1", 9, playback.PlayDirect, false) + if err != nil { + t.Fatalf("StartSession: %v", err) + } + if err := sm.UpdateProgress(session.ID, 12, false); err != nil { + t.Fatalf("UpdateProgress: %v", err) + } + + got := LiveLocalSessions(sm, "local") + if len(got) != 1 { + t.Fatalf("sessions = %+v", got) + } + if got[0].LastServedAt != "" { + t.Fatalf("LastServedAt = %q, want empty for a never-served session", got[0].LastServedAt) } }