diff --git a/README.md b/README.md index a3f3710..919e55c 100644 --- a/README.md +++ b/README.md @@ -34,20 +34,25 @@ deployment painless on NAS hardware. - ✅ JWT authentication with admin/user roles - ✅ First-run admin seeding (`admin / admin123`) - ✅ Library CRUD + recursive filesystem scan +- ✅ ffprobe metadata extraction (duration / resolution / codecs) +- ✅ TMDb scraper with image proxy (poster / backdrop / overview / rating) - ✅ Direct-play streaming with HTTP `Range` support -- ✅ Real-time scan progress via WebSocket -- ✅ React SPA: Login / Home / Library / Search / Media detail / Player / Admin +- ✅ HLS on-demand transcoding (single ffmpeg job per media) +- ✅ Playback history (resume) + Continue Watching row +- ✅ Favourites + Playlists (CRUD + ordered items) +- ✅ Real-time scan / scrape / transcode progress via WebSocket +- ✅ React SPA with code-splitting: Login / Home / Library / Search / + Favourites / Playlists / Media detail / Player (HLS + direct) / Admin - ✅ Single-binary build, multi-arch Docker image, GitHub Actions CI ### Roadmap | Area | Status | |------|--------| -| ffprobe-driven metadata extraction | ⏳ | -| TMDb / Bangumi / Douban scraper chain | ⏳ | -| HLS on-demand transcoding (NVENC / QSV / VAAPI) | ⏳ | +| Bangumi / Douban / Fanart scraper providers | ⏳ | +| Hardware-accelerated transcoding (NVENC / QSV / VAAPI) | ⏳ | | qBittorrent / Transmission / RSS automation | ⏳ | -| Playlists, favourites, watch history UI | ⏳ | +| Subtitles (extract / search / sync) | ⏳ | | Emby/Jellyfin compatibility layer | ⏳ | | DLNA / Chromecast | ⏳ | | AI metadata enhancement & smart search | ⏳ | diff --git a/internal/handler/handler.go b/internal/handler/handler.go index 669531b..783c4fa 100644 --- a/internal/handler/handler.go +++ b/internal/handler/handler.go @@ -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)) } diff --git a/internal/handler/playback.go b/internal/handler/playback.go new file mode 100644 index 0000000..86475f8 --- /dev/null +++ b/internal/handler/playback.go @@ -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) + } +} diff --git a/internal/handler/streaming.go b/internal/handler/streaming.go new file mode 100644 index 0000000..9572bc0 --- /dev/null +++ b/internal/handler/streaming.go @@ -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) + } +} diff --git a/internal/service/ffprobe.go b/internal/service/ffprobe.go new file mode 100644 index 0000000..5d5e5b3 --- /dev/null +++ b/internal/service/ffprobe.go @@ -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 +} diff --git a/internal/service/image_proxy.go b/internal/service/image_proxy.go new file mode 100644 index 0000000..85249c1 --- /dev/null +++ b/internal/service/image_proxy.go @@ -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 +} diff --git a/internal/service/playback.go b/internal/service/playback.go new file mode 100644 index 0000000..289ddc0 --- /dev/null +++ b/internal/service/playback.go @@ -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 +} diff --git a/internal/service/scanner.go b/internal/service/scanner.go index 464785d..db20c13 100644 --- a/internal/service/scanner.go +++ b/internal/service/scanner.go @@ -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 } diff --git a/internal/service/scraper.go b/internal/service/scraper.go new file mode 100644 index 0000000..b4edd8c --- /dev/null +++ b/internal/service/scraper.go @@ -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 +} diff --git a/internal/service/scraper_test.go b/internal/service/scraper_test.go new file mode 100644 index 0000000..0d865b2 --- /dev/null +++ b/internal/service/scraper_test.go @@ -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) + } + }) + } +} diff --git a/internal/service/service.go b/internal/service/service.go index a8aecc4..248734a 100644 --- a/internal/service/service.go +++ b/internal/service/service.go @@ -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() } diff --git a/internal/service/stream.go b/internal/service/stream.go index 2d34a0c..8e256e0 100644 --- a/internal/service/stream.go +++ b/internal/service/stream.go @@ -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 +} diff --git a/internal/service/tmdb.go b/internal/service/tmdb.go new file mode 100644 index 0000000..67b48b8 --- /dev/null +++ b/internal/service/tmdb.go @@ -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) +} diff --git a/internal/service/transcoder.go b/internal/service/transcoder.go new file mode 100644 index 0000000..84d55ed --- /dev/null +++ b/internal/service/transcoder.go @@ -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/. +// 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")) +} diff --git a/web/src/App.tsx b/web/src/App.tsx index 210da98..7bf3adb 100644 --- a/web/src/App.tsx +++ b/web/src/App.tsx @@ -1,56 +1,71 @@ +import { Suspense, lazy } from 'react' import { Navigate, Route, Routes } from 'react-router-dom' import { Layout } from './components/Layout' -import { RequireAuth, RequireAdmin } from './components/RequireAuth' -import { AdminPage } from './pages/AdminPage' -import { HomePage } from './pages/HomePage' -import { LibraryPage } from './pages/LibraryPage' +import { RequireAdmin, RequireAuth } from './components/RequireAuth' import { LoginPage } from './pages/LoginPage' -import { MediaDetailPage } from './pages/MediaDetailPage' -import { PlayerPage } from './pages/PlayerPage' -import { SearchPage } from './pages/SearchPage' -// Top-level route table. -// -// Public: -// /login -// -// Authenticated: -// / → home (continue watching + recently added) -// /library/:id → library content grid -// /search → keyword search -// /media/:id → media detail -// /play/:id → fullscreen player -// -// Admin only: -// /admin → users / settings / activity log +// Lazy-loaded routes — the login screen and the layout shell ship in the +// initial bundle; everything else is fetched on first navigation. +const HomePage = lazy(() => import('./pages/HomePage').then((m) => ({ default: m.HomePage }))) +const LibraryPage = lazy(() => + import('./pages/LibraryPage').then((m) => ({ default: m.LibraryPage })), +) +const SearchPage = lazy(() => + import('./pages/SearchPage').then((m) => ({ default: m.SearchPage })), +) +const FavouritesPage = lazy(() => + import('./pages/FavouritesPage').then((m) => ({ default: m.FavouritesPage })), +) +const PlaylistsPage = lazy(() => + import('./pages/PlaylistsPage').then((m) => ({ default: m.PlaylistsPage })), +) +const PlaylistDetailPage = lazy(() => + import('./pages/PlaylistDetailPage').then((m) => ({ default: m.PlaylistDetailPage })), +) +const MediaDetailPage = lazy(() => + import('./pages/MediaDetailPage').then((m) => ({ default: m.MediaDetailPage })), +) +const PlayerPage = lazy(() => + import('./pages/PlayerPage').then((m) => ({ default: m.PlayerPage })), +) +const AdminPage = lazy(() => import('./pages/AdminPage').then((m) => ({ default: m.AdminPage }))) + +// Fallback shown while a chunk is loading. +const Loading = () =>

加载中…

+ export default function App() { return ( - - } /> - - - - } - > - } /> - } /> - } /> - } /> - } /> + }> + + } /> - - + + + } - /> - - } /> - + > + } /> + } /> + } /> + } /> + } /> + } /> + } /> + } /> + + + + } + /> + + } /> + + ) } diff --git a/web/src/api/client.ts b/web/src/api/client.ts index a6f2777..993187e 100644 --- a/web/src/api/client.ts +++ b/web/src/api/client.ts @@ -31,9 +31,27 @@ api.interceptors.response.use( }, ) -// Helper that returns a streaming URL with the JWT appended as a query -// parameter so