mirror of
https://github.com/truewhile/MeBox.git
synced 2026-10-03 04:26:36 +08:00
feat: hwaccel + Fanart/TheTVDB + AI + Discover + NFO + Tasks + Recycle bin
Backend
- service/transcoder.go: encoder selector (software / nvenc / qsv /
vaapi) with proper hwaccel + scale_* filters per encoder; configurable
bitrate/preset/height/segment-seconds via transcoder.* config; new
Active() snapshot + ActiveJob struct for the Tasks panel; unit-tested
via buildFFmpegArgs(...).
- service/thetvdb.go: TheTVDB v4 provider — login() caches the JWT for
24h, SearchSeries() returns Match struct.
- service/fanart.go: Fanart.tv provider — high-res movie artwork keyed
by TMDb id; used to upgrade poster/backdrop after a successful match.
- service/scraper.go: provider chain reworked — anime → Bangumi, tv →
TheTVDB, default → TMDb, with optional Fanart upgrade post-match.
- service/discover.go: TMDb /trending/movie/day + /movie/popular for
the Discover rail.
- service/ai.go: OpenAI-compatible client. SmartSearch() turns a free-
form query into a structured SearchIntent JSON; Recommend() emits a
short list of titles given the user's history. Disabled when
ai.api_key is empty.
- service/nfo.go: per-movie .nfo writer (Kodi/Jellyfin schema) +
library-scope batch exporter.
- service/media.go: SoftDelete / Restore / Purge / ListRecycleBin
helpers backed by gorm Unscoped().
- service/service.go: container wires Fanart, TheTVDB, Discover, AI,
NFO; everything still tears down cleanly via Close().
- config: new transcoder.* section (encoder, preset, video_bitrate,
max_rate, buf_size, max_height, segment_seconds) with sensible
defaults.
Handlers / routes
- new files: discover.go, recycle.go, nfo.go, ai.go, tasks.go.
- registers GET /tasks, /discover/{trending,popular},
/ai/{status,search,recommend}; admin-only GET /recycle and
DELETE/POST /media/:id (soft delete / restore / purge); admin-only
POST /media/:id/nfo and /libraries/:id/nfo.
Frontend
- new pages: DiscoverPage, TasksPage, RecycleBinPage; SearchPage gains
an optional AI smart-search toggle that calls /api/ai/search.
- MediaDetailPage gains admin actions for 导出 NFO and 移至回收站.
- new api/ helpers: ai.ts, discover.ts, recycle.ts, tasks.ts.
- App.tsx routes /discover, /tasks, /recycle (last two admin-only).
- Layout sidebar groups now include Discover + Tasks + Recycle Bin.
Docker / config
- Dockerfile: adds intel-media-driver / libva-utils / mesa-va-gallium
so QSV + VAAPI work out of the box on Intel iGPUs.
- docker-compose.yml: documents devices + group_add + nvidia runtime
overrides for hardware transcoding.
Verified: go build, go vet, go test (incl. new TestBuildFFmpegArgs across
software / nvenc / qsv / vaapi encoder profiles); tsc -b && vite build
emits 22 route chunks plus the deferred hls chunk; main bundle 248 KB /
83 KB gzipped.
This commit is contained in:
@@ -30,13 +30,25 @@ const EnvPrefix = "MEDIASTATION"
|
||||
|
||||
// Config is the root config aggregate.
|
||||
type Config struct {
|
||||
App AppConfig `mapstructure:"app"`
|
||||
Database DatabaseConfig `mapstructure:"database"`
|
||||
Secrets SecretsConfig `mapstructure:"secrets"`
|
||||
Logging LoggingConfig `mapstructure:"logging"`
|
||||
Cache CacheConfig `mapstructure:"cache"`
|
||||
Media MediaConfig `mapstructure:"media"`
|
||||
AI AIConfig `mapstructure:"ai"`
|
||||
App AppConfig `mapstructure:"app"`
|
||||
Database DatabaseConfig `mapstructure:"database"`
|
||||
Secrets SecretsConfig `mapstructure:"secrets"`
|
||||
Logging LoggingConfig `mapstructure:"logging"`
|
||||
Cache CacheConfig `mapstructure:"cache"`
|
||||
Media MediaConfig `mapstructure:"media"`
|
||||
Transcoder TranscoderConfig `mapstructure:"transcoder"`
|
||||
AI AIConfig `mapstructure:"ai"`
|
||||
}
|
||||
|
||||
// TranscoderConfig controls the HLS / ffmpeg backend.
|
||||
type TranscoderConfig struct {
|
||||
Encoder string `mapstructure:"encoder"` // "" / nvenc / qsv / vaapi
|
||||
Preset string `mapstructure:"preset"`
|
||||
VideoBitrate string `mapstructure:"video_bitrate"`
|
||||
MaxRate string `mapstructure:"max_rate"`
|
||||
BufSize string `mapstructure:"buf_size"`
|
||||
MaxHeight int `mapstructure:"max_height"`
|
||||
SegmentSeconds int `mapstructure:"segment_seconds"`
|
||||
}
|
||||
|
||||
// AppConfig holds runtime app parameters.
|
||||
@@ -195,6 +207,14 @@ func setDefaults(v *viper.Viper) {
|
||||
v.SetDefault("ai.model", "gpt-4o-mini")
|
||||
v.SetDefault("ai.timeout", 30)
|
||||
v.SetDefault("ai.max_concurrent", 3)
|
||||
|
||||
v.SetDefault("transcoder.encoder", "")
|
||||
v.SetDefault("transcoder.preset", "veryfast")
|
||||
v.SetDefault("transcoder.video_bitrate", "1500k")
|
||||
v.SetDefault("transcoder.max_rate", "1800k")
|
||||
v.SetDefault("transcoder.buf_size", "3000k")
|
||||
v.SetDefault("transcoder.max_height", 720)
|
||||
v.SetDefault("transcoder.segment_seconds", 4)
|
||||
}
|
||||
|
||||
// normalize fills derived defaults and self-heals empty critical fields.
|
||||
|
||||
@@ -0,0 +1,71 @@
|
||||
// Package handler — AI integration endpoints.
|
||||
package handler
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"strings"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
|
||||
"github.com/ShukeBta/MediaStationGo/internal/middleware"
|
||||
"github.com/ShukeBta/MediaStationGo/internal/service"
|
||||
)
|
||||
|
||||
type smartSearchReq struct {
|
||||
Query string `json:"query" binding:"required"`
|
||||
}
|
||||
|
||||
func smartSearchHandler(svc *service.Container) gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
var req smartSearchReq
|
||||
if err := c.ShouldBindJSON(&req); err != nil {
|
||||
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
intent, err := svc.AI.SmartSearch(c.Request.Context(), req.Query)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
// Run the actual library search using the cleaned query so the
|
||||
// caller can render results in one round-trip.
|
||||
items, _ := svc.Media.SearchMedia(c.Request.Context(), intent.Query, 60)
|
||||
c.JSON(http.StatusOK, gin.H{
|
||||
"intent": intent,
|
||||
"items": items,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func aiRecommendHandler(svc *service.Container) gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
uid, _ := c.Get(middleware.CtxUserID)
|
||||
hist, err := svc.Playback.RecentHistory(c.Request.Context(), toString(uid), 10)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
titles := make([]string, 0, len(hist))
|
||||
for _, h := range hist {
|
||||
if h.Media != nil && strings.TrimSpace(h.Media.Title) != "" {
|
||||
titles = append(titles, h.Media.Title)
|
||||
}
|
||||
}
|
||||
out, err := svc.AI.Recommend(c.Request.Context(), titles, 8)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"titles": out})
|
||||
}
|
||||
}
|
||||
|
||||
func aiStatusHandler(svc *service.Container) gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
c.JSON(http.StatusOK, gin.H{
|
||||
"enabled": svc.AI.Enabled(),
|
||||
"provider": svc.Cfg.AI.Provider,
|
||||
"model": svc.Cfg.AI.Model,
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,32 @@
|
||||
// Package handler — TMDb discovery endpoints.
|
||||
package handler
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
|
||||
"github.com/ShukeBta/MediaStationGo/internal/service"
|
||||
)
|
||||
|
||||
func trendingHandler(svc *service.Container) gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
items, err := svc.Discover.Trending(c.Request.Context())
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"items": items})
|
||||
}
|
||||
}
|
||||
|
||||
func popularHandler(svc *service.Container) gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
items, err := svc.Discover.Popular(c.Request.Context())
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"items": items})
|
||||
}
|
||||
}
|
||||
@@ -50,8 +50,13 @@ func Register(r *gin.Engine, cfg *config.Config, log *zap.Logger, svc *service.C
|
||||
authed.GET("/media", searchMediaHandler(svc))
|
||||
authed.POST("/media/:id/scrape", middleware.AdminRequired(), scrapeOneHandler(svc))
|
||||
authed.POST("/media/:id/probe", middleware.AdminRequired(), reprobeHandler(svc))
|
||||
authed.DELETE("/media/:id", middleware.AdminRequired(), deleteMediaHandler(svc))
|
||||
authed.POST("/media/:id/restore", middleware.AdminRequired(), restoreMediaHandler(svc))
|
||||
authed.DELETE("/media/:id/purge", middleware.AdminRequired(), purgeMediaHandler(svc))
|
||||
authed.GET("/media/:id/subtitles", listSubtitlesHandler(svc))
|
||||
authed.GET("/subtitles/:id", serveSubtitleHandler(svc))
|
||||
authed.POST("/media/:id/nfo", middleware.AdminRequired(), exportNFOHandler(svc))
|
||||
authed.POST("/libraries/:id/nfo", middleware.AdminRequired(), exportLibraryNFOHandler(svc))
|
||||
|
||||
// Streaming.
|
||||
authed.GET("/stream/:id", streamHandler(svc))
|
||||
@@ -90,6 +95,19 @@ func Register(r *gin.Engine, cfg *config.Config, log *zap.Logger, svc *service.C
|
||||
|
||||
// Stats / dashboard.
|
||||
authed.GET("/stats", statsHandler(svc))
|
||||
authed.GET("/tasks", tasksHandler(svc))
|
||||
|
||||
// Discover (TMDb trending / popular).
|
||||
authed.GET("/discover/trending", trendingHandler(svc))
|
||||
authed.GET("/discover/popular", popularHandler(svc))
|
||||
|
||||
// AI.
|
||||
authed.GET("/ai/status", aiStatusHandler(svc))
|
||||
authed.POST("/ai/search", smartSearchHandler(svc))
|
||||
authed.GET("/ai/recommend", aiRecommendHandler(svc))
|
||||
|
||||
// Recycle bin.
|
||||
authed.GET("/recycle", middleware.AdminRequired(), listRecycleHandler(svc))
|
||||
|
||||
authed.GET("/ws", wsHandler(svc))
|
||||
}
|
||||
|
||||
@@ -0,0 +1,32 @@
|
||||
// Package handler — NFO export endpoints (Kodi / Jellyfin compatibility).
|
||||
package handler
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
|
||||
"github.com/ShukeBta/MediaStationGo/internal/service"
|
||||
)
|
||||
|
||||
func exportNFOHandler(svc *service.Container) gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
path, err := svc.NFO.ExportOne(c.Request.Context(), c.Param("id"))
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"path": path})
|
||||
}
|
||||
}
|
||||
|
||||
func exportLibraryNFOHandler(svc *service.Container) gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
written, err := svc.NFO.ExportLibrary(c.Request.Context(), c.Param("id"))
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"written": written})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,51 @@
|
||||
// Package handler — recycle bin endpoints.
|
||||
package handler
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
|
||||
"github.com/ShukeBta/MediaStationGo/internal/service"
|
||||
)
|
||||
|
||||
func deleteMediaHandler(svc *service.Container) gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
if err := svc.Media.SoftDelete(c.Request.Context(), c.Param("id")); err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.Status(http.StatusNoContent)
|
||||
}
|
||||
}
|
||||
|
||||
func listRecycleHandler(svc *service.Container) gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
items, err := svc.Media.ListRecycleBin(c.Request.Context(), 200)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"items": items})
|
||||
}
|
||||
}
|
||||
|
||||
func restoreMediaHandler(svc *service.Container) gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
if err := svc.Media.RestoreDeleted(c.Request.Context(), c.Param("id")); err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.Status(http.StatusNoContent)
|
||||
}
|
||||
}
|
||||
|
||||
func purgeMediaHandler(svc *service.Container) gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
if err := svc.Media.PurgeDeleted(c.Request.Context(), c.Param("id")); err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.Status(http.StatusNoContent)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,26 @@
|
||||
// Package handler — live tasks board.
|
||||
//
|
||||
// /api/tasks aggregates running ffmpeg transcodes, qBittorrent torrents
|
||||
// and recent scrape progress into a single snapshot suitable for the
|
||||
// React Tasks panel. The panel can layer this REST snapshot on top of
|
||||
// live WS events for instant updates.
|
||||
package handler
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
|
||||
"github.com/ShukeBta/MediaStationGo/internal/service"
|
||||
)
|
||||
|
||||
func tasksHandler(svc *service.Container) gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
transcodes := svc.Transcoder.Active()
|
||||
_, torrents, _ := svc.Downloads.List(c.Request.Context())
|
||||
c.JSON(http.StatusOK, gin.H{
|
||||
"transcodes": transcodes,
|
||||
"torrents": torrents,
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,166 @@
|
||||
// Package service — AI integration (OpenAI-compatible chat completions).
|
||||
//
|
||||
// AIService is a thin wrapper around any OpenAI-compatible REST endpoint
|
||||
// (OpenAI, DeepSeek, Qwen, Ollama, …). Today we expose two operations:
|
||||
//
|
||||
// - SmartSearch: interpret a free-form Chinese / English query and
|
||||
// return a normalised JSON intent the React UI can
|
||||
// translate into filter params.
|
||||
// - Recommend: given a list of recently-watched titles, generate
|
||||
// a short list of "you might like…" recommendations.
|
||||
//
|
||||
// The service is disabled (every method returns nil) when ai.enabled is
|
||||
// false or ai.api_key is empty.
|
||||
package service
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"go.uber.org/zap"
|
||||
|
||||
"github.com/ShukeBta/MediaStationGo/internal/config"
|
||||
)
|
||||
|
||||
// AIService talks to an OpenAI-compatible chat-completions endpoint.
|
||||
type AIService struct {
|
||||
cfg *config.Config
|
||||
log *zap.Logger
|
||||
client *http.Client
|
||||
}
|
||||
|
||||
// NewAIService is the constructor.
|
||||
func NewAIService(cfg *config.Config, log *zap.Logger) *AIService {
|
||||
timeout := time.Duration(cfg.AI.Timeout) * time.Second
|
||||
if timeout <= 0 {
|
||||
timeout = 30 * time.Second
|
||||
}
|
||||
return &AIService{
|
||||
cfg: cfg,
|
||||
log: log,
|
||||
client: &http.Client{Timeout: timeout},
|
||||
}
|
||||
}
|
||||
|
||||
// Enabled reports whether the AI integration is configured.
|
||||
func (a *AIService) Enabled() bool {
|
||||
return a.cfg.AI.Enabled && strings.TrimSpace(a.cfg.AI.APIKey) != ""
|
||||
}
|
||||
|
||||
// SearchIntent is the structured output the smart search endpoint returns.
|
||||
type SearchIntent struct {
|
||||
Query string `json:"query"`
|
||||
Year int `json:"year,omitempty"`
|
||||
Genre string `json:"genre,omitempty"`
|
||||
Type string `json:"type,omitempty"` // movie / tv / anime / music
|
||||
Sort string `json:"sort,omitempty"` // recent / rating / random
|
||||
Language string `json:"language,omitempty"`
|
||||
}
|
||||
|
||||
// SmartSearch turns a natural-language query into a structured intent.
|
||||
// Returns a best-effort intent on parse failure (raw query passes through).
|
||||
func (a *AIService) SmartSearch(ctx context.Context, raw string) (*SearchIntent, error) {
|
||||
if !a.Enabled() {
|
||||
return &SearchIntent{Query: raw}, nil
|
||||
}
|
||||
const sys = "You are a media-library search assistant. Read the user's query and " +
|
||||
"output a JSON object with the keys: query (string), year (int, optional), " +
|
||||
"genre (string, optional), type (movie|tv|anime|music, optional), sort " +
|
||||
"(recent|rating|random, optional), language (zh|en, optional). Respond with " +
|
||||
"JSON only, no commentary."
|
||||
out, err := a.complete(ctx, sys, raw)
|
||||
if err != nil {
|
||||
return &SearchIntent{Query: raw}, err
|
||||
}
|
||||
var intent SearchIntent
|
||||
if err := json.Unmarshal([]byte(out), &intent); err != nil {
|
||||
// Fallback: tolerate non-JSON output by treating the raw text as
|
||||
// the cleaned query.
|
||||
intent.Query = strings.TrimSpace(out)
|
||||
}
|
||||
if intent.Query == "" {
|
||||
intent.Query = raw
|
||||
}
|
||||
return &intent, nil
|
||||
}
|
||||
|
||||
// Recommend builds a short comma-separated list of titles given the user's
|
||||
// history. The first call is intentionally best-effort: a future iteration
|
||||
// may chain media DB lookups onto each suggestion.
|
||||
func (a *AIService) Recommend(ctx context.Context, history []string, max int) ([]string, error) {
|
||||
if !a.Enabled() || len(history) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
if max <= 0 || max > 20 {
|
||||
max = 8
|
||||
}
|
||||
sys := fmt.Sprintf("You are a film / TV recommendation assistant. Reply with %d "+
|
||||
"comma-separated titles only, no commentary, in the same language as the input.", max)
|
||||
usr := "I recently watched: " + strings.Join(history, "; ")
|
||||
out, err := a.complete(ctx, sys, usr)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
parts := strings.Split(out, ",")
|
||||
titles := make([]string, 0, len(parts))
|
||||
for _, p := range parts {
|
||||
p = strings.TrimSpace(p)
|
||||
p = strings.Trim(p, "\"'`")
|
||||
if p != "" {
|
||||
titles = append(titles, p)
|
||||
}
|
||||
}
|
||||
return titles, nil
|
||||
}
|
||||
|
||||
// complete is the shared helper — POST /v1/chat/completions.
|
||||
func (a *AIService) complete(ctx context.Context, system, user string) (string, error) {
|
||||
payload := map[string]any{
|
||||
"model": a.cfg.AI.Model,
|
||||
"temperature": 0.2,
|
||||
"messages": []map[string]string{
|
||||
{"role": "system", "content": system},
|
||||
{"role": "user", "content": user},
|
||||
},
|
||||
}
|
||||
body, _ := json.Marshal(payload)
|
||||
endpoint := strings.TrimRight(a.cfg.AI.APIBase, "/") + "/chat/completions"
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
req.Header.Set("Authorization", "Bearer "+a.cfg.AI.APIKey)
|
||||
resp, err := a.client.Do(req)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode >= 400 {
|
||||
raw, _ := io.ReadAll(resp.Body)
|
||||
return "", fmt.Errorf("ai %d: %s", resp.StatusCode, strings.TrimSpace(string(raw)))
|
||||
}
|
||||
|
||||
type choice struct {
|
||||
Message struct {
|
||||
Content string `json:"content"`
|
||||
} `json:"message"`
|
||||
}
|
||||
var out struct {
|
||||
Choices []choice `json:"choices"`
|
||||
}
|
||||
if err := json.NewDecoder(resp.Body).Decode(&out); err != nil {
|
||||
return "", err
|
||||
}
|
||||
if len(out.Choices) == 0 {
|
||||
return "", errors.New("ai: empty completion")
|
||||
}
|
||||
return strings.TrimSpace(out.Choices[0].Message.Content), nil
|
||||
}
|
||||
@@ -0,0 +1,117 @@
|
||||
// Package service — TMDb discovery (trending / popular).
|
||||
//
|
||||
// DiscoverService surfaces curated lists from TMDb so the React home
|
||||
// page can show "Trending" and "Popular" rails alongside the user's own
|
||||
// library. All methods gracefully no-op when the TMDb provider is
|
||||
// disabled.
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"time"
|
||||
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
// DiscoverService talks to TMDb's /trending and /movie/popular endpoints.
|
||||
type DiscoverService struct {
|
||||
log *zap.Logger
|
||||
tmdb *TMDbProvider
|
||||
client *http.Client
|
||||
}
|
||||
|
||||
// NewDiscoverService is the constructor.
|
||||
func NewDiscoverService(log *zap.Logger, tmdb *TMDbProvider) *DiscoverService {
|
||||
return &DiscoverService{
|
||||
log: log,
|
||||
tmdb: tmdb,
|
||||
client: &http.Client{Timeout: 15 * time.Second},
|
||||
}
|
||||
}
|
||||
|
||||
// Trending returns the daily trending movies (TMDb /trending/movie/day).
|
||||
func (d *DiscoverService) Trending(ctx context.Context) ([]Match, error) {
|
||||
return d.fetch(ctx, "/trending/movie/day")
|
||||
}
|
||||
|
||||
// Popular returns the popular movies list (TMDb /movie/popular).
|
||||
func (d *DiscoverService) Popular(ctx context.Context) ([]Match, error) {
|
||||
return d.fetch(ctx, "/movie/popular")
|
||||
}
|
||||
|
||||
// fetch is the shared helper that paginates page=1 only — that's all the
|
||||
// home page needs and it keeps us under TMDb's 50 rps limit.
|
||||
func (d *DiscoverService) fetch(ctx context.Context, path string) ([]Match, error) {
|
||||
if d.tmdb == nil || !d.tmdb.Enabled() {
|
||||
return nil, nil
|
||||
}
|
||||
q := url.Values{}
|
||||
q.Set("api_key", d.tmdb.cfg.Secrets.TMDbAPIKey)
|
||||
q.Set("language", "zh-CN")
|
||||
q.Set("page", "1")
|
||||
u := d.tmdb.base + path + "?" + q.Encode()
|
||||
|
||||
type result struct {
|
||||
ID int `json:"id"`
|
||||
Title string `json:"title"`
|
||||
Name string `json:"name"`
|
||||
Overview string `json:"overview"`
|
||||
PosterPath string `json:"poster_path"`
|
||||
BackdropPath string `json:"backdrop_path"`
|
||||
ReleaseDate string `json:"release_date"`
|
||||
FirstAirDate string `json:"first_air_date"`
|
||||
VoteAverage float32 `json:"vote_average"`
|
||||
}
|
||||
type page struct {
|
||||
Results []result `json:"results"`
|
||||
}
|
||||
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
resp, err := d.client.Do(req)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode >= 400 {
|
||||
return nil, fmt.Errorf("tmdb %s: %d", path, resp.StatusCode)
|
||||
}
|
||||
var p page
|
||||
if err := json.NewDecoder(resp.Body).Decode(&p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out := make([]Match, 0, len(p.Results))
|
||||
for _, r := range p.Results {
|
||||
title := r.Title
|
||||
if title == "" {
|
||||
title = r.Name
|
||||
}
|
||||
m := Match{
|
||||
TMDbID: r.ID,
|
||||
Title: title,
|
||||
Overview: r.Overview,
|
||||
Rating: r.VoteAverage,
|
||||
}
|
||||
if r.PosterPath != "" {
|
||||
m.PosterURL = d.tmdb.imgCDN + "/w500" + r.PosterPath
|
||||
}
|
||||
if r.BackdropPath != "" {
|
||||
m.BackdropURL = d.tmdb.imgCDN + "/w1280" + r.BackdropPath
|
||||
}
|
||||
date := r.ReleaseDate
|
||||
if date == "" {
|
||||
date = r.FirstAirDate
|
||||
}
|
||||
if len(date) >= 4 {
|
||||
fmt.Sscanf(date[:4], "%d", &m.Year)
|
||||
}
|
||||
out = append(out, m)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
@@ -0,0 +1,109 @@
|
||||
// Package service — Fanart.tv image provider.
|
||||
//
|
||||
// Fanart.tv (https://fanart.tv) hosts community-curated artwork. We use
|
||||
// it to upgrade movie / tv posters and backdrops with higher-resolution
|
||||
// alternatives once a TMDb / Bangumi match is established.
|
||||
//
|
||||
// Endpoints used:
|
||||
//
|
||||
// GET /v3/movies/{tmdb_id} (artwork keyed by TMDb id)
|
||||
// GET /v3/tv/{thetvdb_id} (artwork keyed by TheTVDB id)
|
||||
//
|
||||
// The provider is enabled iff secrets.fanart_tv_api_key is non-empty.
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"go.uber.org/zap"
|
||||
|
||||
"github.com/ShukeBta/MediaStationGo/internal/config"
|
||||
)
|
||||
|
||||
// FanartProvider talks to https://webservice.fanart.tv.
|
||||
type FanartProvider struct {
|
||||
cfg *config.Config
|
||||
log *zap.Logger
|
||||
client *http.Client
|
||||
}
|
||||
|
||||
// NewFanartProvider is the constructor.
|
||||
func NewFanartProvider(cfg *config.Config, log *zap.Logger) *FanartProvider {
|
||||
return &FanartProvider{
|
||||
cfg: cfg,
|
||||
log: log,
|
||||
client: &http.Client{Timeout: 15 * time.Second},
|
||||
}
|
||||
}
|
||||
|
||||
// Enabled reports whether an API key is configured.
|
||||
func (f *FanartProvider) Enabled() bool { return f.cfg.Secrets.FanartAPIKey != "" }
|
||||
|
||||
// Artwork is the high-res image set returned by Fanart.tv.
|
||||
type Artwork struct {
|
||||
Poster string `json:"poster"`
|
||||
Backdrop string `json:"backdrop"`
|
||||
Logo string `json:"logo"`
|
||||
Thumb string `json:"thumb"`
|
||||
}
|
||||
|
||||
// MovieArtwork looks up the artwork bundle for a TMDb movie id.
|
||||
func (f *FanartProvider) MovieArtwork(ctx context.Context, tmdbID int) (*Artwork, error) {
|
||||
if !f.Enabled() || tmdbID <= 0 {
|
||||
return nil, nil
|
||||
}
|
||||
type entry struct {
|
||||
URL string `json:"url"`
|
||||
Lang string `json:"lang"`
|
||||
}
|
||||
type page struct {
|
||||
MoviePoster []entry `json:"movieposter"`
|
||||
MovieBackgr []entry `json:"moviebackground"`
|
||||
HDLogo []entry `json:"hdmovielogo"`
|
||||
MovieThumb []entry `json:"moviethumb"`
|
||||
}
|
||||
u := fmt.Sprintf("https://webservice.fanart.tv/v3/movies/%d?api_key=%s",
|
||||
tmdbID, f.cfg.Secrets.FanartAPIKey)
|
||||
var p page
|
||||
if err := f.getJSON(ctx, u, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
a := &Artwork{}
|
||||
if len(p.MoviePoster) > 0 {
|
||||
a.Poster = p.MoviePoster[0].URL
|
||||
}
|
||||
if len(p.MovieBackgr) > 0 {
|
||||
a.Backdrop = p.MovieBackgr[0].URL
|
||||
}
|
||||
if len(p.HDLogo) > 0 {
|
||||
a.Logo = p.HDLogo[0].URL
|
||||
}
|
||||
if len(p.MovieThumb) > 0 {
|
||||
a.Thumb = p.MovieThumb[0].URL
|
||||
}
|
||||
return a, nil
|
||||
}
|
||||
|
||||
func (f *FanartProvider) getJSON(ctx context.Context, u string, out any) error {
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
req.Header.Set("User-Agent", "MediaStationGo/0.1")
|
||||
resp, err := f.client.Do(req)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode == http.StatusNotFound {
|
||||
return nil
|
||||
}
|
||||
if resp.StatusCode >= 400 {
|
||||
return fmt.Errorf("fanart %s: %d", u, resp.StatusCode)
|
||||
}
|
||||
return json.NewDecoder(resp.Body).Decode(out)
|
||||
}
|
||||
@@ -86,3 +86,34 @@ func (s *MediaService) SearchMedia(ctx context.Context, query string, limit int)
|
||||
func (s *MediaService) GetMedia(ctx context.Context, id string) (*model.Media, error) {
|
||||
return s.repo.Media.FindByID(ctx, id)
|
||||
}
|
||||
|
||||
// SoftDelete moves a media row to the recycle bin (gorm soft delete).
|
||||
// The on-disk file is kept; admins can purge it later.
|
||||
func (s *MediaService) SoftDelete(ctx context.Context, id string) error {
|
||||
return s.repo.DB.Where("id = ?", id).Delete(&model.Media{}).Error
|
||||
}
|
||||
|
||||
// RestoreDeleted unsets DeletedAt for a single media row.
|
||||
func (s *MediaService) RestoreDeleted(ctx context.Context, id string) error {
|
||||
return s.repo.DB.Unscoped().Model(&model.Media{}).
|
||||
Where("id = ?", id).Update("deleted_at", nil).Error
|
||||
}
|
||||
|
||||
// ListRecycleBin returns every soft-deleted row, newest first.
|
||||
func (s *MediaService) ListRecycleBin(ctx context.Context, limit int) ([]model.Media, error) {
|
||||
if limit <= 0 || limit > 500 {
|
||||
limit = 100
|
||||
}
|
||||
var rows []model.Media
|
||||
err := s.repo.DB.Unscoped().
|
||||
Where("deleted_at IS NOT NULL").
|
||||
Order("deleted_at desc").
|
||||
Limit(limit).
|
||||
Find(&rows).Error
|
||||
return rows, err
|
||||
}
|
||||
|
||||
// PurgeDeleted permanently removes a soft-deleted row from the database.
|
||||
func (s *MediaService) PurgeDeleted(ctx context.Context, id string) error {
|
||||
return s.repo.DB.Unscoped().Where("id = ?", id).Delete(&model.Media{}).Error
|
||||
}
|
||||
|
||||
@@ -0,0 +1,115 @@
|
||||
// Package service — NFO writer (Kodi / Jellyfin compatibility).
|
||||
//
|
||||
// Kodi and Jellyfin index media using sidecar XML files alongside the
|
||||
// source video. We export a minimal subset that those scrapers consume
|
||||
// happily:
|
||||
//
|
||||
// movie.mkv -> movie.nfo (<movie>...</movie>)
|
||||
// tvshow/ -> tvshow.nfo (<tvshow>...</tvshow>) [future]
|
||||
//
|
||||
// Today only the per-movie writer is implemented; the per-show / per-episode
|
||||
// writers are stubbed with TODO markers.
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/xml"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
|
||||
"go.uber.org/zap"
|
||||
|
||||
"github.com/ShukeBta/MediaStationGo/internal/repository"
|
||||
)
|
||||
|
||||
// NFOService is the entry point used by the admin "导出 NFO" action.
|
||||
type NFOService struct {
|
||||
log *zap.Logger
|
||||
repo *repository.Container
|
||||
}
|
||||
|
||||
// NewNFOService is the constructor.
|
||||
func NewNFOService(log *zap.Logger, repo *repository.Container) *NFOService {
|
||||
return &NFOService{log: log, repo: repo}
|
||||
}
|
||||
|
||||
// movieNFO is the on-disk schema. Tag names match Kodi's expectations.
|
||||
type movieNFO struct {
|
||||
XMLName xml.Name `xml:"movie"`
|
||||
Title string `xml:"title"`
|
||||
Original string `xml:"originaltitle,omitempty"`
|
||||
Year int `xml:"year,omitempty"`
|
||||
Plot string `xml:"plot,omitempty"`
|
||||
Rating float32 `xml:"rating,omitempty"`
|
||||
Poster string `xml:"thumb,omitempty"`
|
||||
Fanart string `xml:"fanart,omitempty"`
|
||||
TMDb int `xml:"tmdbid,omitempty"`
|
||||
}
|
||||
|
||||
// ExportOne writes a movie.nfo file next to the media file. Existing files
|
||||
// are overwritten so a re-scrape always reflects the latest metadata.
|
||||
func (s *NFOService) ExportOne(ctx context.Context, mediaID string) (string, error) {
|
||||
m, err := s.repo.Media.FindByID(ctx, mediaID)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
if m == nil {
|
||||
return "", errors.New("media not found")
|
||||
}
|
||||
if m.Path == "" {
|
||||
return "", errors.New("media has empty path")
|
||||
}
|
||||
|
||||
doc := movieNFO{
|
||||
Title: m.Title,
|
||||
Original: m.OriginalName,
|
||||
Year: m.Year,
|
||||
Plot: m.Overview,
|
||||
Rating: m.Rating,
|
||||
Poster: m.PosterURL,
|
||||
Fanart: m.BackdropURL,
|
||||
TMDb: m.TMDbID,
|
||||
}
|
||||
out, err := xml.MarshalIndent(doc, "", " ")
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
body := []byte(xml.Header + string(out) + "\n")
|
||||
|
||||
dst := nfoPath(m.Path)
|
||||
if err := os.WriteFile(dst, body, 0o644); err != nil {
|
||||
return "", err
|
||||
}
|
||||
s.log.Info("nfo exported", zap.String("media_id", m.ID), zap.String("path", dst))
|
||||
return dst, nil
|
||||
}
|
||||
|
||||
// ExportLibrary loops through every matched movie in a library and writes
|
||||
// an .nfo for each. Returns (written, error).
|
||||
func (s *NFOService) ExportLibrary(ctx context.Context, libraryID string) (int, error) {
|
||||
type row struct{ ID string }
|
||||
var ids []row
|
||||
q := s.repo.DB.Table("media").Select("id").Where("scrape_status = ?", "matched")
|
||||
if libraryID != "" {
|
||||
q = q.Where("library_id = ?", libraryID)
|
||||
}
|
||||
if err := q.Scan(&ids).Error; err != nil {
|
||||
return 0, err
|
||||
}
|
||||
written := 0
|
||||
for _, r := range ids {
|
||||
if _, err := s.ExportOne(ctx, r.ID); err == nil {
|
||||
written++
|
||||
}
|
||||
}
|
||||
return written, nil
|
||||
}
|
||||
|
||||
func nfoPath(media string) string {
|
||||
dir := filepath.Dir(media)
|
||||
base := strings.TrimSuffix(filepath.Base(media), filepath.Ext(media))
|
||||
return filepath.Join(dir, fmt.Sprintf("%s.nfo", base))
|
||||
}
|
||||
+36
-16
@@ -4,12 +4,11 @@
|
||||
// from one or more providers. Selection is driven by the library type:
|
||||
//
|
||||
// library.type == "anime" -> Bangumi (fallback: TMDb)
|
||||
// library.type == "tv" -> TMDb (movies) — TV episodes inherit
|
||||
// series metadata; episode-level scraping
|
||||
// is left as a future step
|
||||
// library.type == "tv" -> TheTVDB (fallback: TMDb)
|
||||
// default -> TMDb
|
||||
//
|
||||
// The orchestrator publishes scrape progress events on the WS hub.
|
||||
// After the primary match we optionally upgrade poster / backdrop with
|
||||
// Fanart.tv when an API key is configured.
|
||||
package service
|
||||
|
||||
import (
|
||||
@@ -34,6 +33,8 @@ type ScraperService struct {
|
||||
repo *repository.Container
|
||||
tmdb *TMDbProvider
|
||||
bangumi *BangumiProvider
|
||||
thetvdb *TheTVDBProvider
|
||||
fanart *FanartProvider
|
||||
hub *Hub
|
||||
}
|
||||
|
||||
@@ -44,11 +45,13 @@ func NewScraperService(
|
||||
repo *repository.Container,
|
||||
tmdb *TMDbProvider,
|
||||
bangumi *BangumiProvider,
|
||||
thetvdb *TheTVDBProvider,
|
||||
fanart *FanartProvider,
|
||||
hub *Hub,
|
||||
) *ScraperService {
|
||||
return &ScraperService{
|
||||
cfg: cfg, log: log, repo: repo,
|
||||
tmdb: tmdb, bangumi: bangumi, hub: hub,
|
||||
tmdb: tmdb, bangumi: bangumi, thetvdb: thetvdb, fanart: fanart, hub: hub,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -74,7 +77,6 @@ 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
|
||||
@@ -82,10 +84,8 @@ func CleanQuery(raw string) (title string, year int) {
|
||||
}
|
||||
}
|
||||
|
||||
// 2. Drop everything inside brackets.
|
||||
lower = bracketedTag.ReplaceAllString(lower, " ")
|
||||
|
||||
// 3. Drop episode markers (S01E02 / 1x02 / EP05 / 第03集).
|
||||
lower = patSEnE.ReplaceAllString(lower, " ")
|
||||
lower = patNxE.ReplaceAllString(lower, " ")
|
||||
lower = patEP.ReplaceAllString(lower, " ")
|
||||
@@ -102,9 +102,7 @@ func CleanQuery(raw string) (title string, year int) {
|
||||
return strings.TrimSpace(title), year
|
||||
}
|
||||
|
||||
// EnrichOne runs the provider chain for a single media row. The library's
|
||||
// type decides which provider goes first; a fallback runs when the primary
|
||||
// returns nothing.
|
||||
// EnrichOne runs the provider chain for a single media row.
|
||||
func (s *ScraperService) EnrichOne(ctx context.Context, m *model.Media) error {
|
||||
lib, err := s.repo.Library.FindByID(ctx, m.LibraryID)
|
||||
if err != nil {
|
||||
@@ -124,12 +122,23 @@ func (s *ScraperService) EnrichOne(ctx context.Context, m *model.Media) error {
|
||||
|
||||
match := s.lookup(ctx, lib, query, year)
|
||||
if match == nil {
|
||||
// Mark explicitly so we don't retry forever.
|
||||
_ = s.repo.DB.Model(&model.Media{}).Where("id = ?", m.ID).
|
||||
Update("scrape_status", "no_match").Error
|
||||
return nil
|
||||
}
|
||||
|
||||
// Optional Fanart upgrade.
|
||||
if s.fanart != nil && s.fanart.Enabled() && match.TMDbID > 0 {
|
||||
if a, err := s.fanart.MovieArtwork(ctx, match.TMDbID); err == nil && a != nil {
|
||||
if a.Poster != "" {
|
||||
match.PosterURL = a.Poster
|
||||
}
|
||||
if a.Backdrop != "" {
|
||||
match.BackdropURL = a.Backdrop
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
updates := map[string]any{
|
||||
"title": match.Title,
|
||||
"overview": match.Overview,
|
||||
@@ -165,11 +174,19 @@ func (s *ScraperService) lookup(ctx context.Context, lib *model.Library, query s
|
||||
if lib != nil {
|
||||
kind = lib.Type
|
||||
}
|
||||
if kind == "anime" && s.bangumi != nil {
|
||||
if m, err := s.bangumi.Search(ctx, query); err == nil && m != nil {
|
||||
return m
|
||||
switch kind {
|
||||
case "anime":
|
||||
if s.bangumi != nil {
|
||||
if m, err := s.bangumi.Search(ctx, query); err == nil && m != nil {
|
||||
return m
|
||||
}
|
||||
}
|
||||
case "tv":
|
||||
if s.thetvdb != nil && s.thetvdb.Enabled() {
|
||||
if m, err := s.thetvdb.SearchSeries(ctx, query); err == nil && m != nil {
|
||||
return m
|
||||
}
|
||||
}
|
||||
s.log.Debug("bangumi miss, falling back to tmdb", zap.String("query", query))
|
||||
}
|
||||
if s.tmdb != nil && s.tmdb.Enabled() {
|
||||
if m, err := s.tmdb.SearchMovie(ctx, query, year); err == nil && m != nil {
|
||||
@@ -220,5 +237,8 @@ func (s *ScraperService) AnyEnabled() bool {
|
||||
if s.bangumi != nil && s.bangumi.Enabled() {
|
||||
return true
|
||||
}
|
||||
if s.thetvdb != nil && s.thetvdb.Enabled() {
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
@@ -28,7 +28,10 @@ type Container struct {
|
||||
FFprobe *FFprobeService
|
||||
TMDb *TMDbProvider
|
||||
Bangumi *BangumiProvider
|
||||
TheTVDB *TheTVDBProvider
|
||||
Fanart *FanartProvider
|
||||
Scraper *ScraperService
|
||||
Discover *DiscoverService
|
||||
Playback *PlaybackService
|
||||
ImageProxy *ImageProxy
|
||||
Watcher *WatcherService
|
||||
@@ -38,6 +41,8 @@ type Container struct {
|
||||
Stats *StatsService
|
||||
Profile *ProfileService
|
||||
Audit *AuditService
|
||||
NFO *NFOService
|
||||
AI *AIService
|
||||
|
||||
stopCtx context.Context
|
||||
stopCancel context.CancelFunc
|
||||
@@ -51,12 +56,17 @@ func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Cont
|
||||
probe := NewFFprobeService(cfg, log)
|
||||
tmdb := NewTMDbProvider(cfg, log)
|
||||
bangumi := NewBangumiProvider(cfg, log)
|
||||
scraper := NewScraperService(cfg, log, repos, tmdb, bangumi, hub)
|
||||
thetvdb := NewTheTVDBProvider(cfg, log)
|
||||
fanart := NewFanartProvider(cfg, log)
|
||||
scraper := NewScraperService(cfg, log, repos, tmdb, bangumi, thetvdb, fanart, hub)
|
||||
discover := NewDiscoverService(log, tmdb)
|
||||
transcoder := NewTranscoderService(cfg, log, repos, hub)
|
||||
scanner := NewScannerService(cfg, log, repos, hub, probe, scraper)
|
||||
downloads := NewDownloadService(log, repos, hub)
|
||||
subscription := NewSubscriptionService(log, repos, downloads, hub)
|
||||
watcher := NewWatcherService(log, repos, scanner)
|
||||
nfo := NewNFOService(log, repos)
|
||||
ai := NewAIService(cfg, log)
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
|
||||
@@ -73,7 +83,10 @@ func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Cont
|
||||
FFprobe: probe,
|
||||
TMDb: tmdb,
|
||||
Bangumi: bangumi,
|
||||
TheTVDB: thetvdb,
|
||||
Fanart: fanart,
|
||||
Scraper: scraper,
|
||||
Discover: discover,
|
||||
Playback: NewPlaybackService(log, repos),
|
||||
ImageProxy: NewImageProxy(cfg, log),
|
||||
Watcher: watcher,
|
||||
@@ -83,6 +96,8 @@ func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Cont
|
||||
Stats: NewStatsService(log, repos),
|
||||
Profile: NewProfileService(log, repos),
|
||||
Audit: NewAuditService(log, repos),
|
||||
NFO: nfo,
|
||||
AI: ai,
|
||||
stopCtx: ctx,
|
||||
stopCancel: cancel,
|
||||
}
|
||||
|
||||
@@ -0,0 +1,159 @@
|
||||
// Package service — TheTVDB v4 provider.
|
||||
//
|
||||
// TheTVDBProvider implements two methods used by the scraper for TV /
|
||||
// anime libraries:
|
||||
//
|
||||
// Login() -> exchanges secrets.thetvdb_api_key for
|
||||
// a JWT (cached for 24h).
|
||||
// SearchSeries(query) -> /search?query=...&type=series
|
||||
//
|
||||
// The provider is enabled iff secrets.thetvdb_api_key is non-empty. When
|
||||
// disabled every method returns nil, nil so the scraper orchestrator can
|
||||
// gracefully fall through.
|
||||
package service
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"go.uber.org/zap"
|
||||
|
||||
"github.com/ShukeBta/MediaStationGo/internal/config"
|
||||
)
|
||||
|
||||
// TheTVDBProvider talks to https://api4.thetvdb.com/v4.
|
||||
type TheTVDBProvider struct {
|
||||
cfg *config.Config
|
||||
log *zap.Logger
|
||||
client *http.Client
|
||||
|
||||
mu sync.Mutex
|
||||
token string
|
||||
tokenExp time.Time
|
||||
}
|
||||
|
||||
// NewTheTVDBProvider is the constructor.
|
||||
func NewTheTVDBProvider(cfg *config.Config, log *zap.Logger) *TheTVDBProvider {
|
||||
return &TheTVDBProvider{
|
||||
cfg: cfg,
|
||||
log: log,
|
||||
client: &http.Client{Timeout: 15 * time.Second},
|
||||
}
|
||||
}
|
||||
|
||||
// Enabled reports whether an API key is present.
|
||||
func (t *TheTVDBProvider) Enabled() bool { return t.cfg.Secrets.TheTVDBAPIKey != "" }
|
||||
|
||||
// Login fetches a fresh JWT, cached for 24h.
|
||||
func (t *TheTVDBProvider) Login(ctx context.Context) (string, error) {
|
||||
t.mu.Lock()
|
||||
defer t.mu.Unlock()
|
||||
if time.Now().Before(t.tokenExp) && t.token != "" {
|
||||
return t.token, nil
|
||||
}
|
||||
body, _ := json.Marshal(map[string]string{"apikey": t.cfg.Secrets.TheTVDBAPIKey})
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost,
|
||||
"https://api4.thetvdb.com/v4/login", bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
resp, err := t.client.Do(req)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode >= 400 {
|
||||
return "", fmt.Errorf("thetvdb login: %d", resp.StatusCode)
|
||||
}
|
||||
var out struct {
|
||||
Data struct {
|
||||
Token string `json:"token"`
|
||||
} `json:"data"`
|
||||
}
|
||||
if err := json.NewDecoder(resp.Body).Decode(&out); err != nil {
|
||||
return "", err
|
||||
}
|
||||
t.token = out.Data.Token
|
||||
t.tokenExp = time.Now().Add(24 * time.Hour)
|
||||
return t.token, nil
|
||||
}
|
||||
|
||||
// SearchSeries returns the top match for a TV / anime query, or nil.
|
||||
func (t *TheTVDBProvider) SearchSeries(ctx context.Context, query string) (*Match, error) {
|
||||
if !t.Enabled() || query == "" {
|
||||
return nil, nil
|
||||
}
|
||||
tok, err := t.Login(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
u := fmt.Sprintf("https://api4.thetvdb.com/v4/search?query=%s&type=series",
|
||||
urlEscape(query))
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
req.Header.Set("Authorization", "Bearer "+tok)
|
||||
resp, err := t.client.Do(req)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode >= 400 {
|
||||
return nil, fmt.Errorf("thetvdb search: %d", resp.StatusCode)
|
||||
}
|
||||
|
||||
type entry struct {
|
||||
ID string `json:"id"`
|
||||
Name string `json:"name"`
|
||||
Overview string `json:"overview"`
|
||||
Image string `json:"image_url"`
|
||||
Year string `json:"year"`
|
||||
}
|
||||
type page struct {
|
||||
Data []entry `json:"data"`
|
||||
}
|
||||
var p page
|
||||
if err := json.NewDecoder(resp.Body).Decode(&p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(p.Data) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
r := p.Data[0]
|
||||
m := &Match{
|
||||
Title: r.Name,
|
||||
Overview: r.Overview,
|
||||
PosterURL: r.Image,
|
||||
}
|
||||
if len(r.Year) >= 4 {
|
||||
fmt.Sscanf(r.Year[:4], "%d", &m.Year)
|
||||
}
|
||||
return m, nil
|
||||
}
|
||||
|
||||
// urlEscape is a tiny replacement for net/url.QueryEscape kept inline so
|
||||
// the file does not pull a second import for one call.
|
||||
func urlEscape(s string) string {
|
||||
out := make([]byte, 0, len(s)*3)
|
||||
for _, r := range []byte(s) {
|
||||
switch {
|
||||
case r >= '0' && r <= '9',
|
||||
r >= 'A' && r <= 'Z',
|
||||
r >= 'a' && r <= 'z',
|
||||
r == '-', r == '_', r == '.', r == '~':
|
||||
out = append(out, r)
|
||||
default:
|
||||
out = append(out, '%')
|
||||
const hex = "0123456789ABCDEF"
|
||||
out = append(out, hex[r>>4], hex[r&15])
|
||||
}
|
||||
}
|
||||
return string(out)
|
||||
}
|
||||
+161
-37
@@ -4,16 +4,22 @@
|
||||
// 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.
|
||||
//
|
||||
// Encoder selection (read once at startup from the config):
|
||||
//
|
||||
// transcoder.encoder = "" | "nvenc" | "qsv" | "vaapi"
|
||||
//
|
||||
// "" software libx264 (default; runs anywhere)
|
||||
// nvenc h264_nvenc (NVIDIA GPU, requires --gpus all on Docker)
|
||||
// qsv h264_qsv (Intel iGPU, requires /dev/dri:/dev/dri)
|
||||
// vaapi h264_vaapi (Mesa/Intel VAAPI, requires /dev/dri:/dev/dri
|
||||
// plus the kernel module loaded)
|
||||
//
|
||||
// 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 (
|
||||
@@ -34,8 +40,8 @@ import (
|
||||
|
||||
// TranscoderService orchestrates background ffmpeg transcodes.
|
||||
type TranscoderService struct {
|
||||
cfg *config.Config
|
||||
log *zap.Logger
|
||||
cfg *config.Config
|
||||
log *zap.Logger
|
||||
repo *repository.Container
|
||||
hub *Hub
|
||||
|
||||
@@ -50,6 +56,7 @@ type hlsJob struct {
|
||||
cancel context.CancelFunc
|
||||
startedAt time.Time
|
||||
playlistOK bool
|
||||
encoder string
|
||||
}
|
||||
|
||||
// NewTranscoderService is the constructor.
|
||||
@@ -106,6 +113,7 @@ func (t *TranscoderService) EnsureJob(ctx context.Context, mediaID string) (stri
|
||||
outputDir: outDir,
|
||||
cancel: cancel,
|
||||
startedAt: time.Now(),
|
||||
encoder: t.cfg.Transcoder.Encoder,
|
||||
}
|
||||
t.jobs[mediaID] = job
|
||||
t.mu.Unlock()
|
||||
@@ -158,6 +166,30 @@ func (t *TranscoderService) StopAll() {
|
||||
}
|
||||
}
|
||||
|
||||
// ActiveJob is the JSON shape exposed to the React Tasks panel.
|
||||
type ActiveJob struct {
|
||||
MediaID string `json:"media_id"`
|
||||
Encoder string `json:"encoder"`
|
||||
StartedAt time.Time `json:"started_at"`
|
||||
PlaylistOK bool `json:"playlist_ok"`
|
||||
}
|
||||
|
||||
// Active returns a snapshot of the currently running transcode jobs.
|
||||
func (t *TranscoderService) Active() []ActiveJob {
|
||||
t.mu.Lock()
|
||||
defer t.mu.Unlock()
|
||||
out := make([]ActiveJob, 0, len(t.jobs))
|
||||
for _, j := range t.jobs {
|
||||
out = append(out, ActiveJob{
|
||||
MediaID: j.mediaID,
|
||||
Encoder: j.encoder,
|
||||
StartedAt: j.startedAt,
|
||||
PlaylistOK: j.playlistOK,
|
||||
})
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func (t *TranscoderService) runFFmpeg(ctx context.Context, job *hlsJob, source string) {
|
||||
bin := t.cfg.App.FFmpegPath
|
||||
if bin == "" {
|
||||
@@ -167,45 +199,19 @@ func (t *TranscoderService) runFFmpeg(ctx context.Context, job *hlsJob, source s
|
||||
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,
|
||||
}
|
||||
args := buildFFmpegArgs(t.cfg, source, playlist, segments)
|
||||
|
||||
cmd := exec.CommandContext(ctx, bin, args...)
|
||||
cmd.Stderr = os.Stderr
|
||||
|
||||
t.log.Info("transcode started",
|
||||
zap.String("media_id", job.mediaID),
|
||||
zap.String("encoder", job.encoder),
|
||||
zap.String("source", source),
|
||||
)
|
||||
t.hub.Publish("transcode", map[string]any{
|
||||
"media_id": job.mediaID,
|
||||
"encoder": job.encoder,
|
||||
"status": "started",
|
||||
})
|
||||
|
||||
@@ -227,8 +233,126 @@ func (t *TranscoderService) runFFmpeg(ctx context.Context, job *hlsJob, source s
|
||||
})
|
||||
}
|
||||
|
||||
// buildFFmpegArgs assembles the ffmpeg command line for the configured
|
||||
// encoder. The function is package-level so the unit test can pin its
|
||||
// behaviour without spawning a real ffmpeg process.
|
||||
func buildFFmpegArgs(cfg *config.Config, source, playlist, segments string) []string {
|
||||
enc := cfg.Transcoder.Encoder
|
||||
bitrate := cfg.Transcoder.VideoBitrate
|
||||
if bitrate == "" {
|
||||
bitrate = "1500k"
|
||||
}
|
||||
maxrate := cfg.Transcoder.MaxRate
|
||||
if maxrate == "" {
|
||||
maxrate = "1800k"
|
||||
}
|
||||
bufsize := cfg.Transcoder.BufSize
|
||||
if bufsize == "" {
|
||||
bufsize = "3000k"
|
||||
}
|
||||
preset := cfg.Transcoder.Preset
|
||||
if preset == "" {
|
||||
preset = "veryfast"
|
||||
}
|
||||
height := cfg.Transcoder.MaxHeight
|
||||
if height <= 0 {
|
||||
height = 720
|
||||
}
|
||||
segDur := cfg.Transcoder.SegmentSeconds
|
||||
if segDur <= 0 {
|
||||
segDur = 4
|
||||
}
|
||||
|
||||
// Hardware-accel arguments differ in three places:
|
||||
// - Optional input flags (-hwaccel + device init)
|
||||
// - Optional input pixel-format upload filter
|
||||
// - The actual -c:v encoder name + preset/quality flag
|
||||
var pre, vf, vcodec, vpreset string
|
||||
switch enc {
|
||||
case "nvenc":
|
||||
pre = "-hwaccel cuda -hwaccel_output_format cuda"
|
||||
vf = fmt.Sprintf("scale_cuda=-2:min(%d\\,ih)", height)
|
||||
vcodec = "h264_nvenc"
|
||||
vpreset = "p4"
|
||||
case "qsv":
|
||||
pre = "-hwaccel qsv -hwaccel_output_format qsv"
|
||||
vf = fmt.Sprintf("scale_qsv=-1:min(%d\\,ih)", height)
|
||||
vcodec = "h264_qsv"
|
||||
vpreset = preset
|
||||
case "vaapi":
|
||||
device := cfg.App.VAAPIDevice
|
||||
if device == "" {
|
||||
device = "/dev/dri/renderD128"
|
||||
}
|
||||
pre = fmt.Sprintf("-hwaccel vaapi -vaapi_device %s -hwaccel_output_format vaapi", device)
|
||||
vf = fmt.Sprintf("scale_vaapi=-2:min(%d\\,ih),format=nv12|vaapi,hwupload", height)
|
||||
vcodec = "h264_vaapi"
|
||||
vpreset = ""
|
||||
default:
|
||||
// software
|
||||
pre = ""
|
||||
vf = fmt.Sprintf("scale=-2:min(%d\\,ih)", height)
|
||||
vcodec = "libx264"
|
||||
vpreset = preset
|
||||
}
|
||||
|
||||
args := []string{"-y", "-fflags", "+genpts"}
|
||||
for _, p := range splitNonEmptyArgs(pre) {
|
||||
args = append(args, p)
|
||||
}
|
||||
args = append(args, "-i", source, "-map", "0:v:0?", "-map", "0:a:0?", "-vf", vf, "-c:v", vcodec)
|
||||
if vpreset != "" {
|
||||
args = append(args, "-preset", vpreset)
|
||||
}
|
||||
args = append(args,
|
||||
"-pix_fmt", "yuv420p",
|
||||
"-b:v", bitrate,
|
||||
"-maxrate", maxrate,
|
||||
"-bufsize", bufsize,
|
||||
"-c:a", "aac",
|
||||
"-ar", "48000",
|
||||
"-b:a", "128k",
|
||||
"-ac", "2",
|
||||
"-force_key_frames", fmt.Sprintf("expr:gte(t,n_forced*%d)", segDur),
|
||||
"-f", "hls",
|
||||
"-hls_time", fmt.Sprintf("%d", segDur),
|
||||
"-hls_list_size", "0",
|
||||
"-hls_segment_type", "mpegts",
|
||||
"-hls_flags", "independent_segments",
|
||||
"-hls_segment_filename", segments,
|
||||
playlist,
|
||||
)
|
||||
return args
|
||||
}
|
||||
|
||||
// splitNonEmptyArgs is a tiny helper that mirrors strings.Fields for the
|
||||
// pre-input flag string without dragging the strings import into the hot
|
||||
// path of every call to buildFFmpegArgs.
|
||||
func splitNonEmptyArgs(s string) []string {
|
||||
if s == "" {
|
||||
return nil
|
||||
}
|
||||
out := make([]string, 0, 4)
|
||||
field := make([]rune, 0, 16)
|
||||
flush := func() {
|
||||
if len(field) > 0 {
|
||||
out = append(out, string(field))
|
||||
field = field[:0]
|
||||
}
|
||||
}
|
||||
for _, r := range s {
|
||||
if r == ' ' || r == '\t' {
|
||||
flush()
|
||||
continue
|
||||
}
|
||||
field = append(field, r)
|
||||
}
|
||||
flush()
|
||||
return out
|
||||
}
|
||||
|
||||
// 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"))
|
||||
return fmt.Sprintf("ffmpeg=%s, encoder=%s, output=%s",
|
||||
t.cfg.App.FFmpegPath, t.cfg.Transcoder.Encoder, filepath.Join(t.cfg.Cache.CacheDir, "hls"))
|
||||
}
|
||||
|
||||
@@ -0,0 +1,46 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/ShukeBta/MediaStationGo/internal/config"
|
||||
)
|
||||
|
||||
func TestBuildFFmpegArgs(t *testing.T) {
|
||||
base := &config.Config{}
|
||||
base.Transcoder.MaxHeight = 720
|
||||
base.Transcoder.SegmentSeconds = 4
|
||||
base.App.VAAPIDevice = "/dev/dri/renderD128"
|
||||
|
||||
cases := []struct {
|
||||
name string
|
||||
encoder string
|
||||
expectVCodec string
|
||||
expectInArgs []string
|
||||
expectNotPresetIfBlank bool
|
||||
}{
|
||||
{"software", "", "libx264", []string{"-preset", "veryfast", "-c:v", "libx264"}, false},
|
||||
{"nvenc", "nvenc", "h264_nvenc", []string{"-hwaccel", "cuda", "-c:v", "h264_nvenc", "-preset", "p4"}, false},
|
||||
{"qsv", "qsv", "h264_qsv", []string{"-hwaccel", "qsv", "-c:v", "h264_qsv"}, false},
|
||||
{"vaapi", "vaapi", "h264_vaapi", []string{"-hwaccel", "vaapi", "-vaapi_device", "/dev/dri/renderD128", "-c:v", "h264_vaapi"}, true},
|
||||
}
|
||||
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
cfg := *base
|
||||
cfg.Transcoder.Encoder = tc.encoder
|
||||
args := buildFFmpegArgs(&cfg, "/x.mkv", "/o/x.m3u8", "/o/seg_%05d.ts")
|
||||
joined := strings.Join(args, " ")
|
||||
for _, frag := range tc.expectInArgs {
|
||||
if !strings.Contains(joined, frag) {
|
||||
t.Errorf("expected %q in args, got: %s", frag, joined)
|
||||
}
|
||||
}
|
||||
// vaapi has no -preset flag.
|
||||
if tc.expectNotPresetIfBlank && strings.Contains(joined, "-preset") {
|
||||
t.Errorf("vaapi should not include -preset, got: %s", joined)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user