Files
silo-server/internal/playback/transcode_restart_guard_test.go

290 lines
9.1 KiB
Go

package playback
import (
"context"
"os"
"os/exec"
"path/filepath"
"testing"
"time"
)
func TestForwardRestartPreservesConfiguredBackBuffer(t *testing.T) {
truePath, err := exec.LookPath("true")
if err != nil {
t.Skipf("`true` not found in PATH: %v", err)
}
dir := t.TempDir()
old := time.Now().Add(-time.Minute)
for _, segment := range []int{3, 40, 64, 65, 69, 70, 100} {
name := filepath.Join(dir, segmentFilename(segment, TranscodeOpts{}))
if err := os.WriteFile(name, []byte("segment"), 0o600); err != nil {
t.Fatal(err)
}
if err := os.Chtimes(name, old, old); err != nil {
t.Fatal(err)
}
}
session := &TranscodeSession{
outputDir: dir,
lastRequestedSegment: 70,
lastCompletedSegment: 70,
lastPruneFloor: -1,
opts: TranscodeOpts{
OutputDir: dir,
TargetCodecVideo: "h264",
SegmentDuration: 4,
SegmentRetentionSeconds: 20,
StartSegmentNumber: 0,
FFmpegPath: truePath,
},
}
if err := session.Restart(context.Background(), 400, 100); err != nil {
t.Fatalf("Restart: %v", err)
}
// Restart itself has no downstream-completion proof for the replacement
// generation, so it must not unlink even old files synchronously.
for _, segment := range []int{3, 40, 64, 65, 69, 70, 100} {
path := filepath.Join(dir, segmentFilename(segment, TranscodeOpts{}))
if _, err := os.Stat(path); err != nil {
t.Errorf("segment %d was removed during restart: %v", segment, err)
}
}
// Once the replacement generation has produced a complete back buffer,
// the preserved range must become eligible again rather than leaking one
// full window after every forward seek.
session.ReportSegmentDownloaded(105)
waitForPrunerFileMissing(t, filepath.Join(dir, segmentFilename(70, TranscodeOpts{})))
for _, segment := range []int{3, 40, 64, 65, 69, 70} {
path := filepath.Join(dir, segmentFilename(segment, TranscodeOpts{}))
if _, err := os.Stat(path); !os.IsNotExist(err) {
t.Errorf("expired pre-restart segment %d survived completed-download pruning: %v", segment, err)
}
}
if _, err := os.Stat(filepath.Join(dir, segmentFilename(100, TranscodeOpts{}))); err != nil {
t.Fatalf("replacement startup segment was removed: %v", err)
}
}
func TestForwardRestartWithoutCompletedDownloadPreservesFreshFiles(t *testing.T) {
truePath, err := exec.LookPath("true")
if err != nil {
t.Skipf("`true` not found in PATH: %v", err)
}
dir := t.TempDir()
for _, segment := range []int{3, 40, 70} {
if err := os.WriteFile(filepath.Join(dir, segmentFilename(segment, TranscodeOpts{})), []byte("segment"), 0o600); err != nil {
t.Fatal(err)
}
}
session := &TranscodeSession{
outputDir: dir,
lastRequestedSegment: 70,
lastCompletedSegment: -1,
lastPruneFloor: -1,
opts: TranscodeOpts{
OutputDir: dir,
TargetCodecVideo: "h264",
SegmentDuration: 4,
SegmentRetentionSeconds: 20,
StartSegmentNumber: 0,
FFmpegPath: truePath,
},
}
if err := session.Restart(context.Background(), 400, 100); err != nil {
t.Fatalf("Restart: %v", err)
}
time.Sleep(20 * time.Millisecond)
for _, segment := range []int{3, 40, 70} {
if _, err := os.Stat(filepath.Join(dir, segmentFilename(segment, TranscodeOpts{}))); err != nil {
t.Errorf("fresh segment %d was removed without a completed download: %v", segment, err)
}
}
}
func TestForwardRestartKeepsSegmentsWhenRetentionDisabled(t *testing.T) {
truePath, err := exec.LookPath("true")
if err != nil {
t.Skipf("`true` not found in PATH: %v", err)
}
dir := t.TempDir()
for _, segment := range []int{3, 40, 70} {
name := filepath.Join(dir, segmentFilename(segment, TranscodeOpts{}))
if err := os.WriteFile(name, []byte("segment"), 0o600); err != nil {
t.Fatal(err)
}
}
session := &TranscodeSession{
outputDir: dir,
lastRequestedSegment: 70,
lastCompletedSegment: 70,
opts: TranscodeOpts{
OutputDir: dir,
TargetCodecVideo: "h264",
SegmentDuration: 4,
StartSegmentNumber: 0,
FFmpegPath: truePath,
},
}
if err := session.Restart(context.Background(), 400, 100); err != nil {
t.Fatalf("Restart: %v", err)
}
for _, segment := range []int{3, 40, 70} {
path := filepath.Join(dir, segmentFilename(segment, TranscodeOpts{}))
if _, err := os.Stat(path); err != nil {
t.Errorf("disabled retention removed segment %d: %v", segment, err)
}
}
}
// TestSegmentRecoveryDecisionWaitsWhileRestarting covers half of issue #243's
// seek-freeze: while a restart is already in flight, a concurrent segment
// request must WAIT for the restart's output rather than trigger another
// restart. Without this, pipelined HLS segment requests spawn dueling ffmpeg
// restarts that keep preempting the segment the player is blocked on.
func TestSegmentRecoveryDecisionWaitsWhileRestarting(t *testing.T) {
session := &TranscodeSession{
outputDir: t.TempDir(),
restarting: true,
opts: TranscodeOpts{
TargetCodecVideo: "h264",
SegmentDuration: 2,
StartSegmentNumber: 0,
},
}
decision := session.SegmentRecoveryDecision(10, time.Now())
if decision.Reason != "transcode_restarting" {
t.Fatalf("Reason = %q, want transcode_restarting", decision.Reason)
}
if !decision.Wait {
t.Error("Wait = false, want true (concurrent requests must wait out an in-flight restart)")
}
if decision.RestartOnTimeout {
t.Error("RestartOnTimeout = true, want false (a timed-out wait must re-decide, not blindly restart)")
}
}
// TestRestartInvokesRestartHook verifies that a successful Restart fires the
// session's restart hook. The API handler uses the hook to re-arm the
// throttler and the exit monitor; firing it from Restart itself keeps every
// restart caller of a hook-wired session (web segment recovery, audio
// switch) consistent instead of each call site remembering to re-arm by
// hand.
func TestRestartInvokesRestartHook(t *testing.T) {
// `true` starts and exits cleanly, standing in for ffmpeg. Resolve it
// via PATH — it lives in /bin on Linux but /usr/bin on macOS.
truePath, err := exec.LookPath("true")
if err != nil {
t.Skipf("`true` not found in PATH: %v", err)
}
session := &TranscodeSession{
outputDir: t.TempDir(),
opts: TranscodeOpts{
TargetCodecVideo: "h264",
SegmentDuration: 2,
StartSegmentNumber: 0,
FFmpegPath: truePath,
},
}
hookFired := make(chan struct{}, 1)
session.SetRestartHook(func(context.Context) {
hookFired <- struct{}{}
})
if err := session.Restart(context.Background(), 20, 10); err != nil {
t.Fatalf("Restart: %v", err)
}
select {
case <-hookFired:
case <-time.After(2 * time.Second):
t.Fatal("restart hook was not invoked after successful restart")
}
}
func TestRestartCopySeekOriginIsReplacedOrCleared(t *testing.T) {
truePath, err := exec.LookPath("true")
if err != nil {
t.Skipf("`true` not found in PATH: %v", err)
}
newSession := func() *TranscodeSession {
return &TranscodeSession{
outputDir: t.TempDir(),
opts: TranscodeOpts{
TargetCodecVideo: "copy",
SegmentDuration: 2,
SeekSeconds: 18,
StreamOriginSeconds: 10,
CopySeekAnchorResolved: true,
StartSegmentNumber: 5,
FFmpegPath: truePath,
},
}
}
resolved := newSession()
if err := resolved.RestartWithCopySeekAnchor(context.Background(), 100, 48, 96); err != nil {
t.Fatalf("RestartWithCopySeekAnchor: %v", err)
}
resolvedOpts := resolved.Opts()
if resolvedOpts.SeekSeconds != 100 || resolvedOpts.StreamOriginSeconds != 96 ||
!resolvedOpts.CopySeekAnchorResolved || resolvedOpts.StartSegmentNumber != 48 {
t.Fatalf("resolved restart opts = %+v", resolvedOpts)
}
unresolved := newSession()
if err := unresolved.Restart(context.Background(), 100, 50); err != nil {
t.Fatalf("Restart: %v", err)
}
unresolvedOpts := unresolved.Opts()
if unresolvedOpts.StreamOriginSeconds != 0 || unresolvedOpts.CopySeekAnchorResolved {
t.Fatalf("generic restart retained stale copy origin: %+v", unresolvedOpts)
}
}
// TestRestartIsSingleFlight covers the other half: Restart must be
// single-flight per session. A second caller arriving while a restart is in
// progress must return immediately without killing the process the first
// restart just started.
func TestRestartIsSingleFlight(t *testing.T) {
session := &TranscodeSession{
outputDir: t.TempDir(),
restarting: true,
opts: TranscodeOpts{
TargetCodecVideo: "h264",
SegmentDuration: 2,
StartSegmentNumber: 0,
// Nonexistent binary: if the guard is missing and Restart
// proceeds, exec fails and the call returns an error, failing
// the assertions below.
FFmpegPath: "/nonexistent/ffmpeg-single-flight-test",
},
}
err := session.Restart(context.Background(), 20, 10)
if err != nil {
t.Fatalf("Restart during in-flight restart = %v, want nil (single-flight no-op)", err)
}
session.mu.Lock()
restartCount := session.restartCount
stillRestarting := session.restarting
session.mu.Unlock()
if restartCount != 0 {
t.Errorf("restartCount = %d, want 0 (second caller must not perform a restart)", restartCount)
}
if !stillRestarting {
t.Error("restarting flag cleared by no-op caller; must be left for the in-flight restart to clear")
}
}