wip: jellyfin sync

This commit is contained in:
IRHM
2024-03-03 22:53:05 +00:00
parent 8e0fb190d1
commit ccadcf2152
6 changed files with 478 additions and 14 deletions
+45 -12
View File
@@ -8,6 +8,9 @@ import (
"log/slog"
"net/http"
"net/url"
"time"
"github.com/gin-gonic/gin"
)
type JellyfinItemSearchResponse struct {
@@ -15,12 +18,27 @@ type JellyfinItemSearchResponse struct {
}
type JellyfinItems struct {
Name string `json:"Name"`
Type string `json:"Type"`
ServerID string `json:"ServerId"`
Id string `json:"Id"`
ProviderIds struct {
Tmdb string `json:"Tmdb"`
} `json:"ProviderIds"`
UserData struct {
Rating float64 `json:"Rating"`
PlayedPercentage float64 `json:"PlayedPercentage"`
UnplayedItemCount int64 `json:"UnplayedItemCount"`
PlaybackPositionTicks int64 `json:"PlaybackPositionTicks"`
PlayCount int64 `json:"PlayCount"`
IsFavorite bool `json:"IsFavorite"`
Likes bool `json:"Likes"`
LastPlayedDate time.Time `json:"LastPlayedDate"`
Played bool `json:"Played"`
Key string `json:"Key"`
ItemId string `json:"ItemId"`
} `json:"UserData"`
RecursiveItemCount int64 `json:"RecursiveItemCount"`
}
type JFContentFindResponse struct {
@@ -28,6 +46,33 @@ type JFContentFindResponse struct {
Url string `json:"url"`
}
// Jellyfin access middleware, ensures user is a jellyfin user.
// To be ran after AuthRequired middleware with extra data.
func JellyfinAccessRequired() gin.HandlerFunc {
return func(c *gin.Context) {
userId := c.MustGet("userId").(uint)
slog.Debug("JellyfinAccessRequired middleware hit", "user_id", userId)
userType := c.MustGet("userType").(UserType)
userThirdPartyId := c.MustGet("userThirdPartyId").(string)
userThirdPartyAuth := c.MustGet("userThirdPartyAuth").(string)
if Config.JELLYFIN_HOST == "" {
slog.Error("JellyfinAccessRequired: Request made to login via Jellyfin, but JELLYFIN_HOST has not been configured.")
c.AbortWithStatus(401)
return
}
if userType != JELLYFIN_USER || userThirdPartyId == "" {
slog.Error("JellyfinAccessRequired: User is not a jellyfin user..", "user_type", userType, "user_third_party_id", userThirdPartyId)
c.AbortWithStatus(401)
return
}
if userThirdPartyAuth == "" {
slog.Error("JellyfinAccessRequired: User has no thirdPartyAuth token..")
c.AbortWithStatus(401)
return
}
}
}
func jellyfinAPIRequest(method string, ep string, p map[string]string, username string, userToken string, resp interface{}) error {
if Config.JELLYFIN_HOST == "" {
slog.Error("jellyfinAPIRequest: JELLYFIN_HOST not configured.")
@@ -97,18 +142,6 @@ func jellyfinContentFind(
contentName string,
contentTmdbId string,
) (JFContentFindResponse, error) {
if Config.JELLYFIN_HOST == "" {
slog.Error("Request made to login via Jellyfin, but JELLYFIN_HOST has not been configured.")
return JFContentFindResponse{}, errors.New("jellyfin login not enabled")
}
if userType != JELLYFIN_USER || userThirdPartyId == "" {
slog.Error("User is not a jellyfin user..", "user_type", userType, "user_third_party_id", userThirdPartyId)
return JFContentFindResponse{}, errors.New("not jellyfin user")
}
if userThirdPartyAuth == "" {
slog.Error("User has no thirdPartyAuth token..")
return JFContentFindResponse{}, errors.New("user has no jellyfin auth token")
}
if contentType == "" || contentName == "" {
slog.Error("Bad request", "content_type", contentType, "content_name", contentName)
return JFContentFindResponse{}, errors.New("content type or name not provided")
+296
View File
@@ -0,0 +1,296 @@
package main
import (
"errors"
"log/slog"
"strconv"
"gorm.io/gorm"
)
type JellyfinSeriesSeasonsResponse struct {
Items []JellyfinSeriesSeasonItem `json:"Items"`
}
type JellyfinSeriesSeasonItem struct {
JellyfinItems
// aka the season number
IndexNumber int `json:"IndexNumber"`
}
type JellyfinSeriesEpisodesResponse struct {
Items []JellyfinSeriesEpisodeItem `json:"Items"`
}
type JellyfinSeriesEpisodeItem struct {
JellyfinItems
// the episode number
IndexNumber int `json:"IndexNumber"`
// the episodes season number
ParentIndexNumber int `json:"ParentIndexNumber"`
}
type JellyfinSyncResponse struct {
JobId string `json:"jobId"`
}
// Perform the jellyfin sync.
// Gets each type of media separately from jellyfin and attempts to import them.
// Errors are added silently to the job.
func startJellyfinSync(
db *gorm.DB,
jobId string,
userId uint,
username string,
userThirdPartyId string,
userThirdPartyAuth string,
) {
// Get played movies
updateJobCurrentTask(jobId, userId, "syncing movies")
playedMovies := new(JellyfinItemSearchResponse)
err := jellyfinAPIRequest(
"GET",
"/Users/"+userThirdPartyId+"/Items",
map[string]string{
"Filters": "IsPlayed",
"IncludeItemTypes": "Movie",
"Fields": "ProviderIds",
"Recursive": "true",
},
username,
userThirdPartyAuth,
&playedMovies,
)
if err != nil {
slog.Error("jellyfinSyncWatched: Jellyfin API request failed", "error", err)
addJobError(jobId, userId, "failed to get jellyfin response for movies")
} else {
if len(playedMovies.Items) <= 0 {
slog.Info("jellyfinSyncWatched: User has no played movies.", "user_id", userId)
} else {
for _, v := range playedMovies.Items {
slog.Info("jellyfinSyncWatched: Importing played movie.", "movie_name", v.Name, "user_id", userId)
slog.Debug("jellyfinSyncWatched: Importing played movie.", "full_item", v, "user_id", userId)
// 1. Ensure we have a tmdbId
if v.ProviderIds.Tmdb == "" {
slog.Error("jellyfinSyncWatched: Movie to import does not have a tmdb id.", "movie_name", v.Name, "movie_ids", v.ProviderIds, "user_id", userId)
addJobError(jobId, userId, "movie could not be imported (no tmdbId present): "+v.Name)
continue
}
tmdbId, err := strconv.Atoi(v.ProviderIds.Tmdb)
if err != nil {
slog.Error("jellyfinSyncWatched: Movie to import does not have a parseable (to int) tmdb id.", "movie_name", v.Name, "movie_ids", v.ProviderIds, "user_id", userId)
addJobError(jobId, userId, "movie could not be imported (tmdbId was not parseable): "+v.Name)
continue
}
// 2. Imported watched movie
w, err := addWatched(db, userId, WatchedAddRequest{
Status: FINISHED,
ContentID: tmdbId,
ContentType: MOVIE,
WatchedDate: v.UserData.LastPlayedDate,
}, IMPORTED_WATCHED)
if err != nil {
if err.Error() == "content already on watched list" {
slog.Error("jellyfinSyncWatched: Unique constraint hit.. content must already be on watch list.", "movie_name", v.Name, "movie_ids", v.ProviderIds, "user_id", userId)
}
slog.Error("jellyfinSyncWatched: Movie failed to import.", "movie_name", v.Name, "movie_ids", v.ProviderIds, "user_id", userId)
addJobError(jobId, userId, "movie could not be imported (failed when adding to watched list): "+v.Name)
}
// 3. Add IMPORTED_ADDED_WATCHED activity
if !v.UserData.LastPlayedDate.IsZero() {
_, err := addActivity(db, userId, ActivityAddRequest{WatchedID: w.ID, Type: IMPORTED_ADDED_WATCHED, CustomDate: &v.UserData.LastPlayedDate})
if err != nil {
slog.Error("jellyfinSyncWatched: Failed to add dateswatched activity.", "movie_name", v.Name,
"movie_ids", v.ProviderIds, "user_id", userId, "date", v.UserData.LastPlayedDate, "error", err)
}
}
}
}
}
// Get played series
// Can't rely on IsPlayed filter, since we want to get partially played series too.
updateJobCurrentTask(jobId, userId, "syncing series")
allSeries := new(JellyfinItemSearchResponse)
err = jellyfinAPIRequest(
"GET",
"/Users/"+userThirdPartyId+"/Items",
map[string]string{
"IncludeItemTypes": "Series",
"Fields": "ProviderIds,RecursiveItemCount",
"Recursive": "true",
"IsPlaceHolder": "false",
},
username,
userThirdPartyAuth,
&allSeries,
)
if err != nil {
slog.Error("jellyfinSyncWatched: Jellyfin API request failed", "error", err)
addJobError(jobId, userId, "failed to get jellyfin response for series")
} else {
if len(allSeries.Items) <= 0 {
slog.Info("jellyfinSyncWatched: No series found.", "user_id", userId)
} else {
// Import series
for _, v := range allSeries.Items {
slog.Info("jellyfinSyncWatched: Processing series.", "series_name", v.Name, "user_id", userId)
slog.Debug("jellyfinSyncWatched: Processing series.", "full_item", v, "user_id", userId)
// 1. Make sure show is watched or at least partially watched
if !v.UserData.Played && v.UserData.PlayedPercentage <= 0 && v.RecursiveItemCount == v.UserData.UnplayedItemCount {
slog.Debug("jellyfinSyncWatched: Skipping unwatched series:", "series_name", v.Name, "user_id", userId)
continue
}
// 1.1. Ensure we have a tmdbId
if v.ProviderIds.Tmdb == "" {
slog.Error("jellyfinSyncWatched: Series to import does not have a tmdb id.", "series_name", v.Name, "series_ids", v.ProviderIds, "user_id", userId)
addJobError(jobId, userId, "series could not be imported (no tmdbId present): "+v.Name)
continue
}
tmdbId, err := strconv.Atoi(v.ProviderIds.Tmdb)
if err != nil {
slog.Error("jellyfinSyncWatched: Series to import does not have a parseable (to int) tmdb id.", "series_name", v.Name, "series_ids", v.ProviderIds, "user_id", userId)
addJobError(jobId, userId, "series could not be imported (tmdbId was not parseable): "+v.Name)
continue
}
updateJobCurrentTask(jobId, userId, "syncing serie "+v.Name)
// 2. Imported watched series
w, err := addWatched(db, userId, WatchedAddRequest{
Status: FINISHED,
ContentID: tmdbId,
ContentType: SHOW,
WatchedDate: v.UserData.LastPlayedDate,
}, IMPORTED_WATCHED)
if err != nil {
if err.Error() == "content already on watched list" {
slog.Info("jellyfinSyncWatched: Unique constraint hit.. content must already be on watch list.",
"series_name", v.Name, "series_ids", v.ProviderIds, "user_id", userId, "watched_id", w.ID)
} else {
slog.Error("jellyfinSyncWatched: Series failed to import.", "series_name", v.Name, "series_ids", v.ProviderIds, "user_id", userId)
addJobError(jobId, userId, "series could not be imported (failed when adding to watched list): "+v.Name)
}
} else {
// 3. Add IMPORTED_ADDED_WATCHED activity (only if no err above, show also must not have already been on our list)
if !v.UserData.LastPlayedDate.IsZero() {
_, err := addActivity(db, userId, ActivityAddRequest{WatchedID: w.ID, Type: IMPORTED_ADDED_WATCHED, CustomDate: &v.UserData.LastPlayedDate})
if err != nil {
slog.Error("jellyfinSyncWatched: Failed to add dateswatched activity.", "series_name", v.Name,
"series_ids", v.ProviderIds, "user_id", userId, "date", v.UserData.LastPlayedDate, "error", err)
}
}
}
// 4. Import watched seasons for this serie
// Get all show seasons (filtering isPlayed doesn't seem to be a thing, so we will have to do that ourselves)
seriesSeasons := new(JellyfinSeriesSeasonsResponse)
err = jellyfinAPIRequest(
"GET",
"/Shows/"+v.Id+"/Seasons",
map[string]string{
"UserId": userThirdPartyId,
"Fields": "ProviderIds",
"IsPlaceHolder": "false",
},
username,
userThirdPartyAuth,
&seriesSeasons,
)
if err != nil {
slog.Error("jellyfinSyncWatched: Failed to fetch series seasons.", "series_name", v.Name, "series_ids", v.ProviderIds, "user_id", userId)
addJobError(jobId, userId, "series seasons could not be imported (request failed): "+v.Name)
} else if len(seriesSeasons.Items) <= 0 {
slog.Info("jellyfinSyncWatched: Series has no seasons.", "series_name", v.Name, "series_ids", v.ProviderIds, "user_id", userId)
} else {
for _, vs := range seriesSeasons.Items {
if !vs.UserData.Played {
slog.Debug("jellyfinSyncWatched: Skipping import of unplayed season.", "series_name", v.Name, "season_num", vs.IndexNumber, "user_id", userId)
continue
}
updateJobCurrentTask(jobId, userId, "syncing "+v.Name+" season "+strconv.Itoa(vs.IndexNumber))
_, err = addWatchedSeason(db, userId, WatchedSeasonAddRequest{WatchedID: w.ID, SeasonNumber: vs.IndexNumber, Status: FINISHED})
if err != nil {
slog.Error("jellyfinSyncWatched: Failed to fetch series seasons.", "series_name", v.Name, "series_ids", v.ProviderIds, "user_id", userId)
addJobError(jobId, userId, "series season could not be imported (addWatchedSeason request failed): "+v.Name+" season "+strconv.Itoa(vs.IndexNumber))
}
}
}
// 5. Import watched episodes for this serie
// Gets all show episodes (filtering isPlayed doesn't seem to be a thing, so we will have to do that ourselves)
seriesEpisodes := new(JellyfinSeriesEpisodesResponse)
err = jellyfinAPIRequest(
"GET",
"/Shows/"+v.Id+"/Episodes",
map[string]string{
"UserId": userThirdPartyId,
"Fields": "ProviderIds",
"IsPlaceHolder": "false",
},
username,
userThirdPartyAuth,
&seriesEpisodes,
)
if err != nil {
slog.Error("jellyfinSyncWatched: Failed to fetch series episodes.", "series_name", v.Name, "series_ids", v.ProviderIds, "user_id", userId)
addJobError(jobId, userId, "series episodes could not be imported (request failed): "+v.Name)
} else if len(seriesEpisodes.Items) <= 0 {
slog.Info("jellyfinSyncWatched: Series has no episodes.", "series_name", v.Name, "series_ids", v.ProviderIds, "user_id", userId)
} else {
for _, vs := range seriesEpisodes.Items {
if !vs.UserData.Played {
slog.Debug("jellyfinSyncWatched: Skipping import of unplayed episode.", "series_name", v.Name, "season_num", vs.ParentIndexNumber, "episode_num", vs.IndexNumber, "user_id", userId)
continue
}
updateJobCurrentTask(jobId, userId, "syncing "+v.Name+" season "+strconv.Itoa(vs.ParentIndexNumber)+" episode "+strconv.Itoa(vs.IndexNumber))
_, err = addWatchedEpisodes(db, userId, WatchedEpisodeAddRequest{
WatchedID: w.ID,
SeasonNumber: vs.ParentIndexNumber,
EpisodeNumber: vs.IndexNumber,
Status: FINISHED,
})
if err != nil {
slog.Error("jellyfinSyncWatched: Failed to import series episode.", "series_name", v.Name, "season_num", vs.ParentIndexNumber, "episode_num", vs.IndexNumber, "user_id", userId)
addJobError(jobId, userId, "series episode could not be imported (addWatchedEpisode request failed): "+v.Name+" "+vs.Name)
}
}
}
}
}
}
}
func jellyfinSyncWatched(
db *gorm.DB,
userId uint,
userType UserType,
username string,
userThirdPartyId string,
userThirdPartyAuth string,
) (JellyfinSyncResponse, error) {
jobId, err := addJob("jf_sync", userId)
if err != nil {
slog.Error("jellyfinSyncWatched: Failed to create a job", "error", err)
return JellyfinSyncResponse{}, errors.New("failed to create job")
}
updateJobStatus(jobId, userId, JOB_RUNNING)
go startJellyfinSync(
db,
jobId,
userId,
username,
userThirdPartyId,
userThirdPartyAuth,
)
return JellyfinSyncResponse{JobId: jobId}, nil
}
+104
View File
@@ -0,0 +1,104 @@
// We will use jobs only for storing/retrieving active job statuses for the client.
// Running the job will be done wherever needed, but isn't handled here.
// When starting a job elsewhere, we should first add a job here as active to get an `id`,
// this id should be used to update the active job so the client can request job status updates.
package main
import (
"errors"
"log/slog"
)
type JobStatus string
var (
JOB_CREATED JobStatus = "CREATED"
JOB_RUNNING JobStatus = "RUNNING"
JOB_COMPLETED JobStatus = "COMPLETED"
JOB_FAILED JobStatus = "FAILED"
)
type Job struct {
// We can give the job a name simply for showing on the client.
Name string `json:"name"`
// The current status of the job.
Status JobStatus `json:"status"`
// The current task we are performing inside the job.
// Just so we can portray progress on the client by displaying the current task.
CurrentTask string `json:"currentTask,omitempty"`
// Errors that occurred in the task
Errors []string `json:"errors"`
// Stored for access control.
UserId uint `json:"-"`
}
var activeJobs = make(map[string]*Job)
// Add a job to our activeJobs map.
// Returns id of job on success, or error if failed to add.
func addJob(name string, userId uint) (string, error) {
idk, err := generateString(8)
if err != nil {
return "", err
}
_, ok := activeJobs[idk]
if ok {
// Lets just hope this doesn't happen, may the odds be with us.
return "", errors.New("job already exists with id generated, please try again")
}
activeJobs[idk] = &Job{
Name: name,
Status: JOB_CREATED,
UserId: userId,
}
return idk, nil
}
// Get a job.
// Returns job if found, otherwise errors if job does not exist.
func getJob(id string, userId uint) (*Job, error) {
j, ok := activeJobs[id]
if ok {
// Ensure user requesting a job, owns the job.
if j.UserId != userId {
slog.Warn("getJob: A user tried to access a job they do not own.", "user_id", userId, "job_id", id)
return &Job{}, errors.New("job does not exist")
}
return j, nil
}
return &Job{}, errors.New("job does not exist")
}
// Update a jobs status.
func updateJobStatus(id string, userId uint, status JobStatus) error {
j, err := getJob(id, userId)
if err != nil {
slog.Error("updateJobStatus: Failed!", "status", status, "error", err)
return err
}
j.Status = status
return nil
}
// Update a jobs current task.
func updateJobCurrentTask(id string, userId uint, ct string) error {
j, err := getJob(id, userId)
if err != nil {
slog.Error("updateJobCurrentTask: Failed!", "ct", ct, "error", err)
return err
}
j.CurrentTask = ct
return nil
}
// Add an error to a job.
func addJobError(id string, userId uint, e string) error {
j, err := getJob(id, userId)
if err != nil {
slog.Error("updateJobCurrentTask: Failed!", "e", e, "error", err)
return err
}
j.Errors = append(j.Errors, e)
return nil
}
+30 -1
View File
@@ -663,7 +663,7 @@ func (b *BaseRouter) addProfileRoutes() {
}
func (b *BaseRouter) addJellyfinRoutes() {
jf := b.rg.Group("/jellyfin").Use(AuthRequired(b.db))
jf := b.rg.Group("/jellyfin").Use(AuthRequired(b.db), JellyfinAccessRequired())
// Check if jf has item
jf.GET("/:type/:name/:tmdbId", func(c *gin.Context) {
@@ -679,6 +679,21 @@ func (b *BaseRouter) addJellyfinRoutes() {
}
c.JSON(http.StatusOK, response)
})
// Sync users jellyfin watched items to watchlist
jf.GET("/sync", func(c *gin.Context) {
userId := c.MustGet("userId").(uint)
userType := c.MustGet("userType").(UserType)
username := c.MustGet("username").(string)
userThirdPartyId := c.MustGet("userThirdPartyId").(string)
userThirdPartyAuth := c.MustGet("userThirdPartyAuth").(string)
response, err := jellyfinSyncWatched(b.db, userId, userType, username, userThirdPartyId, userThirdPartyAuth)
if err != nil {
c.JSON(http.StatusForbidden, ErrorResponse{Error: err.Error()})
return
}
c.JSON(http.StatusOK, response)
})
}
func (b *BaseRouter) addUserRoutes() {
@@ -1109,3 +1124,17 @@ func (b *BaseRouter) addRadarrRoutes() {
c.AbortWithStatusJSON(http.StatusBadRequest, ErrorResponse{Error: err.Error()})
})
}
func (b *BaseRouter) addJobRoutes() {
job := b.rg.Group("/job").Use(AuthRequired(nil))
job.GET("/:id", func(c *gin.Context) {
userId := c.MustGet("userId").(uint)
response, err := getJob(c.Param("id"), userId)
if err != nil {
c.JSON(http.StatusForbidden, ErrorResponse{Error: err.Error()})
return
}
c.JSON(http.StatusOK, *response)
})
}
+1
View File
@@ -138,6 +138,7 @@ func main() {
br.addFeatureRoutes()
br.addSonarrRoutes()
br.addRadarrRoutes()
br.addJobRoutes()
br.rg.Static("/img", path.Join(DataPath, "img"))
go setupTasks(db)
+2 -1
View File
@@ -163,6 +163,7 @@ func addWatched(db *gorm.DB, userId uint, ar WatchedAddRequest, at ActivityType)
}
// If custom WatchedDate passed, set CreatedAt and UpdatedAt fields to it.
if !ar.WatchedDate.IsZero() {
slog.Debug("Adding watched item: The provided WatchedDate is valid.", "watched_date", ar.WatchedDate, "userId", userId, "contentType", ar.ContentType, "contentId", ar.ContentID)
watched.CreatedAt = ar.WatchedDate
watched.UpdatedAt = ar.WatchedDate
}
@@ -174,7 +175,7 @@ func addWatched(db *gorm.DB, userId uint, ar WatchedAddRequest, at ActivityType)
return Watched{}, errors.New("content already on watched list. errored checking for soft deleted record")
}
if watched.DeletedAt.Time.IsZero() {
return Watched{}, errors.New("content already on watched list")
return watched, errors.New("content already on watched list")
} else {
slog.Info("addWatched: Watched list item for this content exists as soft deleted record.. attempting to restore")
res = db.Model(&Watched{}).Unscoped().Where("user_id = ? AND content_id = ?", userId, watched.ContentID).Updates(map[string]interface{}{"status": ar.Status, "rating": ar.Rating, "deleted_at": nil})