From c489374b61aa2121d49769cfec3b21a088d12453 Mon Sep 17 00:00:00 2001 From: truewhile <62226914+truewhile@users.noreply.github.com> Date: Mon, 7 Sep 2026 23:45:59 +0800 Subject: [PATCH] =?UTF-8?q?bug=E5=A4=84=E7=90=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .cursor/mcp.json | 16 ++++++ internal/service/transcoder.go | 61 +++++++++++++++++++-- internal/service/transcoder_cmd_unix.go | 14 +++++ internal/service/transcoder_cmd_windows.go | 9 +++ internal/service/transcoder_ffmpeg_probe.go | 14 +++-- internal/service/transcoder_jobs.go | 27 +++++---- web/src/pages/PlayerPage.tsx | 27 +++++---- 7 files changed, 138 insertions(+), 30 deletions(-) create mode 100644 .cursor/mcp.json create mode 100644 internal/service/transcoder_cmd_unix.go create mode 100644 internal/service/transcoder_cmd_windows.go diff --git a/.cursor/mcp.json b/.cursor/mcp.json new file mode 100644 index 0000000..4075ad5 --- /dev/null +++ b/.cursor/mcp.json @@ -0,0 +1,16 @@ +{ + "mcpServers": { + "ssh": { + "command": "cmd.exe", + "args": [ + "/c", + "npx", + "-y", + "@aiondadotcom/mcp-ssh" + ], + "env": { + "ProgramData": "C:\\ProgramData" + } + } + } +} diff --git a/internal/service/transcoder.go b/internal/service/transcoder.go index f22f75c..eec72a1 100644 --- a/internal/service/transcoder.go +++ b/internal/service/transcoder.go @@ -49,6 +49,9 @@ type TranscoderService struct { mu sync.Mutex jobs map[string]*hlsJob + // startGates serializes EnsureJobFrom / StopJob per media so concurrent + // playlist hits cannot spawn multiple ffmpeg writers into one HLS dir. + startGates sync.Map // mediaID -> *sync.Mutex strmResolve func(ctx context.Context, raw string) (*StrmPlayResult, error) probe *FFprobeService } @@ -65,6 +68,8 @@ 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 + // done is closed when the ffmpeg goroutine fully exits (after process death). + done chan struct{} } var ( @@ -122,6 +127,10 @@ func (t *TranscoderService) EnsureJobFrom(ctx context.Context, mediaID string, s return "", ErrMediaNotFound } + gate := t.mediaStartGate(mediaID) + gate.Lock() + defer gate.Unlock() + t.mu.Lock() if existing, ok := t.jobs[mediaID]; ok { if sameHLSStart(existing.startSec, startSec) { @@ -129,10 +138,12 @@ func (t *TranscoderService) EnsureJobFrom(ctx context.Context, mediaID string, s t.mu.Unlock() return t.PlaylistPath(mediaID), nil } - existing.cancel() - delete(t.jobs, mediaID) + prev := t.detachJobLocked(mediaID) + t.mu.Unlock() + waitJobExit(prev, 12*time.Second) + } else { + t.mu.Unlock() } - t.mu.Unlock() input, err := t.resolveTranscodeInput(ctx, m) if err != nil { @@ -146,19 +157,24 @@ func (t *TranscoderService) EnsureJobFrom(ctx context.Context, mediaID string, s outDir := t.HLSDir(mediaID) // Wipe prior segments so a mid-file restart cannot serve stale early chunks. + // Only safe after the previous ffmpeg has exited (waited above / via gate). if err := resetHLSDir(outDir); err != nil { return "", err } 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 } - existing.cancel() - delete(t.jobs, mediaID) + prev := t.detachJobLocked(mediaID) + t.mu.Unlock() + waitJobExit(prev, 12*time.Second) + t.mu.Lock() } if max := t.maxConcurrent(); max > 0 && len(t.jobs) >= max { t.mu.Unlock() @@ -174,15 +190,48 @@ func (t *TranscoderService) EnsureJobFrom(ctx context.Context, mediaID string, s lastAccess: time.Now(), encoder: t.effectiveEncoder(), startSec: startSec, + done: make(chan struct{}), } t.jobs[mediaID] = job t.mu.Unlock() helper.Go(t.log, "transcoder.monitorIdle", func() { t.monitorIdle(jobCtx, job) }) - helper.Go(t.log, "transcoder.ffmpeg", func() { t.runFFmpeg(jobCtx, job, input) }) + helper.Go(t.log, "transcoder.ffmpeg", func() { + defer close(job.done) + t.runFFmpeg(jobCtx, job, input) + }) return t.PlaylistPath(mediaID), nil } +func (t *TranscoderService) mediaStartGate(mediaID string) *sync.Mutex { + v, _ := t.startGates.LoadOrStore(mediaID, &sync.Mutex{}) + return v.(*sync.Mutex) +} + +// detachJobLocked cancels and removes a job from the map without waiting. +// Caller must hold t.mu and must waitJobExit afterwards before wiping the HLS dir. +func (t *TranscoderService) detachJobLocked(mediaID string) *hlsJob { + j, ok := t.jobs[mediaID] + if !ok { + return nil + } + j.cancel() + delete(t.jobs, mediaID) + return j +} + +func waitJobExit(job *hlsJob, timeout time.Duration) { + if job == nil || job.done == nil { + return + } + timer := time.NewTimer(timeout) + defer timer.Stop() + select { + case <-job.done: + case <-timer.C: + } +} + func sameHLSStart(a, b float64) bool { const tol = 0.75 if a < b { diff --git a/internal/service/transcoder_cmd_unix.go b/internal/service/transcoder_cmd_unix.go new file mode 100644 index 0000000..ac040a1 --- /dev/null +++ b/internal/service/transcoder_cmd_unix.go @@ -0,0 +1,14 @@ +//go:build unix + +package service + +import ( + "os/exec" + "syscall" +) + +// setFFmpegSysProcAttr puts ffmpeg in its own process group so cancel can +// tear down the whole group (and any helpers) reliably on Linux/macOS. +func setFFmpegSysProcAttr(cmd *exec.Cmd) { + cmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true} +} diff --git a/internal/service/transcoder_cmd_windows.go b/internal/service/transcoder_cmd_windows.go new file mode 100644 index 0000000..a5b303d --- /dev/null +++ b/internal/service/transcoder_cmd_windows.go @@ -0,0 +1,9 @@ +//go:build windows + +package service + +import "os/exec" + +func setFFmpegSysProcAttr(cmd *exec.Cmd) { + // Windows: CommandContext cancel is enough for the single ffmpeg process. +} diff --git a/internal/service/transcoder_ffmpeg_probe.go b/internal/service/transcoder_ffmpeg_probe.go index 57c592a..dee811b 100644 --- a/internal/service/transcoder_ffmpeg_probe.go +++ b/internal/service/transcoder_ffmpeg_probe.go @@ -34,6 +34,7 @@ func (t *TranscoderService) runFFmpeg(ctx context.Context, job *hlsJob, input tr args := buildFFmpegArgsForInput(t.cfg, input, playlist, segments) cmd := exec.CommandContext(ctx, bin, args...) // #nosec G204 -- bin is resolved by resolveFFmpegPath and args are passed without a shell. + setFFmpegSysProcAttr(cmd) cmd.Stderr = os.Stderr t.log.Info("transcode started", @@ -43,20 +44,25 @@ func (t *TranscoderService) runFFmpeg(ctx context.Context, job *hlsJob, input tr zap.Float64("start_sec", input.StartSec), ) t.hub.Publish("transcode", map[string]any{ - "media_id": job.mediaID, - "encoder": job.encoder, - "status": "started", + "media_id": job.mediaID, + "encoder": job.encoder, + "status": "started", + "start_sec": input.StartSec, }) if err := cmd.Run(); err != nil && !errors.Is(ctx.Err(), context.Canceled) { t.log.Warn("ffmpeg exited", zap.String("media_id", job.mediaID), + zap.Float64("start_sec", input.StartSec), zap.Error(err), ) } t.mu.Lock() - delete(t.jobs, job.mediaID) + // Only drop the map entry if we are still the registered generation. + if cur, ok := t.jobs[job.mediaID]; ok && cur == job { + delete(t.jobs, job.mediaID) + } t.mu.Unlock() t.hub.Publish("transcode", map[string]any{ diff --git a/internal/service/transcoder_jobs.go b/internal/service/transcoder_jobs.go index 375dd82..daedb3e 100644 --- a/internal/service/transcoder_jobs.go +++ b/internal/service/transcoder_jobs.go @@ -45,14 +45,18 @@ func (t *TranscoderService) WaitReady(ctx context.Context, mediaID string, timeo } } -// StopJob cancels a running ffmpeg process for mediaID, if any. +// StopJob cancels a running ffmpeg process for mediaID, if any, and waits +// briefly for it to exit so a subsequent EnsureJobFrom cannot race-write the +// same HLS directory. func (t *TranscoderService) StopJob(mediaID string) { + gate := t.mediaStartGate(mediaID) + gate.Lock() + defer gate.Unlock() + t.mu.Lock() - defer t.mu.Unlock() - if j, ok := t.jobs[mediaID]; ok { - j.cancel() - delete(t.jobs, mediaID) - } + prev := t.detachJobLocked(mediaID) + t.mu.Unlock() + waitJobExit(prev, 12*time.Second) } // TouchJob records client activity for the HLS playlist or segment. The idle @@ -73,10 +77,13 @@ func (t *TranscoderService) touchJobLocked(mediaID string) { // StopAll terminates every running transcode (called on graceful shutdown). func (t *TranscoderService) StopAll() { t.mu.Lock() - defer t.mu.Unlock() - for id, j := range t.jobs { - j.cancel() - delete(t.jobs, id) + pending := make([]*hlsJob, 0, len(t.jobs)) + for id := range t.jobs { + pending = append(pending, t.detachJobLocked(id)) + } + t.mu.Unlock() + for _, j := range pending { + waitJobExit(j, 5*time.Second) } } diff --git a/web/src/pages/PlayerPage.tsx b/web/src/pages/PlayerPage.tsx index cae0302..2153a7e 100644 --- a/web/src/pages/PlayerPage.tsx +++ b/web/src/pages/PlayerPage.tsx @@ -242,13 +242,21 @@ export function PlayerPage() { }, [id, modeParam, directOnly]) // Wire up the actual