mirror of
https://github.com/truewhile/MeBox.git
synced 2026-09-28 11:16:37 +08:00
Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| c489374b61 |
@@ -0,0 +1,16 @@
|
||||
{
|
||||
"mcpServers": {
|
||||
"ssh": {
|
||||
"command": "cmd.exe",
|
||||
"args": [
|
||||
"/c",
|
||||
"npx",
|
||||
"-y",
|
||||
"@aiondadotcom/mcp-ssh"
|
||||
],
|
||||
"env": {
|
||||
"ProgramData": "C:\\ProgramData"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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}
|
||||
}
|
||||
@@ -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.
|
||||
}
|
||||
@@ -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{
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -242,13 +242,21 @@ export function PlayerPage() {
|
||||
}, [id, modeParam, directOnly])
|
||||
|
||||
// Wire up the actual <video> element when we know the mode.
|
||||
// Depend on media.id (not the media object): refreshing duration after
|
||||
// MANIFEST_PARSED must not remount HLS or it storms EnsureJob / DELETE.
|
||||
const mediaId = media?.id
|
||||
const mediaRef = useRef(media)
|
||||
mediaRef.current = media
|
||||
useEffect(() => {
|
||||
if (!media || !ref.current) return
|
||||
if (!mediaId || !ref.current) return
|
||||
const currentMedia = mediaRef.current
|
||||
if (!currentMedia) return
|
||||
teardownHls()
|
||||
|
||||
const video = ref.current
|
||||
const durationSec = currentMedia.duration_sec || 0
|
||||
if (mode === 'hls') {
|
||||
const url = hlsURL(media.id, hlsStartSec)
|
||||
const url = hlsURL(mediaId, hlsStartSec)
|
||||
void import('hls.js').then(({ default: HlsCtor }) => {
|
||||
if (HlsCtor.isSupported()) {
|
||||
const hls = new HlsCtor({ enableWorker: true, lowLatencyMode: false })
|
||||
@@ -256,9 +264,9 @@ export function PlayerPage() {
|
||||
hls.attachMedia(video)
|
||||
hls.on(HlsCtor.Events.MANIFEST_PARSED, () => {
|
||||
// .strm 入库时常缺 duration;转码启动时会补探测,这里刷新一次给进度条总时长。
|
||||
if ((media.duration_sec || 0) > 0) return
|
||||
if (durationSec > 0) return
|
||||
mediaAPI
|
||||
.get(media.id)
|
||||
.get(mediaId)
|
||||
.then((fresh) => {
|
||||
if ((fresh.duration_sec || 0) > 0) setMedia(fresh)
|
||||
})
|
||||
@@ -290,23 +298,22 @@ export function PlayerPage() {
|
||||
setMode('direct')
|
||||
})
|
||||
} else {
|
||||
video.src = streamURL(media.id)
|
||||
if (hlsUnavailable && needsTranscodeForBrowser(media)) {
|
||||
video.src = streamURL(mediaId)
|
||||
if (hlsUnavailable && needsTranscodeForBrowser(currentMedia)) {
|
||||
setPlayerError('当前正在直连播放原始文件;此封装或音轨浏览器兼容性有限,可能只有画面没有声音。请配置本机 ffmpeg 后切回 HLS 转码播放。')
|
||||
}
|
||||
void video.play().catch(() => undefined)
|
||||
}
|
||||
return () => teardownHls()
|
||||
}, [hlsUnavailable, hlsStartSec, media, mode, params, setParams, teardownHls])
|
||||
}, [hlsUnavailable, hlsStartSec, mediaId, mode, params, setParams, teardownHls])
|
||||
|
||||
// Stop the host ffmpeg job only when leaving HLS for this media (not on mid-file seek restarts).
|
||||
useEffect(() => {
|
||||
if (!media || mode !== 'hls') return
|
||||
const mediaId = media.id
|
||||
if (!mediaId || mode !== 'hls') return
|
||||
return () => {
|
||||
api.delete(`/hls/${encodeURIComponent(mediaId)}`).catch(() => undefined)
|
||||
}
|
||||
}, [media, mode])
|
||||
}, [mediaId, mode])
|
||||
|
||||
// 自动拉取已有的播放进度并恢复播放位置
|
||||
useEffect(() => {
|
||||
|
||||
Reference in New Issue
Block a user