Extracting an embedded subtitle track walks the interleaved container:
subtitle packets sit between video and audio across clusters, so
harvesting a few KB of text means demuxing that stretch of a multi-GB
file off CephFS. Over 41.6h the web player saw p95 120,021ms / max
121,470ms on /api/v1/stream/{session_id}/subtitles/{track}, with 45 of
220 fetches over 10s. The 120s ceiling is the server's absolute
WriteTimeout (cmd/silo/main.go:2389) cutting the body mid-flight -- and
because WriteHeader(200) already ran, those truncations are logged
status=200 and are invisible in error metrics.
The server believes a 600s window bounds this. It does not: -t is passed
as an input option and ffmpeg silently ignores it for these extracts, so
every request runs from the seek point to EOF. Measured against the
production binary, `-ss 4000 -t 30` and `-ss 4000 -t 300` are
byte-identical to passing no -t at all (last cue 02:10:04, end of film).
Only -ss works, so cost tracks (duration - seek) x bitrate.
The window cannot simply be turned on. silo-apple and silo-android both
fetch a track once and depend on receiving the whole thing, so bounding
the output would silently kill subtitles ~10min into every film on both
platforms. The accidental whole-track behaviour is the de-facto contract.
So: keep whole-track delivery, make it cheap. SubtitleCache already had
the right shape for PGS; text was excluded only by the assumption that
"VTT is already windowed and fast", which the inert -t makes false.
Windowing a cached 83KB VTT costs 52ms versus 16s against the original
27GB remux.
Routing on AllowWindow would have poisoned the cache: it is set only in
the PGS branch, so it is always false for text, while streamExtractArgs
applies -ss to any non-ASS/non-PGS source regardless. A seeked subrip
request would take the full-track path, emit seek->EOF, exit cleanly and
publish that partial as canonical -- and every later viewer from 0 would
lose all cues before it (118 of 143 production requests carry a non-zero
seek). Canonicality is now derived from the effective argv instead:
streamExtractPlanFor is the single source of truth for both the argv and
the partial() predicate, so only a seek=0/duration=0 extract can fill.
- key: adds a schema version and the resolved output profile (an ASS
source is reachable as both .ass and .vtt, so format must be keyed)
- removeStaleSiblings now groups by profile, so committing .vtt no longer
deletes a valid .ass sibling; cleanup and eviction no longer hardcode
.sup, so text entries are reclaimed and counted
- InputIsExtractedTrack carries the cached input's format
- text hits use a plain copy with no-store, matching cold-path HTTP
semantics; ServeContent stays on the PGS path only, since Media3 uses
range-capable data sources and responses must not vary with cache warmth
- renames SUP-specific identifiers now that the cache carries text
The inert -t is deliberately retained and its comment corrected in place:
removing it or moving it after -i would bound the output and break the
native clients.
Latency-only: for every (source codec, requested format, seek) the bytes
a client receives are unchanged.
Rationale, measurements, four superseded revisions and the dead ends are
recorded in docs/superpowers/plans/2026-07-17-subtitle-extract-cache.md.
AI-use disclosure: investigated and implemented with AI assistance
(Claude Code + Codex gpt-5.6-sol); Codex's review caught the cache-
poisoning bug above.
875 lines
26 KiB
Go
875 lines
26 KiB
Go
package playback
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"net/http"
|
|
"net/http/httptest"
|
|
"os"
|
|
"path/filepath"
|
|
"strings"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
)
|
|
|
|
// newTestCache builds a cache rooted under a temp transcode dir and returns
|
|
// it with the path of a fake source media file.
|
|
func newTestCache(t *testing.T) (*SubtitleCache, string) {
|
|
t.Helper()
|
|
base := t.TempDir()
|
|
source := filepath.Join(base, "movie.mkv")
|
|
if err := os.WriteFile(source, []byte("fake mkv contents"), 0o644); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
return NewSubtitleCache(func() string { return base }), source
|
|
}
|
|
|
|
// fillEntry populates the cache for source+track with the given payload via
|
|
// the real BeginFill → Tee → Commit path.
|
|
func fillEntry(t *testing.T, c *SubtitleCache, source string, track int, payload string) {
|
|
t.Helper()
|
|
fillOpts(t, c, supExtractOpts(source, track), payload)
|
|
}
|
|
|
|
func fillOpts(t *testing.T, c *SubtitleCache, opts StreamExtractOpts, payload string) {
|
|
t.Helper()
|
|
fill := c.beginFill(opts)
|
|
if fill == nil {
|
|
t.Fatalf("beginFill returned nil for track %d", opts.TrackIndex)
|
|
}
|
|
if _, err := fill.Tee(io.Discard).Write([]byte(payload)); err != nil {
|
|
t.Fatalf("tee write: %v", err)
|
|
}
|
|
if err := fill.Commit(); err != nil {
|
|
t.Fatalf("commit: %v", err)
|
|
}
|
|
}
|
|
|
|
func vttExtractOpts(source string, track int) StreamExtractOpts {
|
|
return StreamExtractOpts{
|
|
InputPath: source,
|
|
TrackIndex: track,
|
|
SourceCodec: "subrip",
|
|
TargetFormat: "vtt",
|
|
}
|
|
}
|
|
|
|
// supExtractOpts builds the base extract options a handler would pass to
|
|
// ServeExtract for a PGS track.
|
|
func supExtractOpts(source string, track int) StreamExtractOpts {
|
|
return StreamExtractOpts{
|
|
InputPath: source,
|
|
TrackIndex: track,
|
|
SourceCodec: "hdmv_pgs_subtitle",
|
|
}
|
|
}
|
|
|
|
// windowedSupOpts builds options for a windowed (?windowed=1) PGS request.
|
|
func windowedSupOpts(source string, track int, seek, duration float64) StreamExtractOpts {
|
|
opts := supExtractOpts(source, track)
|
|
opts.AllowWindow = true
|
|
opts.SeekSeconds = seek
|
|
opts.DurationSeconds = duration
|
|
return opts
|
|
}
|
|
|
|
// waitForCacheEntry polls until the cache holds a committed entry for
|
|
// source+track — used to observe asynchronous background warms.
|
|
func waitForCacheEntry(t *testing.T, c *SubtitleCache, source string, track int) {
|
|
t.Helper()
|
|
deadline := time.Now().Add(5 * time.Second)
|
|
for time.Now().Before(deadline) {
|
|
if f, _, ok := c.Lookup(source, track); ok {
|
|
_ = f.Close()
|
|
return
|
|
}
|
|
time.Sleep(5 * time.Millisecond)
|
|
}
|
|
t.Fatalf("cache entry for track %d never appeared", track)
|
|
}
|
|
|
|
func readAllAndClose(t *testing.T, f *os.File) string {
|
|
t.Helper()
|
|
defer f.Close()
|
|
data, err := io.ReadAll(f)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
return string(data)
|
|
}
|
|
|
|
func TestSubtitleCacheMissThenHit(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
|
|
if _, _, ok := c.Lookup(source, 0); ok {
|
|
t.Fatal("expected miss on empty cache")
|
|
}
|
|
|
|
fillEntry(t, c, source, 0, "PGS DATA TRACK 0")
|
|
|
|
f, modTime, ok := c.Lookup(source, 0)
|
|
if !ok {
|
|
t.Fatal("expected hit after commit")
|
|
}
|
|
if got := readAllAndClose(t, f); got != "PGS DATA TRACK 0" {
|
|
t.Fatalf("cached content = %q", got)
|
|
}
|
|
src, err := os.Stat(source)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if !modTime.Equal(src.ModTime()) {
|
|
t.Fatalf("hit modTime = %v, want source mtime %v", modTime, src.ModTime())
|
|
}
|
|
|
|
// A different track ordinal is a distinct entry.
|
|
if _, _, ok := c.Lookup(source, 1); ok {
|
|
t.Fatal("expected miss for uncached track ordinal")
|
|
}
|
|
}
|
|
|
|
func TestSubtitleCacheInvalidatedBySourceMtime(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
fillEntry(t, c, source, 0, "old extract")
|
|
|
|
newTime := time.Now().Add(2 * time.Hour)
|
|
if err := os.Chtimes(source, newTime, newTime); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, _, ok := c.Lookup(source, 0); ok {
|
|
t.Fatal("expected miss after source mtime changed")
|
|
}
|
|
|
|
// Re-filling under the new source identity overwrites, and the stale
|
|
// sibling entry is cleaned up.
|
|
fillEntry(t, c, source, 0, "new extract")
|
|
f, _, ok := c.Lookup(source, 0)
|
|
if !ok {
|
|
t.Fatal("expected hit after refill")
|
|
}
|
|
if got := readAllAndClose(t, f); got != "new extract" {
|
|
t.Fatalf("cached content = %q", got)
|
|
}
|
|
if n := countCacheEntries(t, c); n != 1 {
|
|
t.Fatalf("stale sibling not removed: %d entries", n)
|
|
}
|
|
}
|
|
|
|
func TestSubtitleCacheInvalidatedBySourceSize(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
fillEntry(t, c, source, 0, "old extract")
|
|
|
|
src, err := os.Stat(source)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := os.WriteFile(source, []byte("different length contents!"), 0o644); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
// Restore the original mtime so only size differs.
|
|
if err := os.Chtimes(source, src.ModTime(), src.ModTime()); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, _, ok := c.Lookup(source, 0); ok {
|
|
t.Fatal("expected miss after source size changed")
|
|
}
|
|
}
|
|
|
|
func TestSubtitleCacheDiscardLeavesNothing(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
fill := c.BeginFill(source, 0)
|
|
if fill == nil {
|
|
t.Fatal("BeginFill returned nil")
|
|
}
|
|
if _, err := fill.Tee(io.Discard).Write([]byte("partial byt")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
fill.Discard()
|
|
|
|
if _, _, ok := c.Lookup(source, 0); ok {
|
|
t.Fatal("discarded fill must not be served")
|
|
}
|
|
if n := countCacheFiles(t, c); n != 0 {
|
|
t.Fatalf("discard left %d files (temp not removed?)", n)
|
|
}
|
|
}
|
|
|
|
func TestSubtitleCacheCommitRefusesChangedSource(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
fill := c.BeginFill(source, 0)
|
|
if fill == nil {
|
|
t.Fatal("BeginFill returned nil")
|
|
}
|
|
if _, err := fill.Tee(io.Discard).Write([]byte("extract from old source")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
// Source replaced mid-extract.
|
|
newTime := time.Now().Add(time.Hour)
|
|
if err := os.Chtimes(source, newTime, newTime); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := fill.Commit(); err == nil {
|
|
t.Fatal("Commit must refuse when source changed mid-fill")
|
|
}
|
|
if n := countCacheFiles(t, c); n != 0 {
|
|
t.Fatalf("refused commit left %d files", n)
|
|
}
|
|
}
|
|
|
|
func TestSubtitleCacheTeeWriteFailureKeepsServingClient(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
fill := c.BeginFill(source, 0)
|
|
if fill == nil {
|
|
t.Fatal("BeginFill returned nil")
|
|
}
|
|
// Force temp-file writes to fail (simulates disk full).
|
|
_ = fill.tmp.Close()
|
|
|
|
var client strings.Builder
|
|
n, err := fill.Tee(&client).Write([]byte("bytes for the viewer"))
|
|
if err != nil || n != len("bytes for the viewer") {
|
|
t.Fatalf("client write must succeed despite cache failure: n=%d err=%v", n, err)
|
|
}
|
|
if client.String() != "bytes for the viewer" {
|
|
t.Fatalf("client got %q", client.String())
|
|
}
|
|
if err := fill.Commit(); err == nil {
|
|
t.Fatal("Commit must fail after tee write error")
|
|
}
|
|
if _, _, ok := c.Lookup(source, 0); ok {
|
|
t.Fatal("failed fill must not be served")
|
|
}
|
|
}
|
|
|
|
func TestSubtitleCacheEvictionUnderCap(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
c.maxBytes = 25 // each payload below is 10 bytes
|
|
|
|
base := time.Now().Add(-time.Hour)
|
|
for track := 0; track < 3; track++ {
|
|
fillEntry(t, c, source, track, fmt.Sprintf("0123456%03d", track))
|
|
// Pin distinct LRU mtimes: track 0 oldest.
|
|
path := entryPath(t, c, source, track)
|
|
mt := base.Add(time.Duration(track) * time.Minute)
|
|
if err := os.Chtimes(path, mt, mt); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
// 4th commit (10 bytes) pushes the total to 40 > 25; eviction must
|
|
// remove the two oldest entries (tracks 0 and 1) to get back to 20.
|
|
fillEntry(t, c, source, 3, "0123456003")
|
|
|
|
for track, want := range map[int]bool{0: false, 1: false, 2: true, 3: true} {
|
|
_, _, ok := c.Lookup(source, track)
|
|
if ok != want {
|
|
t.Errorf("track %d cached = %v, want %v", track, ok, want)
|
|
}
|
|
}
|
|
}
|
|
|
|
func TestSubtitleCacheCoalescing(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
|
|
first := c.BeginFill(source, 0)
|
|
if first == nil {
|
|
t.Fatal("first BeginFill returned nil")
|
|
}
|
|
if second := c.BeginFill(source, 0); second != nil {
|
|
second.Discard()
|
|
t.Fatal("second BeginFill for in-flight track must return nil")
|
|
}
|
|
// A different track is independent.
|
|
other := c.BeginFill(source, 1)
|
|
if other == nil {
|
|
t.Fatal("BeginFill for a different track must not be blocked")
|
|
}
|
|
other.Discard()
|
|
|
|
if _, err := first.Tee(io.Discard).Write([]byte("data")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := first.Commit(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
// Slot released after commit.
|
|
if again := c.BeginFill(source, 0); again == nil {
|
|
t.Fatal("BeginFill must work again after Commit")
|
|
} else {
|
|
again.Discard()
|
|
}
|
|
}
|
|
|
|
func TestSubtitleCacheCoalescingConcurrent(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
|
|
const workers = 16
|
|
var (
|
|
wg sync.WaitGroup
|
|
mu sync.Mutex
|
|
fills []*SubtitleCacheFill
|
|
)
|
|
start := make(chan struct{})
|
|
for i := 0; i < workers; i++ {
|
|
wg.Add(1)
|
|
go func() {
|
|
defer wg.Done()
|
|
<-start
|
|
if f := c.BeginFill(source, 0); f != nil {
|
|
mu.Lock()
|
|
fills = append(fills, f)
|
|
mu.Unlock()
|
|
}
|
|
}()
|
|
}
|
|
close(start)
|
|
wg.Wait()
|
|
|
|
if len(fills) != 1 {
|
|
t.Fatalf("exactly one concurrent BeginFill must win, got %d", len(fills))
|
|
}
|
|
fills[0].Discard()
|
|
}
|
|
|
|
func TestServeExtractCacheFlow(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
|
|
extractCalls := 0
|
|
extract := func(_ context.Context, opts StreamExtractOpts) error {
|
|
extractCalls++
|
|
_, err := opts.Writer.Write([]byte("SUP PAYLOAD"))
|
|
return err
|
|
}
|
|
|
|
// First request: miss → streamed 200 with no-store, entry committed.
|
|
rec := httptest.NewRecorder()
|
|
req := httptest.NewRequest(http.MethodGet, "/sub.sup", nil)
|
|
if err := c.ServeExtract(rec, req, supExtractOpts(source, 0), extract); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if extractCalls != 1 {
|
|
t.Fatalf("extract calls = %d", extractCalls)
|
|
}
|
|
if rec.Body.String() != "SUP PAYLOAD" {
|
|
t.Fatalf("miss body = %q", rec.Body.String())
|
|
}
|
|
if cc := rec.Header().Get("Cache-Control"); cc != "no-store" {
|
|
t.Fatalf("miss Cache-Control = %q", cc)
|
|
}
|
|
|
|
// Second request: hit → served from cache, no extract, revalidatable.
|
|
rec = httptest.NewRecorder()
|
|
if err := c.ServeExtract(rec, req, supExtractOpts(source, 0), extract); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if extractCalls != 1 {
|
|
t.Fatal("cache hit must not invoke extract")
|
|
}
|
|
if rec.Body.String() != "SUP PAYLOAD" {
|
|
t.Fatalf("hit body = %q", rec.Body.String())
|
|
}
|
|
if cc := rec.Header().Get("Cache-Control"); cc != "private, no-cache" {
|
|
t.Fatalf("hit Cache-Control = %q", cc)
|
|
}
|
|
if rec.Header().Get("Last-Modified") == "" {
|
|
t.Fatal("hit must carry Last-Modified")
|
|
}
|
|
if cl := rec.Header().Get("Content-Length"); cl != "11" {
|
|
t.Fatalf("hit Content-Length = %q", cl)
|
|
}
|
|
|
|
// Range request against the cached entry.
|
|
rec = httptest.NewRecorder()
|
|
rangeReq := httptest.NewRequest(http.MethodGet, "/sub.sup", nil)
|
|
rangeReq.Header.Set("Range", "bytes=4-10")
|
|
if err := c.ServeExtract(rec, rangeReq, supExtractOpts(source, 0), extract); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if rec.Code != http.StatusPartialContent || rec.Body.String() != "PAYLOAD" {
|
|
t.Fatalf("range: code=%d body=%q", rec.Code, rec.Body.String())
|
|
}
|
|
}
|
|
|
|
// A seeked text request has effective -ss/-t argv even though AllowWindow is
|
|
// false. Its successful output is only seek->EOF and must never become the
|
|
// canonical cache artifact used by a later viewer starting at zero.
|
|
func TestServeExtractSeekedVTTMissNeverPoisonsCanonicalEntry(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
// Disable detached warms so the test observes only the request-driven
|
|
// cache behavior and can prove the partial response itself was not stored.
|
|
c.warmSem = make(chan struct{})
|
|
|
|
partial := vttExtractOpts(source, 0)
|
|
partial.SeekSeconds = 900
|
|
partial.DurationSeconds = 600
|
|
extract := func(_ context.Context, opts StreamExtractOpts) error {
|
|
payload := "BEGINNING CUE\nLATE CUE"
|
|
if streamExtractPlanFor(opts).partial() {
|
|
payload = "LATE CUE"
|
|
}
|
|
_, err := opts.Writer.Write([]byte(payload))
|
|
return err
|
|
}
|
|
|
|
rec := httptest.NewRecorder()
|
|
req := httptest.NewRequest(http.MethodGet, "/sub.vtt?position=900", nil)
|
|
if err := c.ServeExtract(rec, req, partial, extract); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if rec.Body.String() != "LATE CUE" {
|
|
t.Fatalf("seeked body = %q", rec.Body.String())
|
|
}
|
|
if _, _, ok := c.lookup(vttExtractOpts(source, 0)); ok {
|
|
t.Fatal("seeked VTT extract must not commit a canonical entry")
|
|
}
|
|
|
|
rec = httptest.NewRecorder()
|
|
if err := c.ServeExtract(rec, httptest.NewRequest(http.MethodGet, "/sub.vtt", nil),
|
|
vttExtractOpts(source, 0), extract); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if rec.Body.String() != "BEGINNING CUE\nLATE CUE" {
|
|
t.Fatalf("seek-0 body lost beginning-of-track cues: %q", rec.Body.String())
|
|
}
|
|
}
|
|
|
|
func TestServeExtractSeekedVTTUsesCachedArtifactAndFormat(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
canonical := vttExtractOpts(source, 0)
|
|
fillOpts(t, c, canonical, "WEBVTT\n\nFULL TRACK")
|
|
entry, _, ok := c.cachedEntryPath(canonical)
|
|
if !ok {
|
|
t.Fatal("canonical VTT entry missing")
|
|
}
|
|
|
|
partial := canonical
|
|
partial.SeekSeconds = 1200
|
|
partial.DurationSeconds = 600
|
|
var got StreamExtractOpts
|
|
rec := httptest.NewRecorder()
|
|
err := c.ServeExtract(rec, httptest.NewRequest(http.MethodGet, "/sub.vtt?position=1200", nil), partial,
|
|
func(_ context.Context, opts StreamExtractOpts) error {
|
|
got = opts
|
|
_, err := opts.Writer.Write([]byte("WINDOW"))
|
|
return err
|
|
})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if got.InputPath != entry {
|
|
t.Fatalf("partial input = %q, want cached artifact %q", got.InputPath, entry)
|
|
}
|
|
if got.ExtractedInputFormat != "webvtt" {
|
|
t.Fatalf("cached input format = %q, want webvtt", got.ExtractedInputFormat)
|
|
}
|
|
args := strings.Join(streamExtractArgs(got), " ")
|
|
if !strings.Contains(args, "-f webvtt -i "+entry) || !strings.Contains(args, "-map 0:s:0") {
|
|
t.Fatalf("cached VTT argv does not force format and sole stream: %s", args)
|
|
}
|
|
}
|
|
|
|
func TestSubtitleCacheOutputProfilesCoexist(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
ass := StreamExtractOpts{InputPath: source, TrackIndex: 0, SourceCodec: "ass"}
|
|
vtt := ass
|
|
vtt.TargetFormat = "vtt"
|
|
|
|
fillOpts(t, c, ass, "ASS TRACK")
|
|
fillOpts(t, c, vtt, "WEBVTT TRACK")
|
|
|
|
assFile, _, assOK := c.lookup(ass)
|
|
if !assOK {
|
|
t.Fatal("ASS entry was removed when VTT sibling committed")
|
|
}
|
|
if got := readAllAndClose(t, assFile); got != "ASS TRACK" {
|
|
t.Fatalf("ASS entry = %q", got)
|
|
}
|
|
vttFile, _, vttOK := c.lookup(vtt)
|
|
if !vttOK {
|
|
t.Fatal("VTT entry missing")
|
|
}
|
|
if got := readAllAndClose(t, vttFile); got != "WEBVTT TRACK" {
|
|
t.Fatalf("VTT entry = %q", got)
|
|
}
|
|
|
|
// In-flight coalescing is profile-specific too.
|
|
assFill := c.beginFill(ass)
|
|
vttFill := c.beginFill(vtt)
|
|
if assFill == nil || vttFill == nil {
|
|
t.Fatal("ASS and VTT fills must not coalesce with each other")
|
|
}
|
|
assFill.Discard()
|
|
vttFill.Discard()
|
|
}
|
|
|
|
func TestSubtitleCacheEvictionCountsTextEntries(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
c.maxBytes = 5
|
|
vtt := vttExtractOpts(source, 0)
|
|
fillOpts(t, c, vtt, "123456")
|
|
if _, _, ok := c.lookup(vtt); ok {
|
|
t.Fatal("VTT entry over the cache budget must be evicted")
|
|
}
|
|
}
|
|
|
|
func TestServeExtractContentTypeOnMissAndHit(t *testing.T) {
|
|
tests := []struct {
|
|
name string
|
|
opts func(string) StreamExtractOpts
|
|
contentType string
|
|
path string
|
|
}{
|
|
{"sup", func(source string) StreamExtractOpts { return supExtractOpts(source, 0) }, "application/octet-stream", "/sub.sup"},
|
|
{"ass", func(source string) StreamExtractOpts {
|
|
return StreamExtractOpts{InputPath: source, SourceCodec: "ass"}
|
|
}, "text/x-ssa; charset=utf-8", "/sub.ass"},
|
|
{"vtt", func(source string) StreamExtractOpts { return vttExtractOpts(source, 0) }, "text/vtt; charset=utf-8", "/sub.vtt"},
|
|
}
|
|
for _, tc := range tests {
|
|
t.Run(tc.name, func(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
opts := tc.opts(source)
|
|
extractCalls := 0
|
|
extract := func(_ context.Context, opts StreamExtractOpts) error {
|
|
extractCalls++
|
|
_, err := opts.Writer.Write([]byte("FULL TRACK"))
|
|
return err
|
|
}
|
|
|
|
for request := 0; request < 2; request++ {
|
|
rec := httptest.NewRecorder()
|
|
req := httptest.NewRequest(http.MethodGet, tc.path, nil)
|
|
if request == 1 && tc.name != "sup" {
|
|
req.Header.Set("Range", "bytes=5-9")
|
|
}
|
|
if err := c.ServeExtract(rec, req, opts, extract); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if got := rec.Header().Get("Content-Type"); got != tc.contentType {
|
|
t.Fatalf("request %d Content-Type = %q, want %q", request, got, tc.contentType)
|
|
}
|
|
if tc.name != "sup" {
|
|
if rec.Code != http.StatusOK || rec.Body.String() != "FULL TRACK" {
|
|
t.Fatalf("text cache warmth changed response: code=%d body=%q", rec.Code, rec.Body.String())
|
|
}
|
|
if rec.Header().Get("Cache-Control") != "no-store" || rec.Header().Get("Content-Length") != "" || rec.Header().Get("Last-Modified") != "" {
|
|
t.Fatalf("text response added cache-dependent HTTP headers: %v", rec.Header())
|
|
}
|
|
}
|
|
}
|
|
if extractCalls != 1 {
|
|
t.Fatalf("extract calls = %d, want one miss then one hit", extractCalls)
|
|
}
|
|
})
|
|
}
|
|
}
|
|
|
|
// A windowed request against a cached track must run its extract with the
|
|
// cached .sup as input (small file → near-instant window) instead of
|
|
// re-demuxing the original media, must never publish its sliced output as a
|
|
// cache entry, and must bump the entry's LRU recency.
|
|
func TestServeExtractWindowedUsesCachedTrack(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
fillEntry(t, c, source, 0, "FULL TRACK")
|
|
|
|
// Age the entry so the LRU recency bump is observable.
|
|
entry := entryPath(t, c, source, 0)
|
|
old := time.Now().Add(-time.Hour)
|
|
if err := os.Chtimes(entry, old, old); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
var got StreamExtractOpts
|
|
extractCalls := 0
|
|
rec := httptest.NewRecorder()
|
|
req := httptest.NewRequest(http.MethodGet, "/sub.sup?windowed=1&position=1200&duration=3600", nil)
|
|
err := c.ServeExtract(rec, req, windowedSupOpts(source, 0, 1200, 3600), func(_ context.Context, opts StreamExtractOpts) error {
|
|
extractCalls++
|
|
got = opts
|
|
_, err := opts.Writer.Write([]byte("WINDOW SLICE"))
|
|
return err
|
|
})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if extractCalls != 1 {
|
|
t.Fatalf("extract calls = %d", extractCalls)
|
|
}
|
|
if rec.Body.String() != "WINDOW SLICE" {
|
|
t.Fatalf("windowed body = %q", rec.Body.String())
|
|
}
|
|
if cc := rec.Header().Get("Cache-Control"); cc != "no-store" {
|
|
t.Fatalf("windowed Cache-Control = %q", cc)
|
|
}
|
|
if got.InputPath != entry {
|
|
t.Fatalf("windowed extract input = %q, want cached entry %q", got.InputPath, entry)
|
|
}
|
|
if got.ExtractedInputFormat != "sup" {
|
|
t.Fatalf("windowed extract cached input format = %q, want sup", got.ExtractedInputFormat)
|
|
}
|
|
if got.SeekSeconds != 1200 || got.DurationSeconds != 3600 || !got.AllowWindow {
|
|
t.Fatalf("window parameters not preserved: %+v", got)
|
|
}
|
|
|
|
// The full-track entry must be untouched, with recency bumped.
|
|
f, _, ok := c.Lookup(source, 0)
|
|
if !ok {
|
|
t.Fatal("full-track entry lost")
|
|
}
|
|
if content := readAllAndClose(t, f); content != "FULL TRACK" {
|
|
t.Fatalf("full-track entry corrupted: %q", content)
|
|
}
|
|
info, err := os.Stat(entry)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if !info.ModTime().After(old.Add(time.Minute)) {
|
|
t.Fatalf("windowed serve must bump LRU recency: mtime = %v", info.ModTime())
|
|
}
|
|
}
|
|
|
|
// A windowed miss must trigger exactly one detached background warm no
|
|
// matter how many windowed requests arrive while it runs, and once the warm
|
|
// commits, the next windowed request extracts from the cached track.
|
|
func TestServeExtractWindowedMissWarmsOnce(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
|
|
var (
|
|
mu sync.Mutex
|
|
warmOpts []StreamExtractOpts
|
|
windowOpts []StreamExtractOpts
|
|
warmRelease = make(chan struct{})
|
|
)
|
|
extract := func(_ context.Context, opts StreamExtractOpts) error {
|
|
if opts.AllowWindow {
|
|
mu.Lock()
|
|
windowOpts = append(windowOpts, opts)
|
|
mu.Unlock()
|
|
_, err := opts.Writer.Write([]byte("WINDOW SLICE"))
|
|
return err
|
|
}
|
|
mu.Lock()
|
|
warmOpts = append(warmOpts, opts)
|
|
mu.Unlock()
|
|
<-warmRelease
|
|
_, err := opts.Writer.Write([]byte("FULL TRACK"))
|
|
return err
|
|
}
|
|
|
|
// N windowed misses: each still streams its own windowed slice from the
|
|
// original file; only the first starts a warm (BeginFill coalescing keeps
|
|
// the rest out — deterministic because the in-flight slot is reserved
|
|
// synchronously before ServeExtract returns).
|
|
for i := 0; i < 4; i++ {
|
|
rec := httptest.NewRecorder()
|
|
req := httptest.NewRequest(http.MethodGet, "/sub.sup?windowed=1&position=100&duration=3600", nil)
|
|
if err := c.ServeExtract(rec, req, windowedSupOpts(source, 0, 100, 3600), extract); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if rec.Body.String() != "WINDOW SLICE" {
|
|
t.Fatalf("windowed body = %q", rec.Body.String())
|
|
}
|
|
}
|
|
close(warmRelease)
|
|
waitForCacheEntry(t, c, source, 0)
|
|
|
|
mu.Lock()
|
|
if len(warmOpts) != 1 {
|
|
t.Fatalf("warm extracts = %d, want exactly 1", len(warmOpts))
|
|
}
|
|
warm := warmOpts[0]
|
|
if warm.InputPath != source || warm.SeekSeconds != 0 || warm.DurationSeconds != 0 || warm.AllowWindow || warm.ExtractedInputFormat != "" {
|
|
t.Fatalf("warm must be a full-track extract of the original file: %+v", warm)
|
|
}
|
|
if len(windowOpts) != 4 {
|
|
t.Fatalf("windowed extracts = %d, want 4", len(windowOpts))
|
|
}
|
|
for _, wo := range windowOpts {
|
|
if wo.InputPath != source || wo.ExtractedInputFormat != "" {
|
|
t.Fatalf("pre-warm windowed extract must read the original file: %+v", wo)
|
|
}
|
|
}
|
|
mu.Unlock()
|
|
|
|
// Warm committed → the next windowed request reads the cached track.
|
|
rec := httptest.NewRecorder()
|
|
req := httptest.NewRequest(http.MethodGet, "/sub.sup?windowed=1&position=200&duration=3600", nil)
|
|
if err := c.ServeExtract(rec, req, windowedSupOpts(source, 0, 200, 3600), extract); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
mu.Lock()
|
|
last := windowOpts[len(windowOpts)-1]
|
|
mu.Unlock()
|
|
if last.InputPath != entryPath(t, c, source, 0) || last.ExtractedInputFormat != "sup" {
|
|
t.Fatalf("post-warm windowed extract must read the cached track: %+v", last)
|
|
}
|
|
}
|
|
|
|
// Warms beyond the server-wide slot budget are dropped, not queued, and a
|
|
// dropped warm must not leave an in-flight reservation behind.
|
|
func TestWarmInBackgroundSemaphoreDrop(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
|
|
release := make(chan struct{})
|
|
extract := func(_ context.Context, opts StreamExtractOpts) error {
|
|
<-release
|
|
_, err := opts.Writer.Write([]byte("FULL TRACK"))
|
|
return err
|
|
}
|
|
|
|
// Occupy every warm slot (slots are acquired synchronously).
|
|
for track := 0; track < subtitleCacheWarmSlots; track++ {
|
|
c.WarmInBackground(supExtractOpts(source, track), extract)
|
|
}
|
|
// One more: dropped without reserving the track's in-flight slot.
|
|
overflow := subtitleCacheWarmSlots
|
|
c.WarmInBackground(supExtractOpts(source, overflow), extract)
|
|
if fill := c.BeginFill(source, overflow); fill == nil {
|
|
t.Fatal("dropped warm must not hold the in-flight slot")
|
|
} else {
|
|
fill.Discard()
|
|
}
|
|
|
|
close(release)
|
|
for track := 0; track < subtitleCacheWarmSlots; track++ {
|
|
waitForCacheEntry(t, c, source, track)
|
|
}
|
|
if _, _, ok := c.Lookup(source, overflow); ok {
|
|
t.Fatal("dropped warm must not populate the cache")
|
|
}
|
|
|
|
// With slots free again, the overflow track's warm goes through.
|
|
c.WarmInBackground(supExtractOpts(source, overflow), extract)
|
|
waitForCacheEntry(t, c, source, overflow)
|
|
}
|
|
|
|
// A warm that races an already-in-flight client fill must skip (BeginFill
|
|
// coalescing) and release its warm slot for other tracks.
|
|
func TestWarmInBackgroundSkipsInFlightFill(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
|
|
clientFill := c.BeginFill(source, 0)
|
|
if clientFill == nil {
|
|
t.Fatal("BeginFill returned nil")
|
|
}
|
|
warmed := make(chan struct{}, 1)
|
|
c.WarmInBackground(supExtractOpts(source, 0), func(_ context.Context, opts StreamExtractOpts) error {
|
|
warmed <- struct{}{}
|
|
_, err := opts.Writer.Write([]byte("WARM"))
|
|
return err
|
|
})
|
|
|
|
// The skipped warm must have released its slot synchronously: all
|
|
// subtitleCacheWarmSlots slots are still available.
|
|
for track := 1; track <= subtitleCacheWarmSlots; track++ {
|
|
c.WarmInBackground(supExtractOpts(source, track), func(_ context.Context, opts StreamExtractOpts) error {
|
|
_, err := opts.Writer.Write([]byte("FULL TRACK"))
|
|
return err
|
|
})
|
|
}
|
|
for track := 1; track <= subtitleCacheWarmSlots; track++ {
|
|
waitForCacheEntry(t, c, source, track)
|
|
}
|
|
|
|
select {
|
|
case <-warmed:
|
|
t.Fatal("warm for an in-flight track must not run")
|
|
default:
|
|
}
|
|
clientFill.Discard()
|
|
}
|
|
|
|
func TestServeExtractDiscardsOnExtractError(t *testing.T) {
|
|
c, source := newTestCache(t)
|
|
opts := vttExtractOpts(source, 0)
|
|
|
|
rec := httptest.NewRecorder()
|
|
req := httptest.NewRequest(http.MethodGet, "/sub.vtt", nil)
|
|
wantErr := errors.New("ffmpeg exploded")
|
|
err := c.ServeExtract(rec, req, opts, func(_ context.Context, opts StreamExtractOpts) error {
|
|
_, _ = opts.Writer.Write([]byte("PARTIAL"))
|
|
return wantErr
|
|
})
|
|
if !errors.Is(err, wantErr) {
|
|
t.Fatalf("err = %v", err)
|
|
}
|
|
if _, _, ok := c.lookup(opts); ok {
|
|
t.Fatal("partial extract must not be cached")
|
|
}
|
|
if n := countCacheFiles(t, c); n != 0 {
|
|
t.Fatalf("failed extract left %d files", n)
|
|
}
|
|
}
|
|
|
|
func TestServeExtractNilCacheStreams(t *testing.T) {
|
|
var c *SubtitleCache
|
|
extract := func(_ context.Context, opts StreamExtractOpts) error {
|
|
_, err := opts.Writer.Write([]byte("UNCACHED"))
|
|
return err
|
|
}
|
|
|
|
rec := httptest.NewRecorder()
|
|
req := httptest.NewRequest(http.MethodGet, "/sub.sup", nil)
|
|
if err := c.ServeExtract(rec, req, supExtractOpts("/nonexistent.mkv", 0), extract); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if rec.Body.String() != "UNCACHED" {
|
|
t.Fatalf("body = %q", rec.Body.String())
|
|
}
|
|
|
|
// Windowed requests on a nil cache stream too (no lookup, no warm).
|
|
rec = httptest.NewRecorder()
|
|
req = httptest.NewRequest(http.MethodGet, "/sub.sup?windowed=1&position=10", nil)
|
|
if err := c.ServeExtract(rec, req, windowedSupOpts("/nonexistent.mkv", 0, 10, 3600), extract); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if rec.Body.String() != "UNCACHED" {
|
|
t.Fatalf("windowed body = %q", rec.Body.String())
|
|
}
|
|
}
|
|
|
|
// entryPath computes the committed entry path for source+track.
|
|
func entryPath(t *testing.T, c *SubtitleCache, source string, track int) string {
|
|
t.Helper()
|
|
src, err := os.Stat(source)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
return filepath.Join(c.dir(), subtitleCacheKey(source, track,
|
|
subtitleProfile(supExtractOpts(source, track)), src.ModTime(), src.Size()))
|
|
}
|
|
|
|
// countCacheEntries counts committed .sup entries in the cache dir.
|
|
func countCacheEntries(t *testing.T, c *SubtitleCache) int {
|
|
t.Helper()
|
|
return countMatching(t, c, func(name string) bool {
|
|
return strings.HasSuffix(name, ".sup") && !strings.Contains(name, ".part-")
|
|
})
|
|
}
|
|
|
|
// countCacheFiles counts every file in the cache dir, temp files included.
|
|
func countCacheFiles(t *testing.T, c *SubtitleCache) int {
|
|
t.Helper()
|
|
return countMatching(t, c, func(string) bool { return true })
|
|
}
|
|
|
|
func countMatching(t *testing.T, c *SubtitleCache, match func(string) bool) int {
|
|
t.Helper()
|
|
entries, err := os.ReadDir(c.dir())
|
|
if os.IsNotExist(err) {
|
|
return 0
|
|
}
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
n := 0
|
|
for _, e := range entries {
|
|
if match(e.Name()) {
|
|
n++
|
|
}
|
|
}
|
|
return n
|
|
}
|