feat(playback): account served bytes and hold transcode liveness server-side
Integrated deployments recorded no served bytes at all: BytesServed was only ever advanced by the edge writer, so a single-node install showed bytes_served: 0 for every stream. Worse, integrated LastServedAt was mapped from Session.LastActivityAt, which UpdateProgress advances — letting a client influence which of its own streams the over-cap enforcer trims first. - internal/playback: Session gains BytesServed and a distinct LastServedAt. LastServedAt advances only from server-observed events (AddServedBytes, BeginTransport, EndTransport) and never from a client progress report. LastActivityAt keeps its exact prior meaning and reaping semantics. - internal/playback/metered_writer: one shared SessionMeteredWriter for every integrated pour. It forwards io.ReaderFrom so the kernel sendfile path survives, guards the fallback copy against re-entering ReadFrom, and implements Unwrap() so the revocation cut's SetWriteDeadline still reaches the socket. Both properties have regressed on this branch before (GAP-3, GAP-9) and are now pinned by a chain test. - Wired at every integrated pour: native direct-play/remux and transcode segment, jellycompat direct/remux and HLS segment. Each site defers the tail flush; the wrapper alone loses the final partial chunk. - Close VERIFY-4 (buffer-ahead evasion): an unpaused transcode gets a bounded 10m grace measured from the server-observed clock. This covers local, cleanly-completed and offloaded transcodes uniformly, unlike an ffmpeg liveness probe, which sees only local processes and reports false once a copy-mode encode finishes ahead of playback. Paused sessions keep their load-bearing 30m grace. - internal/nodesessions: the same idle window one layer out — transcode records idle out at 180s instead of 60s, applied through a single helper used by ActiveCount, Snapshot and refreshAll so the node's count and status view cannot disagree. Part of the stream monitoring & kill-switch epic.
This commit is contained in:
@@ -34,6 +34,10 @@ 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
|
||||
}
|
||||
out = append(out, nodesessions.SessionInfo{
|
||||
SessionID: s.ID,
|
||||
NodeName: nodeName,
|
||||
@@ -48,7 +52,8 @@ func LiveLocalSessions(sm *playback.SessionManager, nodeName string) []nodesessi
|
||||
Resolution: s.TargetResolution,
|
||||
HWAccel: s.TranscodeHWAccel,
|
||||
StartedAt: s.StartedAt.UTC().Format(time.RFC3339),
|
||||
LastServedAt: s.LastActivityAt.UTC().Format(time.RFC3339),
|
||||
LastServedAt: lastServedAt.UTC().Format(time.RFC3339),
|
||||
BytesServed: s.BytesServed,
|
||||
})
|
||||
}
|
||||
return out
|
||||
|
||||
@@ -3784,7 +3784,10 @@ func (h *PlaybackHandler) HandleGetTranscodeSegment(w http.ResponseWriter, r *ht
|
||||
// the hidden-stream window is observable.
|
||||
slog.Warn("begin transport marker failed", "session", sessionID, "segment", segmentPath, "error", err)
|
||||
}
|
||||
http.ServeFile(w, r, segmentPath)
|
||||
recorder, _ := h.sessionMgr.(playback.ServedBytesRecorder)
|
||||
metered := playback.NewSessionMeteredWriter(w, recorder, sessionID)
|
||||
defer func() { _ = metered.Close() }()
|
||||
http.ServeFile(metered, r, segmentPath)
|
||||
}
|
||||
|
||||
// buildProxyManifestURL signs a stream token carrying the session's full
|
||||
|
||||
@@ -143,7 +143,10 @@ func (h *StreamHandler) HandleStream(w http.ResponseWriter, r *http.Request) {
|
||||
_ = h.sessionMgr.EndTransport(sessionID)
|
||||
}()
|
||||
}
|
||||
if err := playback.ServeDirectPlay(w, r, file.FilePath); err != nil {
|
||||
recorder, _ := h.sessionMgr.(playback.ServedBytesRecorder)
|
||||
metered := playback.NewSessionMeteredWriter(w, recorder, sessionID)
|
||||
defer func() { _ = metered.Close() }()
|
||||
if err := playback.ServeDirectPlay(metered, r, file.FilePath); err != nil {
|
||||
h.handleTransportStartFailure(r.Context(), session, file, err)
|
||||
}
|
||||
|
||||
@@ -153,13 +156,16 @@ func (h *StreamHandler) HandleStream(w http.ResponseWriter, r *http.Request) {
|
||||
_ = h.sessionMgr.EndTransport(sessionID)
|
||||
}()
|
||||
}
|
||||
recorder, _ := h.sessionMgr.(playback.ServedBytesRecorder)
|
||||
metered := playback.NewSessionMeteredWriter(w, recorder, sessionID)
|
||||
defer func() { _ = metered.Close() }()
|
||||
seekSeconds := 0.0
|
||||
if seekStr := r.URL.Query().Get("seek"); seekStr != "" {
|
||||
if s, err := strconv.ParseFloat(seekStr, 64); err == nil && s >= 0 {
|
||||
seekSeconds = s
|
||||
}
|
||||
}
|
||||
if err := playback.ServeRemuxWithDVMode(w, r, file.FilePath, "mp4", seekSeconds, session.TranscodeAudio, session.AudioTrackIndex, file.PrimaryDVProfile(), session.RemuxDVMode, h.ffmpegPath()); err != nil {
|
||||
if err := playback.ServeRemuxWithDVMode(metered, r, file.FilePath, "mp4", seekSeconds, session.TranscodeAudio, session.AudioTrackIndex, file.PrimaryDVProfile(), session.RemuxDVMode, h.ffmpegPath()); err != nil {
|
||||
h.handleTransportStartFailure(r.Context(), session, file, err)
|
||||
}
|
||||
|
||||
|
||||
@@ -157,6 +157,9 @@ func (h *PlaybackHandler) HandleVideoStream(w http.ResponseWriter, r *http.Reque
|
||||
stop := h.Revocation.WatchAndCut(w, playSession.UpstreamSessionID, session.StreamAppUserID, time.Now())
|
||||
defer stop()
|
||||
}
|
||||
recorder, _ := h.sessionMgr.(playback.ServedBytesRecorder)
|
||||
metered := playback.NewSessionMeteredWriter(w, recorder, playSession.UpstreamSessionID)
|
||||
defer func() { _ = metered.Close() }()
|
||||
|
||||
switch method {
|
||||
case "remux":
|
||||
@@ -164,9 +167,9 @@ func (h *PlaybackHandler) HandleVideoStream(w http.ResponseWriter, r *http.Reque
|
||||
if resolvedAudioTrackIndex, ok := compatAudioTrackIndex(*source); ok {
|
||||
audioTrackIndex = resolvedAudioTrackIndex
|
||||
}
|
||||
_ = playback.ServeRemux(w, r, file.FilePath, "mp4", seekSeconds, source.TranscodeAudio, audioTrackIndex, file.PrimaryDVProfile())
|
||||
_ = playback.ServeRemux(metered, r, file.FilePath, "mp4", seekSeconds, source.TranscodeAudio, audioTrackIndex, file.PrimaryDVProfile())
|
||||
default:
|
||||
_ = playback.ServeDirectPlay(w, r, file.FilePath)
|
||||
_ = playback.ServeDirectPlay(metered, r, file.FilePath)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -617,7 +620,10 @@ func (h *PlaybackHandler) HandleHLSSegment(w http.ResponseWriter, r *http.Reques
|
||||
}
|
||||
}
|
||||
|
||||
http.ServeFile(w, r, segmentPath)
|
||||
recorder, _ := h.sessionMgr.(playback.ServedBytesRecorder)
|
||||
metered := playback.NewSessionMeteredWriter(w, recorder, playSession.UpstreamSessionID)
|
||||
defer func() { _ = metered.Close() }()
|
||||
http.ServeFile(metered, r, segmentPath)
|
||||
}
|
||||
|
||||
// hlsSegmentErrorResponse maps a segment-retrieval error to a Jellyfin-faithful
|
||||
|
||||
@@ -19,8 +19,10 @@ const (
|
||||
// namespace this tracker writes instead of duplicating the literal.
|
||||
KeyPrefix = "silo:sessions:"
|
||||
|
||||
sessionTTL = 60 * time.Second
|
||||
refreshInt = 30 * time.Second
|
||||
sessionTTL = 60 * time.Second
|
||||
refreshInt = 30 * time.Second
|
||||
transcodeIdleTTL = 180 * time.Second
|
||||
sessionTypeTranscode = "transcode"
|
||||
)
|
||||
|
||||
// SessionInfo represents an active streaming session stored in Redis.
|
||||
@@ -130,7 +132,7 @@ func (tr *Tracker) ActiveCount() int {
|
||||
if _, dup := tr.sessions[id]; dup {
|
||||
continue
|
||||
}
|
||||
if now.Sub(last) <= sessionTTL {
|
||||
if now.Sub(last) <= idleTTLFor(tr.records[id]) {
|
||||
count++
|
||||
}
|
||||
}
|
||||
@@ -150,7 +152,7 @@ func (tr *Tracker) Snapshot() []SessionInfo {
|
||||
// Session-backed entries are always live until Remove; only ephemeral
|
||||
// (non-session) entries age out by idle timeout.
|
||||
if _, isSession := tr.sessions[id]; !isSession {
|
||||
if last, ok := tr.touched[id]; ok && now.Sub(last) > sessionTTL {
|
||||
if last, ok := tr.touched[id]; ok && now.Sub(last) > idleTTLFor(rec) {
|
||||
continue
|
||||
}
|
||||
}
|
||||
@@ -373,7 +375,7 @@ func (tr *Tracker) refreshAll(ctx context.Context) {
|
||||
}
|
||||
for id, last := range tr.touched {
|
||||
_, isSession := tr.sessions[id]
|
||||
if !isSession && now.Sub(last) > sessionTTL {
|
||||
if !isSession && now.Sub(last) > idleTTLFor(tr.records[id]) {
|
||||
// Idle ephemeral session: stop refreshing and let the Redis key
|
||||
// expire on its own. Session-backed entries (direct/remux) are pruned
|
||||
// on Remove, never by idle timeout, so a quiet-but-open pour stays live.
|
||||
@@ -414,3 +416,10 @@ func (tr *Tracker) refreshAll(ctx context.Context) {
|
||||
slog.DebugContext(ctx, "session refresh pipeline failed", "component", "nodesessions", "error", err)
|
||||
}
|
||||
}
|
||||
|
||||
func idleTTLFor(rec SessionInfo) time.Duration {
|
||||
if rec.Type == sessionTypeTranscode {
|
||||
return transcodeIdleTTL
|
||||
}
|
||||
return sessionTTL
|
||||
}
|
||||
|
||||
@@ -0,0 +1,37 @@
|
||||
package nodesessions
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/redis/go-redis/v9"
|
||||
)
|
||||
|
||||
func TestTranscodeIdleTTLAgreesAcrossCountSnapshotAndRefresh(t *testing.T) {
|
||||
rdb := redis.NewClient(&redis.Options{Addr: "127.0.0.1:1", DialTimeout: 10 * time.Millisecond})
|
||||
t.Cleanup(func() { _ = rdb.Close() })
|
||||
tr := NewTracker(rdb, "http://node", "node", "proxy")
|
||||
idle := time.Now().Add(-90 * time.Second)
|
||||
tr.touched[sessionTypeTranscode] = idle
|
||||
tr.records[sessionTypeTranscode] = SessionInfo{SessionID: sessionTypeTranscode, Type: sessionTypeTranscode}
|
||||
tr.touched["direct"] = idle
|
||||
tr.records["direct"] = SessionInfo{SessionID: "direct", Type: "direct_play"}
|
||||
|
||||
if got := tr.ActiveCount(); got != 1 {
|
||||
t.Fatalf("ActiveCount=%d, want only transcode", got)
|
||||
}
|
||||
snapshot := tr.Snapshot()
|
||||
if len(snapshot) != 1 || snapshot[0].SessionID != sessionTypeTranscode {
|
||||
t.Fatalf("Snapshot=%+v, want only transcode", snapshot)
|
||||
}
|
||||
|
||||
tr.refreshAll(context.Background())
|
||||
tr.mu.Lock()
|
||||
_, transcodeRetained := tr.records[sessionTypeTranscode]
|
||||
_, directRetained := tr.records["direct"]
|
||||
tr.mu.Unlock()
|
||||
if !transcodeRetained || directRetained {
|
||||
t.Fatalf("refresh records: transcode=%v direct=%v, want true/false", transcodeRetained, directRetained)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,80 @@
|
||||
package playback
|
||||
|
||||
import (
|
||||
"io"
|
||||
"net/http"
|
||||
)
|
||||
|
||||
const meteredWriterFlushBytes = 1 << 20
|
||||
|
||||
// ServedBytesRecorder receives server-observed byte counts.
|
||||
type ServedBytesRecorder interface {
|
||||
AddServedBytes(sessionID string, n int64) error
|
||||
}
|
||||
|
||||
// SessionMeteredWriter preserves optional ResponseWriter capabilities while
|
||||
// attributing bytes to a playback session in coarse chunks.
|
||||
type SessionMeteredWriter struct {
|
||||
http.ResponseWriter
|
||||
recorder ServedBytesRecorder
|
||||
sessionID string
|
||||
pending int64
|
||||
}
|
||||
|
||||
func NewSessionMeteredWriter(w http.ResponseWriter, recorder ServedBytesRecorder, sessionID string) *SessionMeteredWriter {
|
||||
return &SessionMeteredWriter{ResponseWriter: w, recorder: recorder, sessionID: sessionID}
|
||||
}
|
||||
|
||||
func (w *SessionMeteredWriter) Write(p []byte) (int, error) {
|
||||
n, err := w.ResponseWriter.Write(p)
|
||||
w.account(int64(n))
|
||||
return n, err
|
||||
}
|
||||
|
||||
func (w *SessionMeteredWriter) ReadFrom(src io.Reader) (int64, error) {
|
||||
if rf, ok := w.ResponseWriter.(io.ReaderFrom); ok {
|
||||
n, err := rf.ReadFrom(src)
|
||||
w.account(n)
|
||||
return n, err
|
||||
}
|
||||
return io.Copy(meteredWriteOnly{w}, src)
|
||||
}
|
||||
|
||||
type meteredWriteOnly struct{ io.Writer }
|
||||
|
||||
func (w *SessionMeteredWriter) account(n int64) {
|
||||
if n <= 0 {
|
||||
return
|
||||
}
|
||||
w.pending += n
|
||||
if w.pending >= meteredWriterFlushBytes {
|
||||
w.flush()
|
||||
}
|
||||
}
|
||||
|
||||
func (w *SessionMeteredWriter) flush() {
|
||||
if w.pending <= 0 {
|
||||
return
|
||||
}
|
||||
if w.recorder != nil {
|
||||
_ = w.recorder.AddServedBytes(w.sessionID, w.pending)
|
||||
}
|
||||
w.pending = 0
|
||||
}
|
||||
|
||||
// Close flushes the final partial chunk. Callers must defer it.
|
||||
func (w *SessionMeteredWriter) Close() error {
|
||||
w.flush()
|
||||
return nil
|
||||
}
|
||||
|
||||
func (w *SessionMeteredWriter) Flush() {
|
||||
if f, ok := w.ResponseWriter.(http.Flusher); ok {
|
||||
f.Flush()
|
||||
}
|
||||
}
|
||||
|
||||
// Unwrap lets http.ResponseController reach the underlying connection.
|
||||
func (w *SessionMeteredWriter) Unwrap() http.ResponseWriter {
|
||||
return w.ResponseWriter
|
||||
}
|
||||
@@ -0,0 +1,96 @@
|
||||
package playback
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/Silo-Server/silo-server/internal/httpstream"
|
||||
)
|
||||
|
||||
type byteRecorder struct{ total int64 }
|
||||
|
||||
func (r *byteRecorder) AddServedBytes(_ string, n int64) error {
|
||||
r.total += n
|
||||
return nil
|
||||
}
|
||||
|
||||
type readerFromResponseWriter struct {
|
||||
header http.Header
|
||||
body bytes.Buffer
|
||||
readFromCalls int
|
||||
deadlineSet bool
|
||||
}
|
||||
|
||||
func (w *readerFromResponseWriter) Header() http.Header { return w.header }
|
||||
func (w *readerFromResponseWriter) WriteHeader(int) {}
|
||||
func (w *readerFromResponseWriter) Write(p []byte) (int, error) { return w.body.Write(p) }
|
||||
func (w *readerFromResponseWriter) SetWriteDeadline(time.Time) error {
|
||||
w.deadlineSet = true
|
||||
return nil
|
||||
}
|
||||
func (w *readerFromResponseWriter) ReadFrom(src io.Reader) (int64, error) {
|
||||
w.readFromCalls++
|
||||
return w.body.ReadFrom(src)
|
||||
}
|
||||
|
||||
type readOnly struct{ io.Reader }
|
||||
|
||||
func TestSessionMeteredWriterPreservesReadFromAndUnwrapChain(t *testing.T) {
|
||||
base := &readerFromResponseWriter{header: make(http.Header)}
|
||||
production := httpstream.NewRollingDeadlineWriter(base)
|
||||
recorder := &byteRecorder{}
|
||||
metered := NewSessionMeteredWriter(production, recorder, "s1")
|
||||
|
||||
payload := bytes.Repeat([]byte("x"), 2<<20)
|
||||
n, err := metered.ReadFrom(readOnly{bytes.NewReader(payload)})
|
||||
if err != nil {
|
||||
t.Fatalf("ReadFrom: %v", err)
|
||||
}
|
||||
if n != int64(len(payload)) || base.readFromCalls == 0 {
|
||||
t.Fatalf("ReadFrom n=%d calls=%d, want %d and fast-path calls", n, base.readFromCalls, len(payload))
|
||||
}
|
||||
if err := metered.Close(); err != nil {
|
||||
t.Fatalf("Close: %v", err)
|
||||
}
|
||||
if recorder.total != int64(len(payload)) {
|
||||
t.Fatalf("recorded bytes=%d, want %d", recorder.total, len(payload))
|
||||
}
|
||||
if metered.Unwrap() != production {
|
||||
t.Fatal("Unwrap did not expose the next production writer")
|
||||
}
|
||||
if err := http.NewResponseController(metered).SetWriteDeadline(time.Now()); err != nil {
|
||||
t.Fatalf("SetWriteDeadline through chain: %v", err)
|
||||
}
|
||||
if !base.deadlineSet {
|
||||
t.Fatal("response controller did not reach base writer")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSessionMeteredWriterFallbackDoesNotRecurseAndFlushesTail(t *testing.T) {
|
||||
base := httptest.NewRecorder()
|
||||
recorder := &byteRecorder{}
|
||||
metered := NewSessionMeteredWriter(base, recorder, "s1")
|
||||
payload := []byte("final partial chunk")
|
||||
|
||||
n, err := metered.ReadFrom(readOnly{bytes.NewReader(payload)})
|
||||
if err != nil {
|
||||
t.Fatalf("ReadFrom fallback: %v", err)
|
||||
}
|
||||
if n != int64(len(payload)) {
|
||||
t.Fatalf("ReadFrom n=%d, want %d", n, len(payload))
|
||||
}
|
||||
if recorder.total != 0 {
|
||||
t.Fatalf("tail flushed before Close: %d", recorder.total)
|
||||
}
|
||||
defer func() { _ = metered.Close() }()
|
||||
if err := metered.Close(); err != nil {
|
||||
t.Fatalf("Close: %v", err)
|
||||
}
|
||||
if recorder.total != int64(len(payload)) {
|
||||
t.Fatalf("tail bytes=%d, want %d", recorder.total, len(payload))
|
||||
}
|
||||
}
|
||||
@@ -56,6 +56,8 @@ type Session struct {
|
||||
StartedAt time.Time
|
||||
UpdatedAt time.Time
|
||||
LastActivityAt time.Time
|
||||
BytesServed int64
|
||||
LastServedAt time.Time
|
||||
activeTransportCount int
|
||||
replacementPlayMethod PlayMethod
|
||||
streamRevision uint64
|
||||
@@ -179,6 +181,7 @@ type SessionManager struct {
|
||||
admissionDecider AdmissionDecider
|
||||
activeGrace time.Duration
|
||||
pausedGrace time.Duration
|
||||
transcodeGrace time.Duration
|
||||
expireHook func(*Session)
|
||||
}
|
||||
|
||||
@@ -241,6 +244,10 @@ const (
|
||||
// pressing Play after a long pause freeze the client (issue #243).
|
||||
// Keep in sync with pausedSessionGrace in internal/worker/cleanup.go.
|
||||
DefaultPausedSessionGrace = 30 * time.Minute
|
||||
|
||||
// DefaultTranscodeSessionGrace covers buffer-ahead gaps between segment
|
||||
// requests without depending on local ffmpeg process state.
|
||||
DefaultTranscodeSessionGrace = 10 * time.Minute
|
||||
)
|
||||
|
||||
// NewSessionManager creates a SessionManager with the given concurrency limits.
|
||||
@@ -248,11 +255,12 @@ const (
|
||||
// maxTranscodes limits concurrent transcode streams per user.
|
||||
func NewSessionManager(maxStreams, maxTranscodes int) *SessionManager {
|
||||
return &SessionManager{
|
||||
sessions: make(map[string]*Session),
|
||||
maxStreams: maxStreams,
|
||||
maxTranscodes: maxTranscodes,
|
||||
activeGrace: DefaultActiveSessionGrace,
|
||||
pausedGrace: DefaultPausedSessionGrace,
|
||||
sessions: make(map[string]*Session),
|
||||
maxStreams: maxStreams,
|
||||
maxTranscodes: maxTranscodes,
|
||||
activeGrace: DefaultActiveSessionGrace,
|
||||
pausedGrace: DefaultPausedSessionGrace,
|
||||
transcodeGrace: DefaultTranscodeSessionGrace,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -286,6 +294,15 @@ func (m *SessionManager) SetLivenessGracePeriods(active, paused time.Duration) {
|
||||
}
|
||||
}
|
||||
|
||||
// SetTranscodeLivenessGrace overrides the unpaused transcode idle window.
|
||||
func (m *SessionManager) SetTranscodeLivenessGrace(grace time.Duration) {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
if grace > 0 {
|
||||
m.transcodeGrace = 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)) {
|
||||
@@ -1112,6 +1129,7 @@ func (m *SessionManager) BeginTransport(sessionID string) error {
|
||||
}
|
||||
|
||||
s.activeTransportCount++
|
||||
s.LastServedAt = time.Now()
|
||||
m.touchSessionLocked(s)
|
||||
return nil
|
||||
}
|
||||
@@ -1130,6 +1148,25 @@ func (m *SessionManager) EndTransport(sessionID string) error {
|
||||
if s.activeTransportCount > 0 {
|
||||
s.activeTransportCount--
|
||||
}
|
||||
s.LastServedAt = time.Now()
|
||||
m.touchSessionLocked(s)
|
||||
return nil
|
||||
}
|
||||
|
||||
// AddServedBytes records bytes observed by a server-side media pour. It never
|
||||
// recreates an unknown or already-reaped session.
|
||||
func (m *SessionManager) AddServedBytes(sessionID string, n int64) error {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
s, ok := m.sessions[sessionID]
|
||||
if !ok {
|
||||
return ErrSessionNotFound
|
||||
}
|
||||
if n <= 0 {
|
||||
return nil
|
||||
}
|
||||
s.BytesServed += n
|
||||
s.LastServedAt = time.Now()
|
||||
m.touchSessionLocked(s)
|
||||
return nil
|
||||
}
|
||||
@@ -1348,6 +1385,8 @@ func (m *SessionManager) sessionIsInactiveLocked(s *Session, now time.Time, acti
|
||||
grace := activeIdle
|
||||
if s.IsPaused {
|
||||
grace = pausedIdle
|
||||
} else if s.PlayMethod == PlayTranscode && m.transcodeGrace > grace {
|
||||
grace = m.transcodeGrace
|
||||
}
|
||||
if grace <= 0 {
|
||||
return !lastActivity.After(now)
|
||||
|
||||
@@ -0,0 +1,80 @@
|
||||
package playback
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestAddServedBytesUnknownDoesNotCreateSession(t *testing.T) {
|
||||
m := NewSessionManager(0, 0)
|
||||
if err := m.AddServedBytes("missing", 100); !errors.Is(err, ErrSessionNotFound) {
|
||||
t.Fatalf("AddServedBytes error=%v, want ErrSessionNotFound", err)
|
||||
}
|
||||
if len(m.AllSessions()) != 0 {
|
||||
t.Fatal("AddServedBytes recreated an unknown session")
|
||||
}
|
||||
}
|
||||
|
||||
func TestLastServedAtOnlyAdvancesForServerObservedEvents(t *testing.T) {
|
||||
m := NewSessionManager(0, 0)
|
||||
s, err := m.StartSession(1, "p", 1, PlayDirect, false)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := m.UpdateProgress(s.ID, 10, false); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
afterProgress, _ := m.GetSession(s.ID)
|
||||
if !afterProgress.LastServedAt.IsZero() {
|
||||
t.Fatalf("UpdateProgress advanced LastServedAt: %v", afterProgress.LastServedAt)
|
||||
}
|
||||
if err := m.AddServedBytes(s.ID, 7); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
afterBytes, _ := m.GetSession(s.ID)
|
||||
if afterBytes.LastServedAt.IsZero() || afterBytes.BytesServed != 7 {
|
||||
t.Fatalf("AddServedBytes did not update served state: %+v", afterBytes)
|
||||
}
|
||||
first := afterBytes.LastServedAt
|
||||
time.Sleep(time.Millisecond)
|
||||
if err := m.BeginTransport(s.ID); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
afterBegin, _ := m.GetSession(s.ID)
|
||||
if !afterBegin.LastServedAt.After(first) {
|
||||
t.Fatal("BeginTransport did not advance LastServedAt")
|
||||
}
|
||||
}
|
||||
|
||||
func TestTranscodeLivenessGraceAndPausedPrecedence(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
method PlayMethod
|
||||
paused bool
|
||||
idle time.Duration
|
||||
wantReaped bool
|
||||
}{
|
||||
{name: "transcode five minutes retained", method: PlayTranscode, idle: 5 * time.Minute},
|
||||
{name: "transcode fifteen minutes reaped", method: PlayTranscode, idle: 15 * time.Minute, wantReaped: true},
|
||||
{name: "paused transcode keeps thirty minutes", method: PlayTranscode, paused: true, idle: 15 * time.Minute},
|
||||
{name: "direct keeps active grace", method: PlayDirect, idle: time.Minute, wantReaped: true},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
m := NewSessionManager(0, 0)
|
||||
s, err := m.StartSession(1, "p", 1, tt.method, false)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
m.mu.Lock()
|
||||
m.sessions[s.ID].IsPaused = tt.paused
|
||||
m.sessions[s.ID].LastActivityAt = time.Now().Add(-tt.idle)
|
||||
m.mu.Unlock()
|
||||
reaped := m.CleanStale()
|
||||
if (len(reaped) > 0) != tt.wantReaped {
|
||||
t.Fatalf("reaped=%d, wantReaped=%v", len(reaped), tt.wantReaped)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -1092,6 +1092,7 @@ func TestSetRealtimeConnection(t *testing.T) {
|
||||
func TestSessionManager_LimitCountsIgnoreStaleSessions(t *testing.T) {
|
||||
sm := playback.NewSessionManager(5, 2)
|
||||
sm.SetLivenessGracePeriods(20*time.Millisecond, 40*time.Millisecond)
|
||||
sm.SetTranscodeLivenessGrace(20 * time.Millisecond)
|
||||
|
||||
if _, err := sm.StartSession(1, "profile-1", 100, playback.PlayDirect, false); err != nil {
|
||||
t.Fatalf("StartSession direct: %v", err)
|
||||
@@ -1110,6 +1111,22 @@ func TestSessionManager_LimitCountsIgnoreStaleSessions(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestSessionManager_LivenessGraceDoesNotOverwriteTranscodeGrace(t *testing.T) {
|
||||
sm := playback.NewSessionManager(5, 2)
|
||||
sm.SetTranscodeLivenessGrace(time.Hour)
|
||||
sm.SetLivenessGracePeriods(20*time.Millisecond, 40*time.Millisecond)
|
||||
|
||||
if _, err := sm.StartSession(1, "profile-1", 101, playback.PlayTranscode, false); err != nil {
|
||||
t.Fatalf("StartSession transcode: %v", err)
|
||||
}
|
||||
|
||||
time.Sleep(30 * time.Millisecond)
|
||||
|
||||
if got := sm.TranscodeCount(1); got != 1 {
|
||||
t.Fatalf("TranscodeCount after active grace = %d, want 1", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSessionManager_ActiveTransportKeepsSessionLive(t *testing.T) {
|
||||
sm := playback.NewSessionManager(5, 2)
|
||||
sm.SetLivenessGracePeriods(20*time.Millisecond, 40*time.Millisecond)
|
||||
@@ -1366,6 +1383,7 @@ func TestSessionManager_CleanExpired_PausedGracePeriod(t *testing.T) {
|
||||
func TestCheckReplacementAllowedExcludesOnlyTheReplacedSession(t *testing.T) {
|
||||
sm := playback.NewSessionManager(10, 2)
|
||||
sm.SetLivenessGracePeriods(25*time.Millisecond, time.Hour)
|
||||
sm.SetTranscodeLivenessGrace(25 * time.Millisecond)
|
||||
|
||||
failed, err := sm.StartSession(1, "profile-1", 100, playback.PlayTranscode, false)
|
||||
if err != nil {
|
||||
|
||||
Reference in New Issue
Block a user