Files
MeBox/internal/service/transcoder_jobs.go
T
truewhile c489374b61 bug处理
2026-09-07 23:45:59 +08:00

158 lines
3.8 KiB
Go

package service
import (
"context"
"os"
"time"
"go.uber.org/zap"
)
// WaitReady blocks (with a deadline) until the playlist file shows up on
// disk for the *current* job generation. Stale playlists left behind by a
// failed RemoveAll / still-exiting ffmpeg must not unblock a mid-file restart.
func (t *TranscoderService) WaitReady(ctx context.Context, mediaID string, timeout time.Duration) bool {
deadline := time.Now().Add(timeout)
for {
t.mu.Lock()
job, ok := t.jobs[mediaID]
var started time.Time
if ok {
started = job.startedAt
}
t.mu.Unlock()
if ok {
if info, err := os.Stat(t.PlaylistPath(mediaID)); err == nil {
// Allow a small clock skew; reject anything older than this job.
if !info.ModTime().Before(started.Add(-2 * time.Second)) {
t.mu.Lock()
if j, exists := t.jobs[mediaID]; exists {
j.playlistOK = true
}
t.mu.Unlock()
return true
}
}
}
if time.Now().After(deadline) || ctx.Err() != nil {
return false
}
select {
case <-ctx.Done():
return false
case <-time.After(250 * time.Millisecond):
}
}
}
// 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()
prev := t.detachJobLocked(mediaID)
t.mu.Unlock()
waitJobExit(prev, 12*time.Second)
}
// TouchJob records client activity for the HLS playlist or segment. The idle
// watchdog uses it to stop ffmpeg soon after the player is closed or switches
// back to direct play.
func (t *TranscoderService) TouchJob(mediaID string) {
t.mu.Lock()
defer t.mu.Unlock()
t.touchJobLocked(mediaID)
}
func (t *TranscoderService) touchJobLocked(mediaID string) {
if j, ok := t.jobs[mediaID]; ok {
j.lastAccess = time.Now()
}
}
// StopAll terminates every running transcode (called on graceful shutdown).
func (t *TranscoderService) StopAll() {
t.mu.Lock()
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)
}
}
// ActiveJob is the JSON shape exposed to the React Tasks panel.
type ActiveJob struct {
MediaID string `json:"media_id"`
Encoder string `json:"encoder"`
StartedAt time.Time `json:"started_at"`
PlaylistOK bool `json:"playlist_ok"`
}
// Active returns a snapshot of the currently running transcode jobs.
func (t *TranscoderService) Active() []ActiveJob {
t.mu.Lock()
defer t.mu.Unlock()
out := make([]ActiveJob, 0, len(t.jobs))
for _, j := range t.jobs {
out = append(out, ActiveJob{
MediaID: j.mediaID,
Encoder: j.encoder,
StartedAt: j.startedAt,
PlaylistOK: j.playlistOK,
})
}
return out
}
func (t *TranscoderService) maxConcurrent() int {
if t.cfg.Transcoder.MaxConcurrent <= 0 {
return 1
}
return t.cfg.Transcoder.MaxConcurrent
}
func (t *TranscoderService) idleTimeout() time.Duration {
if t.cfg.Transcoder.IdleTimeoutSeconds <= 0 {
return 120 * time.Second
}
return time.Duration(t.cfg.Transcoder.IdleTimeoutSeconds) * time.Second
}
func (t *TranscoderService) monitorIdle(ctx context.Context, job *hlsJob) {
timeout := t.idleTimeout()
ticker := time.NewTicker(15 * time.Second)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
t.mu.Lock()
current, ok := t.jobs[job.mediaID]
if !ok {
t.mu.Unlock()
return
}
idleFor := time.Since(current.lastAccess)
t.mu.Unlock()
if idleFor >= timeout {
t.log.Info("transcode idle timeout",
zap.String("media_id", job.mediaID),
zap.Duration("idle_for", idleFor),
zap.Duration("timeout", timeout),
)
t.StopJob(job.mediaID)
return
}
}
}
}