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") } }