feat: ffprobe + TMDb scrape + HLS transcode + history/favourites/playlists

Backend
  - service/ffprobe.go: thin ffprobe wrapper, parses duration / resolution /
    codecs into a typed ProbeResult. 30s per-file timeout.
  - service/tmdb.go: minimal TMDb provider (search/movie). Disabled when no
    api key; supports tmdb_api_proxy / tmdb_image_proxy overrides for users
    behind a firewall.
  - service/scraper.go: filename cleaner (handles bracketed tags, scene
    noise tokens, year extraction), per-row + per-library enrichment with a
    4 RPS throttle and WS hub progress events. Unit-tested.
  - service/scanner.go: now invokes ffprobe per file and kicks the TMDb
    scraper in the background once a library scan finishes.
  - service/transcoder.go: per-media ffmpeg HLS job manager; outputs
    index.m3u8 + seg_NNNNN.ts under cache/hls/<id>; cancels jobs on
    shutdown; publishes 'transcode' WS events.
  - service/stream.go: serves HLS playlist (with 30s wait-for-ready) and
    .ts segments with path-traversal protection. Adds Probe() helper used
    by the admin 'reprobe' button.
  - service/image_proxy.go: cached, host-allow-listed reverse proxy for
    TMDb / Bangumi / Douban / Fanart / TheTVDB images so the SPA never
    hits a CORS or GFW issue.
  - service/playback.go: history upsert, favourites toggle, playlist CRUD
    + ordered items. RecentHistory joins with model.Media in one extra
    query so the home page can render a 'Continue Watching' row.
  - handler/streaming.go + handler/playback.go: REST endpoints for HLS,
    image proxy, scrape (one + library), reprobe, history, favourites,
    playlists.
  - handler/handler.go: registers /api/hls/:id/{index.m3u8,:seg}, /api/img,
    /api/history, /api/favourites/:id, /api/playlists/* with proper
    auth/admin guards.

Frontend
  - api/client.ts: imageURL() helper; hlsURL() endpoint; reuses the JWT in
    a query parameter for <video src> and <img src>.
  - api/playback.ts: typed helpers for history, favourites, playlists.
  - components/MediaCard.tsx: optional 'progress' prop renders a thin
    bottom progress bar, used by the new Continue Watching row.
  - pages/HomePage.tsx: two rows (Continue Watching + Recently Added);
    falls back to the empty-state hint when both are empty.
  - pages/PlayerPage.tsx: hls.js (lazy-imported) with auto-fallback to
    direct play; ?mode=hls|direct query toggle; resume position written
    every 10s while playing.
  - pages/MediaDetailPage.tsx: heart toggle + admin 'rescrape' / 'reprobe'
    buttons + dedicated 'HLS 转码播放' CTA.
  - pages/FavouritesPage.tsx, PlaylistsPage.tsx, PlaylistDetailPage.tsx:
    new screens.
  - components/Layout.tsx + App.tsx: sidebar links for Favourites and
    Playlists; routes are now lazily code-split via React.lazy + Suspense
    so the initial bundle stays at ~243 KB / 82 KB gzipped (hls.js is
    fetched only on first HLS playback).

Verified: go build, go vet, go test (incl. CleanQuery cases) all pass;
frontend tsc -b && vite build emits 9 route chunks plus a deferred hls
chunk.
This commit is contained in:
Kiro
2026-05-14 15:44:31 +00:00
parent d5cf5fb4b2
commit 0f30c34463
26 changed files with 2205 additions and 139 deletions
+24
View File
@@ -38,12 +38,36 @@ func Register(r *gin.Engine, cfg *config.Config, log *zap.Logger, svc *service.C
authed.POST("/libraries", middleware.AdminRequired(), createLibraryHandler(svc))
authed.DELETE("/libraries/:id", middleware.AdminRequired(), deleteLibraryHandler(svc))
authed.POST("/libraries/:id/scan", middleware.AdminRequired(), scanLibraryHandler(svc))
authed.POST("/libraries/:id/scrape", middleware.AdminRequired(), scrapeLibraryHandler(svc))
authed.GET("/libraries/:id/media", listMediaHandler(svc))
authed.GET("/media/:id", getMediaHandler(svc))
authed.GET("/media", searchMediaHandler(svc))
authed.POST("/media/:id/scrape", middleware.AdminRequired(), scrapeOneHandler(svc))
authed.POST("/media/:id/probe", middleware.AdminRequired(), reprobeHandler(svc))
// Streaming.
authed.GET("/stream/:id", streamHandler(svc))
authed.GET("/hls/:id/index.m3u8", hlsPlaylistHandler(svc))
authed.GET("/hls/:id/:seg", hlsSegmentHandler(svc))
authed.DELETE("/hls/:id", stopTranscodeHandler(svc))
// Image proxy (URL passed as ?url=...).
authed.GET("/img", imageProxyHandler(svc))
// History / favourites / playlists.
authed.GET("/history", recentHistoryHandler(svc))
authed.POST("/history", recordProgressHandler(svc))
authed.GET("/favourites", listFavouritesHandler(svc))
authed.POST("/favourites/:id", toggleFavouriteHandler(svc))
authed.GET("/playlists", listPlaylistsHandler(svc))
authed.POST("/playlists", createPlaylistHandler(svc))
authed.GET("/playlists/:id", getPlaylistHandler(svc))
authed.POST("/playlists/:id/items", addPlaylistItemHandler(svc))
authed.DELETE("/playlists/:id/items/:media_id", removePlaylistItemHandler(svc))
authed.DELETE("/playlists/:id", deletePlaylistHandler(svc))
authed.GET("/ws", wsHandler(svc))
}
+177
View File
@@ -0,0 +1,177 @@
// Package handler — playback history / favourites / playlists endpoints.
package handler
import (
"net/http"
"github.com/gin-gonic/gin"
"github.com/ShukeBta/MediaStationGo/internal/middleware"
"github.com/ShukeBta/MediaStationGo/internal/service"
)
// ─── History ────────────────────────────────────────────────────────────────
type progressReq struct {
MediaID string `json:"media_id" binding:"required"`
PositionMs int64 `json:"position_ms"`
DurationMs int64 `json:"duration_ms"`
}
func recordProgressHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
var req progressReq
if err := c.ShouldBindJSON(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
uid, _ := c.Get(middleware.CtxUserID)
if err := svc.Playback.RecordProgress(
c.Request.Context(), uid.(string), req.MediaID, req.PositionMs, req.DurationMs,
); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.Status(http.StatusNoContent)
}
}
func recentHistoryHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
uid, _ := c.Get(middleware.CtxUserID)
items, err := svc.Playback.RecentHistory(c.Request.Context(), uid.(string), 30)
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"items": items})
}
}
// ─── Favourites ─────────────────────────────────────────────────────────────
func toggleFavouriteHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
uid, _ := c.Get(middleware.CtxUserID)
state, err := svc.Playback.ToggleFavourite(
c.Request.Context(), uid.(string), c.Param("id"),
)
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"favourite": state})
}
}
func listFavouritesHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
uid, _ := c.Get(middleware.CtxUserID)
items, err := svc.Playback.ListFavourites(c.Request.Context(), uid.(string))
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"items": items})
}
}
// ─── Playlists ──────────────────────────────────────────────────────────────
type createPlaylistReq struct {
Name string `json:"name" binding:"required"`
IsPublic bool `json:"is_public"`
}
func createPlaylistHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
var req createPlaylistReq
if err := c.ShouldBindJSON(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
uid, _ := c.Get(middleware.CtxUserID)
pl, err := svc.Playback.CreatePlaylist(
c.Request.Context(), uid.(string), req.Name, req.IsPublic,
)
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, pl)
}
}
func listPlaylistsHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
uid, _ := c.Get(middleware.CtxUserID)
items, err := svc.Playback.ListPlaylists(c.Request.Context(), uid.(string))
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"items": items})
}
}
func getPlaylistHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
detail, err := svc.Playback.GetPlaylist(c.Request.Context(), c.Param("id"))
if err != nil {
c.JSON(http.StatusNotFound, gin.H{"error": err.Error()})
return
}
uid, _ := c.Get(middleware.CtxUserID)
role, _ := c.Get(middleware.CtxUserRole)
if !detail.Playlist.IsPublic && detail.Playlist.UserID != uid.(string) && role != "admin" {
c.JSON(http.StatusForbidden, gin.H{"error": "forbidden"})
return
}
c.JSON(http.StatusOK, detail)
}
}
type playlistItemReq struct {
MediaID string `json:"media_id" binding:"required"`
}
func addPlaylistItemHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
var req playlistItemReq
if err := c.ShouldBindJSON(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
if err := svc.Playback.AddToPlaylist(
c.Request.Context(), c.Param("id"), req.MediaID,
); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.Status(http.StatusNoContent)
}
}
func removePlaylistItemHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
if err := svc.Playback.RemoveFromPlaylist(
c.Request.Context(), c.Param("id"), c.Param("media_id"),
); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.Status(http.StatusNoContent)
}
}
func deletePlaylistHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
if err := svc.Playback.DeletePlaylist(
c.Request.Context(), c.Param("id"),
); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.Status(http.StatusNoContent)
}
}
+104
View File
@@ -0,0 +1,104 @@
// Package handler — HLS / image-proxy / scrape endpoints.
package handler
import (
"errors"
"net/http"
"github.com/gin-gonic/gin"
"github.com/ShukeBta/MediaStationGo/internal/service"
)
func hlsPlaylistHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
err := svc.Stream.ServeHLSPlaylist(c.Writer, c.Request, c.Param("id"))
if errors.Is(err, service.ErrMediaNotFound) {
c.JSON(http.StatusNotFound, gin.H{"error": "not found"})
return
}
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
}
}
func hlsSegmentHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
err := svc.Stream.ServeHLSSegment(c.Writer, c.Request, c.Param("id"), c.Param("seg"))
if err != nil {
c.JSON(http.StatusNotFound, gin.H{"error": err.Error()})
return
}
}
}
func stopTranscodeHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
svc.Transcoder.StopJob(c.Param("id"))
c.Status(http.StatusNoContent)
}
}
func imageProxyHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
raw := c.Query("url")
if err := svc.ImageProxy.Serve(c.Request.Context(), c.Writer, raw); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
}
}
// scrapeOneHandler enriches a single media via TMDb. Admin-only.
func scrapeOneHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
m, err := svc.Repo.Media.FindByID(c.Request.Context(), c.Param("id"))
if err != nil || m == nil {
c.JSON(http.StatusNotFound, gin.H{"error": "not found"})
return
}
if !svc.TMDb.Enabled() {
c.JSON(http.StatusPreconditionFailed, gin.H{"error": "tmdb api key not configured"})
return
}
if err := svc.Scraper.EnrichOne(c.Request.Context(), m); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
refreshed, _ := svc.Repo.Media.FindByID(c.Request.Context(), m.ID)
c.JSON(http.StatusOK, refreshed)
}
}
// scrapeLibraryHandler enriches every pending media in a library. Admin-only.
func scrapeLibraryHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
if !svc.TMDb.Enabled() {
c.JSON(http.StatusPreconditionFailed, gin.H{"error": "tmdb api key not configured"})
return
}
// Run in the background so HTTP returns instantly; the WS hub
// pushes per-item progress on the "scrape" topic.
go func(libID string) {
_, _ = svc.Scraper.EnrichLibrary(c.Copy().Request.Context(), libID)
}(c.Param("id"))
c.JSON(http.StatusAccepted, gin.H{"status": "scraping"})
}
}
// reprobeHandler re-runs ffprobe against a single media. Admin-only.
func reprobeHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
if err := svc.Stream.Probe(c.Request.Context(), c.Param("id"), svc.FFprobe); err != nil {
if errors.Is(err, service.ErrMediaNotFound) {
c.JSON(http.StatusNotFound, gin.H{"error": "not found"})
return
}
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.Status(http.StatusNoContent)
}
}
+110
View File
@@ -0,0 +1,110 @@
// Package service — ffprobe wrapper.
//
// FFprobeService shells out to the `ffprobe` binary configured in
// app.ffprobe_path and parses its JSON output into a typed struct. It is
// intentionally minimal: we only extract the fields needed to populate
// model.Media (duration, resolution, video / audio codec) so a fresh scan
// can show meaningful metadata even before the TMDb scraper has run.
package service
import (
"context"
"encoding/json"
"errors"
"fmt"
"os/exec"
"strconv"
"time"
"go.uber.org/zap"
"github.com/ShukeBta/MediaStationGo/internal/config"
)
// FFprobeService wraps the external ffprobe binary.
type FFprobeService struct {
cfg *config.Config
log *zap.Logger
}
// NewFFprobeService is the constructor.
func NewFFprobeService(cfg *config.Config, log *zap.Logger) *FFprobeService {
return &FFprobeService{cfg: cfg, log: log}
}
// ProbeResult is the subset of ffprobe output consumed by the scanner.
type ProbeResult struct {
DurationSec int
Width int
Height int
VideoCodec string
AudioCodec string
Container string
}
// Probe runs ffprobe against path and returns a typed result. A 30s timeout
// is applied so a single broken file does not hang the scanner.
func (f *FFprobeService) Probe(ctx context.Context, path string) (*ProbeResult, error) {
if f == nil {
return nil, errors.New("ffprobe service nil")
}
bin := f.cfg.App.FFprobePath
if bin == "" {
bin = "ffprobe"
}
probeCtx, cancel := context.WithTimeout(ctx, 30*time.Second)
defer cancel()
cmd := exec.CommandContext(probeCtx, bin,
"-v", "error",
"-print_format", "json",
"-show_format",
"-show_streams",
path,
)
out, err := cmd.Output()
if err != nil {
return nil, fmt.Errorf("ffprobe %s: %w", path, err)
}
return parseProbeJSON(out)
}
// rawProbe mirrors the relevant fields of `ffprobe -show_format -show_streams`.
type rawProbe struct {
Format struct {
Duration string `json:"duration"`
FormatName string `json:"format_name"`
} `json:"format"`
Streams []struct {
CodecType string `json:"codec_type"`
CodecName string `json:"codec_name"`
Width int `json:"width"`
Height int `json:"height"`
} `json:"streams"`
}
func parseProbeJSON(data []byte) (*ProbeResult, error) {
var raw rawProbe
if err := json.Unmarshal(data, &raw); err != nil {
return nil, fmt.Errorf("parse ffprobe json: %w", err)
}
res := &ProbeResult{Container: raw.Format.FormatName}
if d, err := strconv.ParseFloat(raw.Format.Duration, 64); err == nil {
res.DurationSec = int(d)
}
for _, s := range raw.Streams {
switch s.CodecType {
case "video":
if res.VideoCodec == "" {
res.VideoCodec = s.CodecName
res.Width = s.Width
res.Height = s.Height
}
case "audio":
if res.AudioCodec == "" {
res.AudioCodec = s.CodecName
}
}
}
return res, nil
}
+143
View File
@@ -0,0 +1,143 @@
// Package service — image proxy.
//
// Some deployments cannot reach image.tmdb.org directly (GFW, internal-only
// networks). ImageProxy fronts a remote image URL so the browser only ever
// talks to the MediaStationGo origin. The proxy:
//
// - validates the URL belongs to a small allow-list of trusted hosts,
// - streams bytes through with a small disk cache under cache/images,
// - falls back to a transparent 1×1 PNG on upstream failure so the UI
// never breaks layout.
package service
import (
"context"
"crypto/sha1"
"encoding/hex"
"errors"
"io"
"net/http"
"net/url"
"os"
"path/filepath"
"strings"
"sync"
"time"
"go.uber.org/zap"
"github.com/ShukeBta/MediaStationGo/internal/config"
)
// ImageProxy fetches and caches remote images on behalf of the browser.
type ImageProxy struct {
cfg *config.Config
log *zap.Logger
client *http.Client
cacheDir string
allowHost map[string]struct{}
mu sync.Mutex
}
// NewImageProxy is the constructor.
func NewImageProxy(cfg *config.Config, log *zap.Logger) *ImageProxy {
return &ImageProxy{
cfg: cfg,
log: log,
cacheDir: filepath.Join(cfg.Cache.CacheDir, "images"),
client: &http.Client{Timeout: 20 * time.Second},
allowHost: map[string]struct{}{
"image.tmdb.org": {},
"www.themoviedb.org": {},
"lain.bgm.tv": {},
"img.bgm.tv": {},
"webdav.bgm.tv": {},
"img1.doubanio.com": {},
"img2.doubanio.com": {},
"img3.doubanio.com": {},
"img9.doubanio.com": {},
"assets.fanart.tv": {},
"artworks.thetvdb.com": {},
},
}
}
// Serve writes the requested image to w. Caller is expected to validate
// the JWT before invoking it.
func (p *ImageProxy) Serve(ctx context.Context, w http.ResponseWriter, raw string) error {
if raw == "" {
return errors.New("missing url")
}
u, err := url.Parse(raw)
if err != nil || u.Scheme == "" || u.Host == "" {
return errors.New("invalid url")
}
if _, ok := p.allowHost[strings.ToLower(u.Host)]; !ok {
return errors.New("host not allowed")
}
// Cache key = sha1(url)
sum := sha1.Sum([]byte(raw))
key := hex.EncodeToString(sum[:])
cachePath := filepath.Join(p.cacheDir, key)
// Cache hit.
if f, err := os.Open(cachePath); err == nil {
defer f.Close()
stat, _ := f.Stat()
w.Header().Set("Cache-Control", "public, max-age=604800")
http.ServeContent(w, &http.Request{}, key, stat.ModTime(), f)
return nil
}
// Cache miss → fetch upstream.
if err := os.MkdirAll(p.cacheDir, 0o755); err != nil {
return err
}
req, err := http.NewRequestWithContext(ctx, http.MethodGet, raw, nil)
if err != nil {
return err
}
req.Header.Set("User-Agent", "MediaStationGo/0.1")
resp, err := p.client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
return errors.New("upstream returned " + resp.Status)
}
// Write to a temp file then rename for atomicity.
tmp, err := os.CreateTemp(p.cacheDir, "img-*.tmp")
if err != nil {
return err
}
if _, err := io.Copy(tmp, resp.Body); err != nil {
tmp.Close()
os.Remove(tmp.Name())
return err
}
tmp.Close()
if err := os.Rename(tmp.Name(), cachePath); err != nil {
os.Remove(tmp.Name())
}
// Now serve the freshly cached file.
f, err := os.Open(cachePath)
if err != nil {
return err
}
defer f.Close()
stat, _ := f.Stat()
for _, h := range []string{"Content-Type", "Content-Length", "ETag", "Last-Modified"} {
if v := resp.Header.Get(h); v != "" {
w.Header().Set(h, v)
}
}
w.Header().Set("Cache-Control", "public, max-age=604800")
http.ServeContent(w, &http.Request{}, key, stat.ModTime(), f)
return nil
}
+198
View File
@@ -0,0 +1,198 @@
// Package service — playback history / favourites / playlists.
//
// These three concerns are intentionally co-located: they all sit between
// "the user" and "a media item" and share the same join-table flavour. A
// dedicated PlaybackService keeps the wiring simple and lets handlers
// dispatch by feature instead of by repository.
package service
import (
"context"
"errors"
"time"
"go.uber.org/zap"
"github.com/ShukeBta/MediaStationGo/internal/model"
"github.com/ShukeBta/MediaStationGo/internal/repository"
)
// PlaybackService bundles history / favourite / playlist business logic.
type PlaybackService struct {
log *zap.Logger
repo *repository.Container
}
// NewPlaybackService is the constructor.
func NewPlaybackService(log *zap.Logger, repo *repository.Container) *PlaybackService {
return &PlaybackService{log: log, repo: repo}
}
// ─── History ────────────────────────────────────────────────────────────────
// RecordProgress upserts the resume position for a (user, media) pair. A
// position within 30 seconds of the duration auto-flags the item as
// completed so the home page can hide it from "Continue Watching".
func (p *PlaybackService) RecordProgress(ctx context.Context, userID, mediaID string, position, duration int64) error {
if userID == "" || mediaID == "" {
return errors.New("missing user or media")
}
completed := duration > 0 && position >= duration-30_000
h := &model.PlaybackHistory{
UserID: userID,
MediaID: mediaID,
PositionMs: position,
DurationMs: duration,
WatchedAt: time.Now(),
Completed: completed,
}
return p.repo.History.Upsert(ctx, h)
}
// HistoryItem joins the playback row with its media so the API consumer
// gets a fully-populated card without a second round-trip.
type HistoryItem struct {
model.PlaybackHistory
Media *model.Media `json:"media,omitempty"`
}
// RecentHistory returns the most recently-watched items for a user. We
// fetch the history rows first then attach each Media row in a single
// follow-up query.
func (p *PlaybackService) RecentHistory(ctx context.Context, userID string, limit int) ([]HistoryItem, error) {
rows, err := p.repo.History.ListByUser(ctx, userID, limit)
if err != nil {
return nil, err
}
items := make([]HistoryItem, 0, len(rows))
for i := range rows {
var m model.Media
if err := p.repo.DB.Where("id = ?", rows[i].MediaID).First(&m).Error; err == nil {
items = append(items, HistoryItem{PlaybackHistory: rows[i], Media: &m})
} else {
items = append(items, HistoryItem{PlaybackHistory: rows[i]})
}
}
return items, nil
}
// ─── Favourites ─────────────────────────────────────────────────────────────
// ToggleFavourite flips the favourite flag and reports the new state.
func (p *PlaybackService) ToggleFavourite(ctx context.Context, userID, mediaID string) (bool, error) {
return p.repo.Favorite.Toggle(ctx, userID, mediaID)
}
// ListFavourites returns every favourited media for a user.
func (p *PlaybackService) ListFavourites(ctx context.Context, userID string) ([]model.Media, error) {
favs, err := p.repo.Favorite.ListByUser(ctx, userID)
if err != nil {
return nil, err
}
if len(favs) == 0 {
return nil, nil
}
ids := make([]string, len(favs))
for i, f := range favs {
ids[i] = f.MediaID
}
var out []model.Media
err = p.repo.DB.Where("id IN ?", ids).
Order("created_at desc").Find(&out).Error
return out, err
}
// ─── Playlists ──────────────────────────────────────────────────────────────
// CreatePlaylist persists a new playlist owned by userID.
func (p *PlaybackService) CreatePlaylist(ctx context.Context, userID, name string, isPublic bool) (*model.Playlist, error) {
if name == "" {
return nil, errors.New("name required")
}
pl := &model.Playlist{UserID: userID, Name: name, IsPublic: isPublic}
if err := p.repo.Playlist.Create(ctx, pl); err != nil {
return nil, err
}
return pl, nil
}
// ListPlaylists returns every playlist owned by userID.
func (p *PlaybackService) ListPlaylists(ctx context.Context, userID string) ([]model.Playlist, error) {
return p.repo.Playlist.ListByUser(ctx, userID)
}
// PlaylistDetail returns the playlist together with its ordered media items.
type PlaylistDetail struct {
Playlist model.Playlist `json:"playlist"`
Items []model.Media `json:"items"`
}
// GetPlaylist returns the playlist + its ordered media. Visibility is
// enforced at the handler level; the service trusts callers.
func (p *PlaybackService) GetPlaylist(ctx context.Context, playlistID string) (*PlaylistDetail, error) {
var pl model.Playlist
if err := p.repo.DB.Where("id = ?", playlistID).First(&pl).Error; err != nil {
return nil, err
}
var rows []model.PlaylistItem
if err := p.repo.DB.
Where("playlist_id = ?", playlistID).
Order("position asc").
Find(&rows).Error; err != nil {
return nil, err
}
if len(rows) == 0 {
return &PlaylistDetail{Playlist: pl}, nil
}
ids := make([]string, len(rows))
for i, r := range rows {
ids[i] = r.MediaID
}
var media []model.Media
if err := p.repo.DB.Where("id IN ?", ids).Find(&media).Error; err != nil {
return nil, err
}
// Preserve playlist order.
byID := make(map[string]model.Media, len(media))
for _, m := range media {
byID[m.ID] = m
}
ordered := make([]model.Media, 0, len(rows))
for _, r := range rows {
if m, ok := byID[r.MediaID]; ok {
ordered = append(ordered, m)
}
}
return &PlaylistDetail{Playlist: pl, Items: ordered}, nil
}
// AddToPlaylist appends a media item to the end of a playlist.
func (p *PlaybackService) AddToPlaylist(ctx context.Context, playlistID, mediaID string) error {
var count int64
if err := p.repo.DB.Model(&model.PlaylistItem{}).
Where("playlist_id = ?", playlistID).Count(&count).Error; err != nil {
return err
}
item := &model.PlaylistItem{
PlaylistID: playlistID,
MediaID: mediaID,
Position: int(count) + 1,
}
return p.repo.DB.Create(item).Error
}
// RemoveFromPlaylist removes a media item from a playlist (idempotent).
func (p *PlaybackService) RemoveFromPlaylist(ctx context.Context, playlistID, mediaID string) error {
return p.repo.DB.
Where("playlist_id = ? AND media_id = ?", playlistID, mediaID).
Delete(&model.PlaylistItem{}).Error
}
// DeletePlaylist removes a playlist and all of its items.
func (p *PlaybackService) DeletePlaylist(ctx context.Context, playlistID string) error {
if err := p.repo.DB.Where("playlist_id = ?", playlistID).
Delete(&model.PlaylistItem{}).Error; err != nil {
return err
}
return p.repo.DB.Where("id = ?", playlistID).Delete(&model.Playlist{}).Error
}
+73 -15
View File
@@ -1,10 +1,12 @@
// Package service — filesystem scanner.
//
// ScannerService walks the configured library roots looking for video files,
// then upserts a model.Media row per file. A future iteration will plug
// ffprobe / a metadata-provider chain on top of this skeleton, but the
// scaffold keeps the surface narrow and synchronous so handlers can call
// "POST /api/libraries/:id/scan" today.
// ScannerService walks the configured library roots looking for video
// files, then upserts a model.Media row per file. Each upsert also runs
// ffprobe (when available) and queues a TMDb lookup for newly added rows.
//
// The scan is synchronous from the HTTP layer's point of view, but it
// publishes WebSocket progress events on the "scan" topic so the React
// UI can render a live counter / spinner.
package service
import (
@@ -39,15 +41,27 @@ var videoExtensions = map[string]struct{}{
// ScannerService walks libraries on disk and upserts model.Media rows.
type ScannerService struct {
cfg *config.Config
log *zap.Logger
repo *repository.Container
hub *Hub
cfg *config.Config
log *zap.Logger
repo *repository.Container
hub *Hub
probe *FFprobeService
scraper *ScraperService
}
// NewScannerService is the constructor.
func NewScannerService(cfg *config.Config, log *zap.Logger, repo *repository.Container, hub *Hub) *ScannerService {
return &ScannerService{cfg: cfg, log: log, repo: repo, hub: hub}
func NewScannerService(
cfg *config.Config,
log *zap.Logger,
repo *repository.Container,
hub *Hub,
probe *FFprobeService,
scraper *ScraperService,
) *ScannerService {
return &ScannerService{
cfg: cfg, log: log, repo: repo, hub: hub,
probe: probe, scraper: scraper,
}
}
// ScanResult summarises a scan run.
@@ -55,19 +69,26 @@ type ScanResult struct {
LibraryID string `json:"library_id"`
Visited int `json:"visited"`
Added int `json:"added"`
Probed int `json:"probed"`
}
// ScanLibrary walks the library root and persists discovered media files.
//
// This is a synchronous skeleton: large libraries should call it in a
// goroutine. WebSocket progress events are pushed to the hub on the
// "scan" topic so the React UI can display a progress indicator.
// Workflow per file:
// 1. fast filename-based title cleanup.
// 2. ffprobe → duration / resolution / codecs (best effort).
// 3. upsert into the media table.
// 4. publish progress over the WS hub.
//
// After the walk we kick off the TMDb scraper for every still-pending
// row in the same library. The scraper has its own throttle.
func (s *ScannerService) ScanLibrary(ctx context.Context, libraryID string) (*ScanResult, error) {
lib, err := s.repo.Library.FindByID(ctx, libraryID)
if err != nil || lib == nil {
return nil, err
}
res := &ScanResult{LibraryID: lib.ID}
walkFn := func(path string, info walkInfo) error {
if info.isDir {
return nil
@@ -77,14 +98,38 @@ func (s *ScannerService) ScanLibrary(ctx context.Context, libraryID string) (*Sc
return nil
}
res.Visited++
title := strings.TrimSuffix(filepath.Base(path), ext)
title, year := CleanQuery(path)
if title == "" {
title = strings.TrimSuffix(filepath.Base(path), ext)
}
m := &model.Media{
LibraryID: lib.ID,
Title: title,
Year: year,
Path: path,
SizeBytes: info.size,
Container: strings.TrimPrefix(ext, "."),
}
// Best-effort ffprobe; failure does not abort the file.
if s.probe != nil {
if probe, err := s.probe.Probe(ctx, path); err == nil && probe != nil {
m.DurationSec = probe.DurationSec
m.Width = probe.Width
m.Height = probe.Height
m.VideoCodec = probe.VideoCodec
m.AudioCodec = probe.AudioCodec
if probe.Container != "" {
m.Container = probe.Container
}
res.Probed++
} else if err != nil {
s.log.Debug("ffprobe failed", zap.String("path", path), zap.Error(err))
}
}
if err := s.repo.Media.Upsert(ctx, m); err != nil {
s.log.Warn("upsert media failed", zap.String("path", path), zap.Error(err))
return nil
@@ -95,17 +140,30 @@ func (s *ScannerService) ScanLibrary(ctx context.Context, libraryID string) (*Sc
"path": path,
"visited": res.Visited,
"added": res.Added,
"probed": res.Probed,
})
return nil
}
if err := walk(lib.Path, walkFn); err != nil {
return res, err
}
s.hub.Publish("scan", map[string]any{
"library_id": lib.ID,
"finished": true,
"visited": res.Visited,
"added": res.Added,
"probed": res.Probed,
})
// Fire-and-forget metadata enrichment when a TMDb key is configured.
if s.scraper != nil && s.scraper.tmdb != nil && s.scraper.tmdb.Enabled() {
go func(libID string) {
if _, err := s.scraper.EnrichLibrary(context.Background(), libID); err != nil {
s.log.Warn("scraper enrich failed", zap.Error(err))
}
}(lib.ID)
}
return res, nil
}
+164
View File
@@ -0,0 +1,164 @@
// Package service — scraper orchestrator.
//
// ScraperService takes a Media row and tries to enrich it with metadata
// from one or more providers (currently TMDb only). It is invoked at the
// end of every scan cycle for media items whose `scrape_status` is still
// "pending"; it can also be re-triggered manually from the admin UI.
//
// The orchestrator is deliberately stateless: it loops media → provider →
// repository, publishing scrape progress events to the WS hub.
package service
import (
"context"
"path/filepath"
"regexp"
"strconv"
"strings"
"time"
"go.uber.org/zap"
"github.com/ShukeBta/MediaStationGo/internal/config"
"github.com/ShukeBta/MediaStationGo/internal/model"
"github.com/ShukeBta/MediaStationGo/internal/repository"
)
// ScraperService coordinates metadata enrichment across providers.
type ScraperService struct {
cfg *config.Config
log *zap.Logger
repo *repository.Container
tmdb *TMDbProvider
hub *Hub
}
// NewScraperService is the constructor.
func NewScraperService(cfg *config.Config, log *zap.Logger, repo *repository.Container, tmdb *TMDbProvider, hub *Hub) *ScraperService {
return &ScraperService{cfg: cfg, log: log, repo: repo, tmdb: tmdb, hub: hub}
}
// yearPattern extracts a 4-digit year from a filename (1900-2099).
var yearPattern = regexp.MustCompile(`(?:^|[^\d])(19\d{2}|20\d{2})(?:[^\d]|$)`)
// noiseTokens are aggressively stripped from filenames before search.
// Keep in sync with nowen-video's filename_parser.go intent.
var noiseTokens = []string{
"1080p", "2160p", "4k", "720p", "480p",
"hdrip", "bluray", "blu-ray", "webrip", "web-dl", "web",
"x264", "x265", "h264", "h265", "hevc", "avc",
"hdr", "sdr", "dts", "ddp", "atmos", "aac", "ac3", "flac",
"remux", "extended", "uncut", "directors-cut", "directors_cut",
"hkfree", "yify", "rarbg", "ettv", "fgt",
}
// bracketedTag matches "[anything]" or "(anything)" segments, which are
// almost always release-group / encoder tags in scene filenames.
var bracketedTag = regexp.MustCompile(`[\[\(][^\]\)]*[\]\)]`)
// CleanQuery converts a filename like "Inception.2010.1080p.BluRay.x264.mkv"
// into a TMDb-friendly title plus an optional year hint.
func CleanQuery(raw string) (title string, year int) {
name := strings.TrimSuffix(filepath.Base(raw), filepath.Ext(raw))
lower := strings.ToLower(name)
// 1. Year first — bracketed years (1999) must survive the next step.
if m := yearPattern.FindStringSubmatch(lower); len(m) >= 2 {
if v, err := strconv.Atoi(m[1]); err == nil {
year = v
lower = strings.ReplaceAll(lower, m[1], " ")
}
}
// 2. Drop everything inside square / round brackets — those are tags.
lower = bracketedTag.ReplaceAllString(lower, " ")
for _, t := range noiseTokens {
lower = strings.ReplaceAll(lower, t, " ")
}
// collapse separators / spaces
for _, sep := range []string{".", "_", "-", "[", "]", "(", ")"} {
lower = strings.ReplaceAll(lower, sep, " ")
}
fields := strings.Fields(lower)
title = strings.Join(fields, " ")
return strings.TrimSpace(title), year
}
// EnrichOne runs the provider chain for a single media row.
func (s *ScraperService) EnrichOne(ctx context.Context, m *model.Media) error {
if s.tmdb == nil || !s.tmdb.Enabled() {
return nil
}
query := m.Title
if query == "" {
query, _ = CleanQuery(m.Path)
} else {
query, _ = CleanQuery(query)
}
year := m.Year
if year == 0 {
_, year = CleanQuery(filepath.Base(m.Path))
}
match, err := s.tmdb.SearchMovie(ctx, query, year)
if err != nil || match == nil {
return err
}
updates := map[string]any{
"title": match.Title,
"overview": match.Overview,
"poster_url": match.PosterURL,
"backdrop_url": match.BackdropURL,
"rating": match.Rating,
"year": match.Year,
"tmdb_id": match.TMDbID,
"scrape_status": "matched",
}
if err := s.repo.DB.Model(&model.Media{}).Where("id = ?", m.ID).
Updates(updates).Error; err != nil {
return err
}
s.hub.Publish("scrape", map[string]any{
"media_id": m.ID,
"title": match.Title,
"tmdb_id": match.TMDbID,
})
return nil
}
// EnrichLibrary runs the provider chain for every "pending" media in a
// library. It throttles to 4 RPS to stay below TMDb's rate limit and
// publishes a summary event when done.
func (s *ScraperService) EnrichLibrary(ctx context.Context, libraryID string) (int, error) {
if s.tmdb == nil || !s.tmdb.Enabled() {
return 0, nil
}
var rows []model.Media
q := s.repo.DB.Where("scrape_status = ?", "pending")
if libraryID != "" {
q = q.Where("library_id = ?", libraryID)
}
if err := q.Find(&rows).Error; err != nil {
return 0, err
}
matched := 0
for i := range rows {
select {
case <-ctx.Done():
return matched, ctx.Err()
default:
}
if err := s.EnrichOne(ctx, &rows[i]); err != nil {
s.log.Warn("enrich failed", zap.String("media", rows[i].ID), zap.Error(err))
continue
}
matched++
time.Sleep(250 * time.Millisecond) // ~4 RPS
}
s.hub.Publish("scrape", map[string]any{
"library_id": libraryID,
"finished": true,
"matched": matched,
})
return matched, nil
}
+26
View File
@@ -0,0 +1,26 @@
package service
import "testing"
func TestCleanQuery(t *testing.T) {
cases := []struct {
in string
wantTitle string
wantYear int
}{
{"Inception.2010.1080p.BluRay.x264.mkv", "inception", 2010},
{"The_Matrix_(1999).1080p.WEB-DL.H265.mp4", "the matrix", 1999},
{"interstellar.2014.4k.hdr.dts.atmos.mkv", "interstellar", 2014},
{"My Movie 2022 [HDR] (1080p) [TGx].mp4", "my movie", 2022},
{"NoYearOrTags.mkv", "noyearortags", 0},
}
for _, tc := range cases {
t.Run(tc.in, func(t *testing.T) {
gotTitle, gotYear := CleanQuery(tc.in)
if gotTitle != tc.wantTitle || gotYear != tc.wantYear {
t.Errorf("CleanQuery(%q) = (%q, %d), want (%q, %d)",
tc.in, gotTitle, gotYear, tc.wantTitle, tc.wantYear)
}
})
}
}
+39 -17
View File
@@ -14,34 +14,56 @@ import (
// Container holds every service initialized at startup. Handlers receive a
// pointer to it and pick the relevant fields.
type Container struct {
Cfg *config.Config
Log *zap.Logger
Repo *repository.Container
WSHub *Hub
Auth *AuthService
Media *MediaService
Scan *ScannerService
Stream *StreamService
Cfg *config.Config
Log *zap.Logger
Repo *repository.Container
WSHub *Hub
Auth *AuthService
Media *MediaService
Scan *ScannerService
Stream *StreamService
Transcoder *TranscoderService
FFprobe *FFprobeService
TMDb *TMDbProvider
Scraper *ScraperService
Playback *PlaybackService
ImageProxy *ImageProxy
}
// New builds the service container.
func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Container {
hub := NewHub(log)
go hub.Run()
probe := NewFFprobeService(cfg, log)
tmdb := NewTMDbProvider(cfg, log)
scraper := NewScraperService(cfg, log, repos, tmdb, hub)
transcoder := NewTranscoderService(cfg, log, repos, hub)
return &Container{
Cfg: cfg,
Log: log,
Repo: repos,
WSHub: hub,
Auth: NewAuthService(cfg, log, repos),
Media: NewMediaService(cfg, log, repos),
Scan: NewScannerService(cfg, log, repos, hub),
Stream: NewStreamService(cfg, log, repos),
Cfg: cfg,
Log: log,
Repo: repos,
WSHub: hub,
Auth: NewAuthService(cfg, log, repos),
Media: NewMediaService(cfg, log, repos),
Scan: NewScannerService(cfg, log, repos, hub, probe, scraper),
Stream: NewStreamService(cfg, log, repos, transcoder),
Transcoder: transcoder,
FFprobe: probe,
TMDb: tmdb,
Scraper: scraper,
Playback: NewPlaybackService(log, repos),
ImageProxy: NewImageProxy(cfg, log),
}
}
// Close releases any resources held by services (e.g. the websocket hub).
// Close releases any resources held by services (websocket hub, ffmpeg
// transcodes).
func (c *Container) Close() {
if c.Transcoder != nil {
c.Transcoder.StopAll()
}
if c.WSHub != nil {
c.WSHub.Stop()
}
+105 -11
View File
@@ -1,10 +1,29 @@
// Package service — direct-play / range request streaming.
// Package service — direct-play / HLS streaming.
//
// StreamService exposes two flavours of playback:
//
// - Direct play: the original file is served with HTTP Range support.
// Works for browser-friendly containers (mp4 / webm / m4v), no ffmpeg
// involved, zero CPU overhead.
// - HLS: when the client opts in (or the source codec / container is
// not browser-friendly), the TranscoderService runs ffmpeg in the
// background and we serve the resulting .m3u8 + .ts files directly.
//
// The HTTP layer decides which mode to use based on the request path:
//
// GET /api/stream/:id → direct play
// GET /api/hls/:id/index.m3u8 → HLS playlist
// GET /api/hls/:id/seg_NNNNN.ts → HLS segment
package service
import (
"context"
"errors"
"net/http"
"os"
"path/filepath"
"strings"
"time"
"go.uber.org/zap"
@@ -14,20 +33,21 @@ import (
// StreamService serves media files with proper Range support so browsers can
// seek into the stream.
//
// HLS / on-demand transcoding is intentionally omitted from this initial
// scaffold. The HTTP handler returns 501 (NotImplemented) for that path,
// while direct-play already works for browser-friendly containers (mp4 /
// webm / m4v).
type StreamService struct {
cfg *config.Config
log *zap.Logger
repo *repository.Container
cfg *config.Config
log *zap.Logger
repo *repository.Container
transcoder *TranscoderService
}
// NewStreamService is the constructor.
func NewStreamService(cfg *config.Config, log *zap.Logger, repo *repository.Container) *StreamService {
return &StreamService{cfg: cfg, log: log, repo: repo}
func NewStreamService(cfg *config.Config, log *zap.Logger, repo *repository.Container, transcoder *TranscoderService) *StreamService {
return &StreamService{
cfg: cfg,
log: log,
repo: repo,
transcoder: transcoder,
}
}
// ErrMediaNotFound is returned when the media row or its file is missing.
@@ -56,3 +76,77 @@ func (s *StreamService) ServeFile(w http.ResponseWriter, r *http.Request, mediaI
http.ServeContent(w, r, stat.Name(), stat.ModTime(), f)
return nil
}
// ServeHLSPlaylist makes sure a transcode is running and writes the m3u8.
// We block (with a 30s timeout) until the playlist file shows up.
func (s *StreamService) ServeHLSPlaylist(w http.ResponseWriter, r *http.Request, mediaID string) error {
if _, err := s.transcoder.EnsureJob(r.Context(), mediaID); err != nil {
return err
}
if !s.transcoder.WaitReady(r.Context(), mediaID, 30*time.Second) {
return errors.New("hls playlist not ready")
}
playlist := s.transcoder.PlaylistPath(mediaID)
f, err := os.Open(playlist)
if err != nil {
return err
}
defer f.Close()
stat, _ := f.Stat()
w.Header().Set("Content-Type", "application/vnd.apple.mpegurl")
w.Header().Set("Cache-Control", "no-cache")
http.ServeContent(w, r, stat.Name(), stat.ModTime(), f)
return nil
}
// ServeHLSSegment writes a single .ts segment from the on-disk cache.
func (s *StreamService) ServeHLSSegment(w http.ResponseWriter, r *http.Request, mediaID, segment string) error {
// Only allow segments that look like seg_NNNNN.ts so we cannot be tricked
// into reading arbitrary files via path traversal.
if !strings.HasPrefix(segment, "seg_") || !strings.HasSuffix(segment, ".ts") {
return errors.New("bad segment")
}
full := filepath.Join(s.transcoder.HLSDir(mediaID), segment)
abs, err := filepath.Abs(full)
if err != nil {
return err
}
dir, _ := filepath.Abs(s.transcoder.HLSDir(mediaID))
if !strings.HasPrefix(abs, dir) {
return errors.New("path escape")
}
f, err := os.Open(abs)
if err != nil {
return err
}
defer f.Close()
stat, _ := f.Stat()
w.Header().Set("Content-Type", "video/mp2t")
w.Header().Set("Cache-Control", "public, max-age=3600")
http.ServeContent(w, r, stat.Name(), stat.ModTime(), f)
return nil
}
// Probe re-runs ffprobe against an existing media row and refreshes the
// extracted metadata. Used by the admin UI's "rescan" button.
func (s *StreamService) Probe(ctx context.Context, mediaID string, probe *FFprobeService) error {
m, err := s.repo.Media.FindByID(ctx, mediaID)
if err != nil || m == nil {
return ErrMediaNotFound
}
res, err := probe.Probe(ctx, m.Path)
if err != nil {
return err
}
updates := map[string]any{
"duration_sec": res.DurationSec,
"width": res.Width,
"height": res.Height,
"video_codec": res.VideoCodec,
"audio_codec": res.AudioCodec,
}
if res.Container != "" {
updates["container"] = res.Container
}
return s.repo.DB.Model(m).Updates(updates).Error
}
+148
View File
@@ -0,0 +1,148 @@
// Package service — TMDb metadata provider.
//
// TMDbProvider implements the (minimal) MetadataProvider interface and uses
// the public The Movie Database REST API. The API key is taken from
// secrets.tmdb_api_key; when empty the provider returns nil from every
// method so the scraper can no-op gracefully.
//
// We only call the two endpoints the scrape pipeline actually needs:
//
// GET /search/movie?query=...&year=...
// GET /movie/{id}?language=zh-CN
//
// TV / anime support follows the same pattern; for the bootstrap we expose
// a single SearchMovie path so that the home page and library gallery can
// show real posters as soon as a TMDb key is configured.
package service
import (
"context"
"encoding/json"
"errors"
"fmt"
"net/http"
"net/url"
"time"
"go.uber.org/zap"
"github.com/ShukeBta/MediaStationGo/internal/config"
)
// TMDbProvider talks to https://api.themoviedb.org/3.
type TMDbProvider struct {
cfg *config.Config
log *zap.Logger
client *http.Client
base string
imgCDN string
}
// NewTMDbProvider is the constructor. APIBase / image CDN can be overridden
// via secrets.tmdb_api_proxy + tmdb_image_proxy for users behind GFW.
func NewTMDbProvider(cfg *config.Config, log *zap.Logger) *TMDbProvider {
base := cfg.Secrets.TMDbAPIProxy
if base == "" {
base = "https://api.themoviedb.org/3"
}
img := cfg.Secrets.TMDbImageProxy
if img == "" {
img = "https://image.tmdb.org/t/p"
}
return &TMDbProvider{
cfg: cfg,
log: log,
base: base,
imgCDN: img,
client: &http.Client{Timeout: 15 * time.Second},
}
}
// Enabled reports whether the operator has supplied an API key.
func (t *TMDbProvider) Enabled() bool { return t.cfg.Secrets.TMDbAPIKey != "" }
// Match describes a successful metadata match.
type Match struct {
TMDbID int `json:"tmdb_id"`
Title string `json:"title"`
Overview string `json:"overview"`
PosterURL string `json:"poster_url"`
BackdropURL string `json:"backdrop_url"`
Year int `json:"year"`
Rating float32 `json:"rating"`
}
// SearchMovie issues `/search/movie` and returns the best match, or nil
// when no result is found. The `year` argument is optional (0 = any).
func (t *TMDbProvider) SearchMovie(ctx context.Context, query string, year int) (*Match, error) {
if !t.Enabled() {
return nil, nil
}
if query == "" {
return nil, errors.New("empty query")
}
q := url.Values{}
q.Set("api_key", t.cfg.Secrets.TMDbAPIKey)
q.Set("query", query)
q.Set("language", "zh-CN")
q.Set("include_adult", "false")
if year > 0 {
q.Set("year", fmt.Sprintf("%d", year))
}
u := t.base + "/search/movie?" + q.Encode()
type result struct {
ID int `json:"id"`
Title string `json:"title"`
Overview string `json:"overview"`
PosterPath string `json:"poster_path"`
BackdropPath string `json:"backdrop_path"`
ReleaseDate string `json:"release_date"`
VoteAverage float32 `json:"vote_average"`
}
type page struct {
Results []result `json:"results"`
}
var p page
if err := t.getJSON(ctx, u, &p); err != nil {
return nil, err
}
if len(p.Results) == 0 {
return nil, nil
}
r := p.Results[0]
m := &Match{
TMDbID: r.ID,
Title: r.Title,
Overview: r.Overview,
Rating: r.VoteAverage,
}
if r.PosterPath != "" {
m.PosterURL = t.imgCDN + "/w500" + r.PosterPath
}
if r.BackdropPath != "" {
m.BackdropURL = t.imgCDN + "/w1280" + r.BackdropPath
}
if len(r.ReleaseDate) >= 4 {
fmt.Sscanf(r.ReleaseDate[:4], "%d", &m.Year)
}
return m, nil
}
func (t *TMDbProvider) getJSON(ctx context.Context, url string, out any) error {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
if err != nil {
return err
}
resp, err := t.client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
return fmt.Errorf("tmdb %s: %d", url, resp.StatusCode)
}
return json.NewDecoder(resp.Body).Decode(out)
}
+234
View File
@@ -0,0 +1,234 @@
// Package service — HLS on-demand transcoder.
//
// TranscoderService spawns ffmpeg processes that segment a source media file
// into HLS (.m3u8 + .ts). The output lives under cache.cache_dir/hls/<id>.
// The HTTP layer serves these files directly with a normal http.FileServer.
//
// Concurrency model:
// - Each Media has at most one active ffmpeg job.
// - jobs[mediaID] tracks the running goroutine + cancel func.
// - Calling Start while a job already exists is a no-op.
// - When the playlist file appears on disk we consider the job "ready"
// and unblock the HTTP handler that was waiting on it.
//
// The transcode profile is intentionally conservative: a single 720p/1.5M
// bitrate, AAC stereo audio, MPEG-TS segments. Hardware acceleration
// (NVENC / QSV / VAAPI) is left as a future config-driven extension point.
package service
import (
"context"
"errors"
"fmt"
"os"
"os/exec"
"path/filepath"
"sync"
"time"
"go.uber.org/zap"
"github.com/ShukeBta/MediaStationGo/internal/config"
"github.com/ShukeBta/MediaStationGo/internal/repository"
)
// TranscoderService orchestrates background ffmpeg transcodes.
type TranscoderService struct {
cfg *config.Config
log *zap.Logger
repo *repository.Container
hub *Hub
mu sync.Mutex
jobs map[string]*hlsJob
}
// hlsJob holds the live state of one ffmpeg run.
type hlsJob struct {
mediaID string
outputDir string
cancel context.CancelFunc
startedAt time.Time
playlistOK bool
}
// NewTranscoderService is the constructor.
func NewTranscoderService(cfg *config.Config, log *zap.Logger, repo *repository.Container, hub *Hub) *TranscoderService {
return &TranscoderService{
cfg: cfg,
log: log,
repo: repo,
hub: hub,
jobs: make(map[string]*hlsJob),
}
}
// HLSDir is the per-media directory that holds index.m3u8 + segment files.
func (t *TranscoderService) HLSDir(mediaID string) string {
return filepath.Join(t.cfg.Cache.CacheDir, "hls", mediaID)
}
// PlaylistPath returns the absolute path of the m3u8 playlist for a media.
func (t *TranscoderService) PlaylistPath(mediaID string) string {
return filepath.Join(t.HLSDir(mediaID), "index.m3u8")
}
// EnsureJob makes sure a transcode is running for mediaID. The function is
// non-blocking: it returns the playlist path immediately. The caller is
// expected to poll until WaitReady reports true.
func (t *TranscoderService) EnsureJob(ctx context.Context, mediaID string) (string, error) {
m, err := t.repo.Media.FindByID(ctx, mediaID)
if err != nil {
return "", err
}
if m == nil {
return "", ErrMediaNotFound
}
if _, err := os.Stat(m.Path); err != nil {
return "", ErrMediaNotFound
}
t.mu.Lock()
if _, ok := t.jobs[mediaID]; ok {
t.mu.Unlock()
return t.PlaylistPath(mediaID), nil
}
outDir := t.HLSDir(mediaID)
if err := os.MkdirAll(outDir, 0o755); err != nil {
t.mu.Unlock()
return "", err
}
jobCtx, cancel := context.WithCancel(context.Background())
job := &hlsJob{
mediaID: mediaID,
outputDir: outDir,
cancel: cancel,
startedAt: time.Now(),
}
t.jobs[mediaID] = job
t.mu.Unlock()
go t.runFFmpeg(jobCtx, job, m.Path)
return t.PlaylistPath(mediaID), nil
}
// WaitReady blocks (with a deadline) until the playlist file shows up on
// disk. Returns true on success.
func (t *TranscoderService) WaitReady(ctx context.Context, mediaID string, timeout time.Duration) bool {
deadline := time.Now().Add(timeout)
for {
if _, err := os.Stat(t.PlaylistPath(mediaID)); err == nil {
t.mu.Lock()
if j, ok := t.jobs[mediaID]; ok {
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.
func (t *TranscoderService) StopJob(mediaID string) {
t.mu.Lock()
defer t.mu.Unlock()
if j, ok := t.jobs[mediaID]; ok {
j.cancel()
delete(t.jobs, mediaID)
}
}
// 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)
}
}
func (t *TranscoderService) runFFmpeg(ctx context.Context, job *hlsJob, source string) {
bin := t.cfg.App.FFmpegPath
if bin == "" {
bin = "ffmpeg"
}
playlist := filepath.Join(job.outputDir, "index.m3u8")
segments := filepath.Join(job.outputDir, "seg_%05d.ts")
args := []string{
"-y",
"-fflags", "+genpts",
"-i", source,
// Single 720p H.264 video rendition.
"-map", "0:v:0?",
"-map", "0:a:0?",
"-vf", "scale=-2:min(720\\,ih)",
"-c:v", "libx264",
"-preset", "veryfast",
"-profile:v", "main",
"-level", "4.0",
"-pix_fmt", "yuv420p",
"-b:v", "1500k",
"-maxrate", "1800k",
"-bufsize", "3000k",
"-c:a", "aac",
"-ar", "48000",
"-b:a", "128k",
"-ac", "2",
"-force_key_frames", "expr:gte(t,n_forced*4)",
"-f", "hls",
"-hls_time", "4",
"-hls_list_size", "0",
"-hls_segment_type", "mpegts",
"-hls_flags", "independent_segments",
"-hls_segment_filename", segments,
playlist,
}
cmd := exec.CommandContext(ctx, bin, args...)
cmd.Stderr = os.Stderr
t.log.Info("transcode started",
zap.String("media_id", job.mediaID),
zap.String("source", source),
)
t.hub.Publish("transcode", map[string]any{
"media_id": job.mediaID,
"status": "started",
})
if err := cmd.Run(); err != nil && !errors.Is(ctx.Err(), context.Canceled) {
t.log.Warn("ffmpeg exited",
zap.String("media_id", job.mediaID),
zap.Error(err),
)
}
t.mu.Lock()
delete(t.jobs, job.mediaID)
t.mu.Unlock()
t.hub.Publish("transcode", map[string]any{
"media_id": job.mediaID,
"status": "stopped",
"duration": time.Since(job.startedAt).Seconds(),
})
}
// HumanFFmpegProfile is exposed for the admin UI / settings view.
func (t *TranscoderService) HumanFFmpegProfile() string {
return fmt.Sprintf("ffmpeg=%s, output=%s",
t.cfg.App.FFmpegPath, filepath.Join(t.cfg.Cache.CacheDir, "hls"))
}