mirror of
https://github.com/truewhile/MeBox.git
synced 2026-10-04 20:46:37 +08:00
bug处理
This commit is contained in:
@@ -24,11 +24,17 @@ func (s *StreamService) ServeHLSPlaylist(w http.ResponseWriter, r *http.Request,
|
||||
return ErrTranscodeDisabled
|
||||
}
|
||||
startSec := parseHLSStartSec(r)
|
||||
if _, err := s.transcoder.EnsureJobFrom(r.Context(), mediaID, startSec); err != nil {
|
||||
seekGen := parseHLSSeekGen(r)
|
||||
if _, err := s.transcoder.EnsureJobFrom(r.Context(), mediaID, startSec, seekGen); err != nil {
|
||||
return err
|
||||
}
|
||||
s.transcoder.TouchJob(mediaID)
|
||||
if !s.transcoder.WaitReady(r.Context(), mediaID, 30*time.Second) {
|
||||
readyTimeout := 45 * time.Second
|
||||
if startSec > 0.05 {
|
||||
// Mid-file HTTP seeks (esp. WMV) need longer before the first segment appears.
|
||||
readyTimeout = 120 * time.Second
|
||||
}
|
||||
if !s.transcoder.WaitReady(r.Context(), mediaID, readyTimeout) {
|
||||
return errors.New("hls playlist not ready")
|
||||
}
|
||||
playlist := s.transcoder.PlaylistPath(mediaID)
|
||||
@@ -69,6 +75,21 @@ func parseHLSStartSec(r *http.Request) float64 {
|
||||
return v
|
||||
}
|
||||
|
||||
func parseHLSSeekGen(r *http.Request) int64 {
|
||||
if r == nil {
|
||||
return 0
|
||||
}
|
||||
raw := strings.TrimSpace(r.URL.Query().Get("_seek"))
|
||||
if raw == "" {
|
||||
return 0
|
||||
}
|
||||
v, err := strconv.ParseInt(raw, 10, 64)
|
||||
if err != nil || v < 0 {
|
||||
return 0
|
||||
}
|
||||
return v
|
||||
}
|
||||
|
||||
func appendQueryToHLSSegments(playlist, rawQuery string) string {
|
||||
if strings.TrimSpace(rawQuery) == "" {
|
||||
return playlist
|
||||
|
||||
@@ -68,6 +68,9 @@ type hlsJob struct {
|
||||
// startSec is the source seek offset fed to ffmpeg (-ss). The HLS
|
||||
// playlist itself always starts at t=0 for that session.
|
||||
startSec float64
|
||||
// seekGen is the client `_seek` token for this job. Newer gens win;
|
||||
// older/missing gens must not cancel a mid-file restart.
|
||||
seekGen int64
|
||||
// done is closed when the ffmpeg goroutine fully exits (after process death).
|
||||
done chan struct{}
|
||||
}
|
||||
@@ -104,7 +107,7 @@ func (t *TranscoderService) PlaylistPath(mediaID string) string {
|
||||
// EnsureJob makes sure a transcode is running for mediaID from the start of
|
||||
// the source. Prefer EnsureJobFrom when the player seeks into the middle.
|
||||
func (t *TranscoderService) EnsureJob(ctx context.Context, mediaID string) (string, error) {
|
||||
return t.EnsureJobFrom(ctx, mediaID, 0)
|
||||
return t.EnsureJobFrom(ctx, mediaID, 0, 0)
|
||||
}
|
||||
|
||||
// EnsureJobFrom starts (or reuses) an HLS job that seeks the source to
|
||||
@@ -112,7 +115,12 @@ func (t *TranscoderService) EnsureJob(ctx context.Context, mediaID string) (stri
|
||||
// matches that offset; otherwise the previous job is cancelled and the HLS
|
||||
// cache dir is wiped so the player can jump without waiting for a full
|
||||
// head-to-tail transcode.
|
||||
func (t *TranscoderService) EnsureJobFrom(ctx context.Context, mediaID string, startSec float64) (string, error) {
|
||||
//
|
||||
// seekGen is the client `_seek` query (unix ms). A newer gen replaces an older
|
||||
// job; an older or missing gen must not clobber a mid-file restart — hls.js
|
||||
// in-flight playlist refreshes from a destroyed player commonly arrive as
|
||||
// start=0 right after a scrub and would otherwise reset playback to the head.
|
||||
func (t *TranscoderService) EnsureJobFrom(ctx context.Context, mediaID string, startSec float64, seekGen int64) (string, error) {
|
||||
if !t.cfg.Transcoder.Enabled {
|
||||
return "", ErrTranscodeDisabled
|
||||
}
|
||||
@@ -138,6 +146,11 @@ func (t *TranscoderService) EnsureJobFrom(ctx context.Context, mediaID string, s
|
||||
t.mu.Unlock()
|
||||
return t.PlaylistPath(mediaID), nil
|
||||
}
|
||||
if !shouldReplaceHLSJob(existing, startSec, seekGen) {
|
||||
t.touchJobLocked(mediaID)
|
||||
t.mu.Unlock()
|
||||
return t.PlaylistPath(mediaID), nil
|
||||
}
|
||||
prev := t.detachJobLocked(mediaID)
|
||||
t.mu.Unlock()
|
||||
waitJobExit(prev, 12*time.Second)
|
||||
@@ -164,13 +177,16 @@ func (t *TranscoderService) EnsureJobFrom(ctx context.Context, mediaID string, s
|
||||
|
||||
t.mu.Lock()
|
||||
if existing, ok := t.jobs[mediaID]; ok {
|
||||
// Another path (StopJob then a peer Ensure) should be rare under the gate;
|
||||
// still reuse or replace safely.
|
||||
if sameHLSStart(existing.startSec, startSec) {
|
||||
t.touchJobLocked(mediaID)
|
||||
t.mu.Unlock()
|
||||
return t.PlaylistPath(mediaID), nil
|
||||
}
|
||||
if !shouldReplaceHLSJob(existing, startSec, seekGen) {
|
||||
t.touchJobLocked(mediaID)
|
||||
t.mu.Unlock()
|
||||
return t.PlaylistPath(mediaID), nil
|
||||
}
|
||||
prev := t.detachJobLocked(mediaID)
|
||||
t.mu.Unlock()
|
||||
waitJobExit(prev, 12*time.Second)
|
||||
@@ -190,6 +206,7 @@ func (t *TranscoderService) EnsureJobFrom(ctx context.Context, mediaID string, s
|
||||
lastAccess: time.Now(),
|
||||
encoder: t.effectiveEncoder(),
|
||||
startSec: startSec,
|
||||
seekGen: seekGen,
|
||||
done: make(chan struct{}),
|
||||
}
|
||||
t.jobs[mediaID] = job
|
||||
@@ -203,6 +220,25 @@ func (t *TranscoderService) EnsureJobFrom(ctx context.Context, mediaID string, s
|
||||
return t.PlaylistPath(mediaID), nil
|
||||
}
|
||||
|
||||
// shouldReplaceHLSJob reports whether an incoming playlist request may cancel
|
||||
// the running job. Stale hls.js refreshes (older/missing `_seek`) must not win.
|
||||
func shouldReplaceHLSJob(existing *hlsJob, startSec float64, seekGen int64) bool {
|
||||
if existing == nil {
|
||||
return true
|
||||
}
|
||||
if sameHLSStart(existing.startSec, startSec) {
|
||||
return false
|
||||
}
|
||||
if seekGen > 0 && existing.seekGen > 0 && seekGen < existing.seekGen {
|
||||
return false
|
||||
}
|
||||
// Untagged request while a seek-tagged job is active: treat as stale.
|
||||
if seekGen == 0 && existing.seekGen > 0 {
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func (t *TranscoderService) mediaStartGate(mediaID string) *sync.Mutex {
|
||||
v, _ := t.startGates.LoadOrStore(mediaID, &sync.Mutex{})
|
||||
return v.(*sync.Mutex)
|
||||
|
||||
@@ -133,16 +133,13 @@ func appendInputAndVideoArgs(args []string, input transcodeInput, settings ffmpe
|
||||
if input.StartSec > 0.05 {
|
||||
ss = strconv.FormatFloat(input.StartSec, 'f', 3, 64)
|
||||
}
|
||||
// Local files: input -ss (byte/keyframe seek). HTTP/STRM: output -ss after
|
||||
// -i — several cloud/WMV demuxers ignore input seeks and would otherwise
|
||||
// restart from t=0. Realtime (-re) is already disabled for StartSec > 0.
|
||||
if ss != "" && !isHTTPSource(input.Source) {
|
||||
// Always use input -ss (before -i). Output -ss on HTTP/WMV decodes from
|
||||
// byte 0 up to the offset and cannot meet playlist WaitReady for deep
|
||||
// scrubbing; CDNs with Range support jump via demuxer seek instead.
|
||||
if ss != "" {
|
||||
args = append(args, "-ss", ss)
|
||||
}
|
||||
args = append(args, "-i", input.Source)
|
||||
if ss != "" && isHTTPSource(input.Source) {
|
||||
args = append(args, "-ss", ss)
|
||||
}
|
||||
args = append(args, "-map", "0:v:0?", "-map", "0:a:0?", "-vf", video.filter, "-c:v", video.codec)
|
||||
if settings.threads > 0 && video.codec == "libx264" {
|
||||
args = append(args, "-threads", strconv.Itoa(settings.threads))
|
||||
|
||||
@@ -232,17 +232,23 @@ func TestBuildFFmpegArgsHTTPSeekAfterDashI(t *testing.T) {
|
||||
if strings.Contains(joined, " -re ") {
|
||||
t.Fatalf("seek restart must disable -re, got: %s", joined)
|
||||
}
|
||||
idxSS, idxI := -1, -1
|
||||
idxSS, idxI, ssCount := -1, -1, 0
|
||||
for i, arg := range args {
|
||||
if arg == "-ss" {
|
||||
idxSS = i
|
||||
ssCount++
|
||||
if idxSS < 0 {
|
||||
idxSS = i
|
||||
}
|
||||
}
|
||||
if arg == "-i" && idxI < 0 {
|
||||
idxI = i
|
||||
}
|
||||
}
|
||||
if idxSS < 0 || idxI < 0 || idxSS < idxI {
|
||||
t.Fatalf("expected http -ss after -i, args=%v", args)
|
||||
if idxSS < 0 || idxI < 0 || idxSS > idxI {
|
||||
t.Fatalf("expected http -ss before -i, args=%v", args)
|
||||
}
|
||||
if ssCount != 1 {
|
||||
t.Fatalf("expected a single -ss, got %d in %v", ssCount, args)
|
||||
}
|
||||
if args[idxSS+1] != "90.000" {
|
||||
t.Fatalf("start = %q", args[idxSS+1])
|
||||
@@ -258,6 +264,25 @@ func TestSameHLSStart(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestShouldReplaceHLSJob(t *testing.T) {
|
||||
existing := &hlsJob{startSec: 120, seekGen: 1000}
|
||||
if shouldReplaceHLSJob(existing, 0, 0) {
|
||||
t.Fatal("untagged start=0 must not clobber seek-tagged job")
|
||||
}
|
||||
if shouldReplaceHLSJob(existing, 0, 900) {
|
||||
t.Fatal("older _seek must not clobber newer job")
|
||||
}
|
||||
if !shouldReplaceHLSJob(existing, 200, 1001) {
|
||||
t.Fatal("newer _seek should replace")
|
||||
}
|
||||
if !shouldReplaceHLSJob(&hlsJob{startSec: 0, seekGen: 0}, 120, 1000) {
|
||||
t.Fatal("seek should replace untagged head job")
|
||||
}
|
||||
if shouldReplaceHLSJob(existing, 120.2, 1001) {
|
||||
t.Fatal("same start should not replace")
|
||||
}
|
||||
}
|
||||
|
||||
func TestFilterHLSSegmentQueryDropsStart(t *testing.T) {
|
||||
got := filterHLSSegmentQuery("token=abc&start=120.5&profile_id=1")
|
||||
if strings.Contains(got, "start=") {
|
||||
|
||||
Reference in New Issue
Block a user