package service import ( "context" "crypto/sha256" "encoding/hex" "errors" "fmt" "os" "strings" "sync" "time" "go.uber.org/zap" "github.com/truewhile/MeBox/internal/helper" "github.com/truewhile/MeBox/internal/model" "github.com/truewhile/MeBox/internal/repository" ) // 探测结果的缓存与预算策略。 const ( // mediaProbeTimeout 是一次后台探测的总预算(含把 strm 目标解析成直链)。 mediaProbeTimeout = 90 * time.Second // mediaProbeFailureRetry 是探测失败后允许重新探测的间隔。失败的结果也会落库, // 否则每次播放都会为一个坏源重跑一次。 mediaProbeFailureRetry = 6 * time.Hour // mediaProbeErrorLimit 限制落库的错误信息长度(列宽 512 字节,错误里可能 // 带 URL,截断同时避免超长)。 mediaProbeErrorLimit = 200 ) // mediaProber 是 MediaProbeService 需要的探测能力。抽成接口是为了在测试里注入 // 桩,避免依赖真实 ffprobe 二进制。 type mediaProber interface { ProbeFull(ctx context.Context, input ProbeInput) (*FullProbeResult, error) } // MediaProbeService 用 ffprobe 提取媒体的基础信息(容器、每路轨道、内嵌章节) // 并落库缓存。 // // 定位是「播放时顺带补齐媒体信息」:STRM / 云盘媒体在扫描阶段拿不到时长,而 // 播放链路要用缓存里的时长来换算「延续到片尾」这类区间,详情页将来也直接读这份 // 媒体信息。 // // 它**不参与片头/片尾判定**:章节标题绝大多数没有语义(生产库实测抽样 64 个 // 文件,命中 0 个),拿它去猜跳过点只会给出错误的位置,时间轴数据仍然只信 // TheIntroDB。 // // 核心约束:一次探测要 2~4.5 秒(远端直链要跨洋跑几次 HTTP 事务),所以只允许 // 异步跑,播放链路永远只读缓存。 type MediaProbeService struct { log *zap.Logger repo *repository.Container probe mediaProber // resolve 把 strm 播放目标解析成最终直链(含绑定 UA 的请求头)。与转码、 // 内嵌字幕发现走同一条换链路径,否则会踩到网盘 CDN 的防盗链 403。 resolve func(ctx context.Context, raw, userAgent string) (*StrmPlayResult, error) mu sync.Mutex inFlight map[string]struct{} } // NewMediaProbeService is the constructor. func NewMediaProbeService(log *zap.Logger, repo *repository.Container, probe *FFprobeService) *MediaProbeService { svc := &MediaProbeService{ log: log, repo: repo, inFlight: make(map[string]struct{}), } if probe != nil { svc.probe = probe } return svc } // SetPlayTargetResolver injects the strm → direct-link resolver. func (s *MediaProbeService) SetPlayTargetResolver(resolve func(ctx context.Context, raw, userAgent string) (*StrmPlayResult, error)) *MediaProbeService { if s != nil { s.resolve = resolve } return s } // EnsureAsync 保证这部媒体的探测已排上队,并立刻返回。 // // 已经有同一条媒体的探测在跑、或服务未配置好时返回 false;这次新排上一条返回 // true——调用方据此告诉客户端「稍后再拉一次」。 func (s *MediaProbeService) EnsureAsync(m *model.Media) bool { if s == nil || s.probe == nil || s.repo == nil || m == nil || strings.TrimSpace(m.ID) == "" { return false } if !s.reserve(m.ID) { return false } // 复制一份媒体行:调用方的对象可能属于请求作用域,后台协程不该继续引用它。 snapshot := *m helper.Go(s.log, "service.mediaProbe", func() { defer s.release(snapshot.ID) s.run(context.Background(), &snapshot) }) return true } func (s *MediaProbeService) reserve(mediaID string) bool { s.mu.Lock() defer s.mu.Unlock() if s.inFlight == nil { s.inFlight = make(map[string]struct{}) } if _, running := s.inFlight[mediaID]; running { return false } s.inFlight[mediaID] = struct{}{} return true } func (s *MediaProbeService) release(mediaID string) { s.mu.Lock() delete(s.inFlight, mediaID) s.mu.Unlock() } // run 执行一次探测并落库。它跑在后台,没有调用方能接收错误,所以任何失败都只 // 记录、不外抛。 func (s *MediaProbeService) run(ctx context.Context, m *model.Media) { ctx, cancel := context.WithTimeout(ctx, mediaProbeTimeout) defer cancel() input, err := s.probeInput(ctx, m) if err != nil { s.markFailure(ctx, m.ID, err) return } result, err := s.probe.ProbeFull(ctx, input) if err != nil { s.markFailure(ctx, m.ID, err) return } if err := s.persistProbe(ctx, m, result); err != nil { s.markFailure(ctx, m.ID, err) } } // probeInput 把媒体行解析成 ffprobe 能直接打开的输入。 // // 本地文件给路径;STRM / 云盘先解析成最终直链并带上绑定的请求头——直链与 UA // 必须配套,用错会被 CDN 拒绝。 func (s *MediaProbeService) probeInput(ctx context.Context, m *model.Media) (ProbeInput, error) { if m == nil { return ProbeInput{}, ErrMediaNotFound } if !isStrmMediaRow(m) { if _, err := os.Stat(m.Path); err != nil { return ProbeInput{}, ErrMediaNotFound } return ProbeInput{Source: m.Path}, nil } raw := strings.TrimSpace(m.STRMURL) if raw == "" && strings.HasSuffix(strings.ToLower(strings.TrimSpace(m.Path)), ".strm") { parsed, err := readLocalSTRMTarget(m.Path) if err == nil { raw = strings.TrimSpace(parsed) } } if raw == "" { return ProbeInput{}, errors.New("strm play target missing") } if s.resolve != nil { resolved, err := s.resolve(ctx, raw, "") if err != nil { return ProbeInput{}, err } in, err := transcodeInputFromPlayResult(resolved) if err != nil { return ProbeInput{}, err } return ProbeInput{Source: in.Source, Headers: in.Headers}, nil } if isHTTPPlaybackTarget(raw) { return ProbeInput{Source: raw}, nil } return ProbeInput{}, errors.New("strm probe source unavailable") } // persistProbe 把一次成功的探测落库。 func (s *MediaProbeService) persistProbe(ctx context.Context, m *model.Media, result *FullProbeResult) error { if s.repo == nil || s.repo.MediaProbe == nil { return errors.New("media probe repository not wired") } payload, err := result.PayloadJSON() if err != nil { return err } row := &model.MediaProbe{ MediaID: m.ID, Signature: mediaProbeSignature(m), Source: mediaProbeInputKind(m), Container: result.Container, DurationSec: result.DurationSec, BitRate: result.BitRate, Width: firstStreamDimension(result, "video", true), Height: firstStreamDimension(result, "video", false), VideoCodec: firstStreamCodec(result, "video"), AudioCodec: firstStreamCodec(result, "audio"), VideoStreams: len(result.StreamsOfType("video")), AudioStreams: len(result.StreamsOfType("audio")), SubtitleStreams: len(result.StreamsOfType("subtitle")), ChapterCount: len(result.Chapters), Payload: payload, ProbedAt: time.Now(), } if err := s.repo.MediaProbe.Upsert(ctx, row); err != nil { return err } s.backfillDuration(ctx, m, result.DurationSec) return nil } // backfillDuration 把探测到的时长补进 media.duration_sec。STRM / 云盘媒体在扫描 // 阶段拿不到时长,而末段区间(end_ms = 0)要靠它才能换算出真实结束时间。 func (s *MediaProbeService) backfillDuration(ctx context.Context, m *model.Media, durationSec int) { if s.repo == nil || s.repo.DB == nil || m == nil || durationSec <= 0 || m.DurationSec > 0 { return } err := s.repo.DB.WithContext(ctx).Model(&model.Media{}). Where("id = ? AND duration_sec <= 0", m.ID). Update("duration_sec", durationSec).Error if err != nil && s.log != nil { s.log.Debug("backfill probed duration failed", zap.String("media_id", m.ID), zap.Error(err)) } } func (s *MediaProbeService) markFailure(ctx context.Context, mediaID string, probeErr error) { if s == nil || s.repo == nil || probeErr == nil { return } message := truncateProbeError(probeErr) if err := s.repo.MediaProbe.MarkFailure(ctx, mediaID, message, time.Now()); err != nil && s.log != nil { s.log.Debug("record media probe failure failed", zap.String("media_id", mediaID), zap.Error(err)) } if s.log != nil { s.log.Debug("media probe failed", zap.String("media_id", mediaID), zap.Error(probeErr)) } } // truncateProbeError 限制错误信息长度。错误里可能带被拒绝的直链,落库时截断, // 避免超长并减少敏感内容。 func truncateProbeError(err error) string { message := strings.TrimSpace(err.Error()) runes := []rune(message) if len(runes) > mediaProbeErrorLimit { return string(runes[:mediaProbeErrorLimit]) } return message } // mediaProbeInputKind 记录输入形态,供详情页判断「这个时长是本地读的还是远端读的」。 func mediaProbeInputKind(m *model.Media) string { if m == nil { return "" } if isStrmMediaRow(m) { return "strm" } return "local" } // mediaProbeSignature 是「探的是哪个文件」的指纹。 // // 本地文件用路径 + 大小 + 修改时间;STRM / 云盘没有本地文件,用固化的播放目标, // 且绝不能用解析后的直链(每次签名都不同,缓存会永远失效)。整体做哈希: // 既固定长度,也不把路径或 pickcode 再抄一份进数据库。 func mediaProbeSignature(m *model.Media) string { if m == nil { return "" } var base string if isStrmMediaRow(m) { target := strings.TrimSpace(m.STRMURL) if target == "" { target = strings.TrimSpace(m.Path) } base = "strm|" + target } else if info, err := os.Stat(m.Path); err == nil { base = fmt.Sprintf("local|%s|%d|%d", m.Path, info.Size(), info.ModTime().Unix()) } else { base = "local|" + m.Path } sum := sha256.Sum256([]byte(base)) return hex.EncodeToString(sum[:16]) } func firstStreamCodec(result *FullProbeResult, kind string) string { if result == nil { return "" } for _, stream := range result.Streams { if stream.Type == kind && stream.Codec != "" { return stream.Codec } } return "" } // firstStreamDimension 取第一路指定类型轨道的宽(width=true)或高。 func firstStreamDimension(result *FullProbeResult, kind string, width bool) int { if result == nil { return 0 } for _, stream := range result.Streams { if stream.Type != kind { continue } if width { return stream.Width } return stream.Height } return 0 }